"""DB Tasks. ========================== Garcon tasks for executing db updates. """ import subprocess from garcon import task from snowflake_connector.etl_connector import SnowflakeSQLExecutor @task.decorate(timeout=72000) def art_relation_sql_to_stdout(activity, query): """Extract data from the art_relation db (MySQL) to stdout. For use in conjunction with other tasks streaming mysql dimension inputs. See () Args: activity (ActivityWorker): The swf activity worker. query (str): query to be piped Return: dict: A dictionary with a Popen object """ activity.logger.info( 'Setting up query to stream to stdout: {query}'.format(query=query)) command = 'use art_relations; {}'.format(query) p1 = subprocess.Popen( ['echo', command], stdout=subprocess.PIPE, close_fds=True) return dict(pipe=p1) @task.decorate(timeout=72000) def snowflake_execute_query(activity, query, sf_config): """Execute a Snowflake query. Args: activity (ActivityWorker): The SWF activity worker. sf_config (dict): Snowflake credentials. query (str): sql statement to be executed on Snowflake. """ with SnowflakeSQLExecutor(sf_config=sf_config) as executor: activity.logger.info('Executing query:\n{}'.format(query)) executor.execute(query) @task.decorate(timeout=72000) def snowflake_execute_query_list(activity, query_list, sf_config): """Execute a list of Snowflake queries. Args: activity (ActivityWorker): the activity worker. query_list (list): list of sql strings to execute on Snowflake sf_config (dict): Snowflake credentials. """ if not query_list: activity.logger.info('No query list sql to run') return activity.logger.info( 'Starting to execute query list {}'.format(query_list)) for query in query_list: with SnowflakeSQLExecutor(sf_config=sf_config) as executor: activity.logger.info('Executing query:\n{}'.format(query)) executor.execute(query) activity.logger.info('Finished executing query list {}'.format(query_list)) @task.decorate(timeout=72000) def snowflake_copy_from_s3( activity, sf_config, aws, table, file_format, s3_path): """Task to load a dimension staging table from s3. Args: activity (ActivityWorker): The activity worker. sf_config (dict): Snowflake credentials. aws (dict): Dictionary stores AWS credentials. Keys should include 'access_key' & 'access_secret'. table (str): name of Snowflake staging table to load data into. s3_path (str): S3 path to COPY data from ex. for s3://foo/bar/000.tz, s3_staging_path=s3://foo/bar/. file_format (str): Snowflake FILE_FORMAT options. """ copy_sql = ( 'COPY INTO {db}.{schema}.{table} FROM {s3_path} ' 'FILE_FORMAT = ( {file_format} ) ' 'CREDENTIALS=(' "AWS_KEY_ID='{access_key}' " "AWS_SECRET_KEY='{secret_key}');".format( table=table, s3_path=s3_path, file_format=file_format, db=sf_config['db'], schema=sf_config['schema'], access_key=aws['access_key'], secret_key=aws['access_secret'])) with SnowflakeSQLExecutor(sf_config=sf_config) as executor: activity.logger.info( 'Executing COPY INTO from S3 to {}...'.format(table)) executor.execute(copy_sql) @task.decorate(timeout=72000) def snowflake_copy_to_s3( activity, sf_config, aws, file_format, s3_path, query): """Task to unload a table to S3. Args: activity (ActivityWorker): The activity worker. sf_config (dict): Snowflake credentials. aws (dict): Dictionary stores AWS credentials. Keys should include 'access_key' & 'access_secret'. query (str): SQL statement to use within COPY INTO. s3_path (str): S3 path to COPY data from ex. for s3://foo/bar/000.tz, s3_staging_path=s3://foo/bar/. file_format (str): Snowflake FILE_FORMAT options. """ copy_sql = ( 'COPY INTO {s3_path} FROM ({query}) ' 'FILE_FORMAT = ( {file_format} ) ' 'CREDENTIALS=(' "AWS_KEY_ID='{access_key}' " "AWS_SECRET_KEY='{secret_key}');".format( s3_path=s3_path, query=query, file_format=file_format, access_key=aws['access_key'], secret_key=aws['access_secret'])) with SnowflakeSQLExecutor(sf_config=sf_config) as executor: activity.logger.info( 'Executing {}...'.format(query)) executor.execute(copy_sql)