misc.python.materialize.zippy.copy_s3_actions
1# Copyright Materialize, Inc. and contributors. All rights reserved. 2# 3# Use of this software is governed by the Business Source License 4# included in the LICENSE file at the root of this repository. 5# 6# As of the Change Date specified in that file, in accordance with 7# the Business Source License, use of this software will be governed 8# by the Apache License, Version 2.0. 9 10import random 11from textwrap import dedent 12 13from materialize.mzcompose.composition import Composition 14from materialize.zippy.balancerd_capabilities import BalancerdIsRunning 15from materialize.zippy.copy_s3_capabilities import S3ObjectExists 16from materialize.zippy.framework import Action, Capabilities, Capability, State 17from materialize.zippy.mz_capabilities import MzIsRunning 18from materialize.zippy.table_capabilities import TableExists 19 20 21class CopyToS3(Action): 22 """Performs a COPY TO S3 and records the resulting S3 object as a capability.""" 23 24 @classmethod 25 def requires(cls) -> set[type[Capability]]: 26 return {BalancerdIsRunning, MzIsRunning, TableExists} 27 28 def __init__(self, capabilities: Capabilities) -> None: 29 self.table = random.choice(capabilities.get(TableExists)) 30 super().__init__(capabilities) 31 self.s3_key = f"zippy/{self.seqno}" 32 self.s3_object = S3ObjectExists( 33 s3_key=self.s3_key, 34 min_val=self.table.watermarks.min, 35 max_val=self.table.watermarks.max, 36 ) 37 38 def provides(self) -> list[Capability]: 39 return [self.s3_object] 40 41 def run(self, c: Composition, state: State) -> None: 42 conn_name = f"zippy_aws_conn_{self.seqno}" 43 secret_name = f"zippy_minio_secret_{self.seqno}" 44 45 c.testdrive( 46 dedent(f""" 47 > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}' 48 > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name}) 49 > COPY (SELECT f1 FROM {self.table.name}) TO 's3://copytos3/{self.s3_key}' WITH (AWS CONNECTION = {conn_name}, FORMAT = 'csv') 50 > DROP CONNECTION {conn_name} 51 > DROP SECRET {secret_name} 52 """), 53 mz_service=state.mz_service, 54 ) 55 self.s3_object.min_val = self.table.watermarks.min 56 self.s3_object.max_val = self.table.watermarks.max 57 58 59class CopyFromS3(Action): 60 """Reads a previously written S3 object back into a staging table and validates the row count.""" 61 62 @classmethod 63 def requires(cls) -> set[type[Capability]]: 64 return {BalancerdIsRunning, MzIsRunning, S3ObjectExists} 65 66 def __init__(self, capabilities: Capabilities) -> None: 67 self.s3_object = random.choice(capabilities.get(S3ObjectExists)) 68 super().__init__(capabilities) 69 70 def run(self, c: Composition, state: State) -> None: 71 conn_name = f"zippy_aws_conn_{self.seqno}" 72 secret_name = f"zippy_minio_secret_{self.seqno}" 73 staging_table = f"zippy_s3_staging_{self.seqno}" 74 75 c.testdrive( 76 dedent(f""" 77 > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}' 78 > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name}) 79 > CREATE TABLE {staging_table} (f1 INTEGER) 80 > COPY INTO {staging_table} FROM 's3://copytos3/{self.s3_object.s3_key}' (FORMAT CSV, AWS CONNECTION = {conn_name}) 81 > SELECT MIN(f1), MAX(f1), COUNT(f1), COUNT(DISTINCT f1) FROM {staging_table} 82 {self.s3_object.min_val} {self.s3_object.max_val} {self.s3_object.max_val - self.s3_object.min_val + 1} {self.s3_object.max_val - self.s3_object.min_val + 1} 83 > DROP TABLE {staging_table} 84 > DROP CONNECTION {conn_name} 85 > DROP SECRET {secret_name} 86 """), 87 mz_service=state.mz_service, 88 )
class
CopyToS3(materialize.zippy.framework.Action):
22class CopyToS3(Action): 23 """Performs a COPY TO S3 and records the resulting S3 object as a capability.""" 24 25 @classmethod 26 def requires(cls) -> set[type[Capability]]: 27 return {BalancerdIsRunning, MzIsRunning, TableExists} 28 29 def __init__(self, capabilities: Capabilities) -> None: 30 self.table = random.choice(capabilities.get(TableExists)) 31 super().__init__(capabilities) 32 self.s3_key = f"zippy/{self.seqno}" 33 self.s3_object = S3ObjectExists( 34 s3_key=self.s3_key, 35 min_val=self.table.watermarks.min, 36 max_val=self.table.watermarks.max, 37 ) 38 39 def provides(self) -> list[Capability]: 40 return [self.s3_object] 41 42 def run(self, c: Composition, state: State) -> None: 43 conn_name = f"zippy_aws_conn_{self.seqno}" 44 secret_name = f"zippy_minio_secret_{self.seqno}" 45 46 c.testdrive( 47 dedent(f""" 48 > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}' 49 > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name}) 50 > COPY (SELECT f1 FROM {self.table.name}) TO 's3://copytos3/{self.s3_key}' WITH (AWS CONNECTION = {conn_name}, FORMAT = 'csv') 51 > DROP CONNECTION {conn_name} 52 > DROP SECRET {secret_name} 53 """), 54 mz_service=state.mz_service, 55 ) 56 self.s3_object.min_val = self.table.watermarks.min 57 self.s3_object.max_val = self.table.watermarks.max
Performs a COPY TO S3 and records the resulting S3 object as a capability.
CopyToS3(capabilities: materialize.zippy.framework.Capabilities)
29 def __init__(self, capabilities: Capabilities) -> None: 30 self.table = random.choice(capabilities.get(TableExists)) 31 super().__init__(capabilities) 32 self.s3_key = f"zippy/{self.seqno}" 33 self.s3_object = S3ObjectExists( 34 s3_key=self.s3_key, 35 min_val=self.table.watermarks.min, 36 max_val=self.table.watermarks.max, 37 )
Construct a new action, possibly conditioning on the available capabilities.
@classmethod
def
requires(cls) -> set[type[materialize.zippy.framework.Capability]]:
25 @classmethod 26 def requires(cls) -> set[type[Capability]]: 27 return {BalancerdIsRunning, MzIsRunning, TableExists}
Compute the capability classes that this action requires.
def
provides(self) -> list[materialize.zippy.framework.Capability]:
Compute the capabilities that this action will make available.
def
run( self, c: materialize.mzcompose.composition.Composition, state: materialize.zippy.framework.State) -> None:
42 def run(self, c: Composition, state: State) -> None: 43 conn_name = f"zippy_aws_conn_{self.seqno}" 44 secret_name = f"zippy_minio_secret_{self.seqno}" 45 46 c.testdrive( 47 dedent(f""" 48 > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}' 49 > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name}) 50 > COPY (SELECT f1 FROM {self.table.name}) TO 's3://copytos3/{self.s3_key}' WITH (AWS CONNECTION = {conn_name}, FORMAT = 'csv') 51 > DROP CONNECTION {conn_name} 52 > DROP SECRET {secret_name} 53 """), 54 mz_service=state.mz_service, 55 ) 56 self.s3_object.min_val = self.table.watermarks.min 57 self.s3_object.max_val = self.table.watermarks.max
Run this action on the provided composition.
class
CopyFromS3(materialize.zippy.framework.Action):
60class CopyFromS3(Action): 61 """Reads a previously written S3 object back into a staging table and validates the row count.""" 62 63 @classmethod 64 def requires(cls) -> set[type[Capability]]: 65 return {BalancerdIsRunning, MzIsRunning, S3ObjectExists} 66 67 def __init__(self, capabilities: Capabilities) -> None: 68 self.s3_object = random.choice(capabilities.get(S3ObjectExists)) 69 super().__init__(capabilities) 70 71 def run(self, c: Composition, state: State) -> None: 72 conn_name = f"zippy_aws_conn_{self.seqno}" 73 secret_name = f"zippy_minio_secret_{self.seqno}" 74 staging_table = f"zippy_s3_staging_{self.seqno}" 75 76 c.testdrive( 77 dedent(f""" 78 > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}' 79 > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name}) 80 > CREATE TABLE {staging_table} (f1 INTEGER) 81 > COPY INTO {staging_table} FROM 's3://copytos3/{self.s3_object.s3_key}' (FORMAT CSV, AWS CONNECTION = {conn_name}) 82 > SELECT MIN(f1), MAX(f1), COUNT(f1), COUNT(DISTINCT f1) FROM {staging_table} 83 {self.s3_object.min_val} {self.s3_object.max_val} {self.s3_object.max_val - self.s3_object.min_val + 1} {self.s3_object.max_val - self.s3_object.min_val + 1} 84 > DROP TABLE {staging_table} 85 > DROP CONNECTION {conn_name} 86 > DROP SECRET {secret_name} 87 """), 88 mz_service=state.mz_service, 89 )
Reads a previously written S3 object back into a staging table and validates the row count.
CopyFromS3(capabilities: materialize.zippy.framework.Capabilities)
67 def __init__(self, capabilities: Capabilities) -> None: 68 self.s3_object = random.choice(capabilities.get(S3ObjectExists)) 69 super().__init__(capabilities)
Construct a new action, possibly conditioning on the available capabilities.
@classmethod
def
requires(cls) -> set[type[materialize.zippy.framework.Capability]]:
63 @classmethod 64 def requires(cls) -> set[type[Capability]]: 65 return {BalancerdIsRunning, MzIsRunning, S3ObjectExists}
Compute the capability classes that this action requires.
def
run( self, c: materialize.mzcompose.composition.Composition, state: materialize.zippy.framework.State) -> None:
71 def run(self, c: Composition, state: State) -> None: 72 conn_name = f"zippy_aws_conn_{self.seqno}" 73 secret_name = f"zippy_minio_secret_{self.seqno}" 74 staging_table = f"zippy_s3_staging_{self.seqno}" 75 76 c.testdrive( 77 dedent(f""" 78 > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}' 79 > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name}) 80 > CREATE TABLE {staging_table} (f1 INTEGER) 81 > COPY INTO {staging_table} FROM 's3://copytos3/{self.s3_object.s3_key}' (FORMAT CSV, AWS CONNECTION = {conn_name}) 82 > SELECT MIN(f1), MAX(f1), COUNT(f1), COUNT(DISTINCT f1) FROM {staging_table} 83 {self.s3_object.min_val} {self.s3_object.max_val} {self.s3_object.max_val - self.s3_object.min_val + 1} {self.s3_object.max_val - self.s3_object.min_val + 1} 84 > DROP TABLE {staging_table} 85 > DROP CONNECTION {conn_name} 86 > DROP SECRET {secret_name} 87 """), 88 mz_service=state.mz_service, 89 )
Run this action on the provided composition.