import json from garcon import task from garcon.contrib import redshift_util as util import psycopg2 import psycopg2.extras @task.decorate(timeout=7200) def unload( activity, query, s3_prefix, options, redshift_host, redshift_port, redshift_db, redshift_user, redshift_password, aws_access_key, aws_access_secret, ): """Unload data from AWS Redshift onto S3. Args: activity (ActivityWorker): The swf activity worker query (str): the query to run on the redshift cluster s3_prefix (str): destination of the query results options (str): options portion of the sql literal redshift_host (str): redshift host redshift_port (int): redshift port redshift_db (str): redshift database redshift_user (str): redshift user login redshift_password (str): redshift user password aws_access_key (str): aws access key aws_access_secret(str): aws access secret key """ activity.logger.info('Starting redshift unload to s3') # prepare the unload sql sql = util.unload_sql( query=query, s3_prefix=s3_prefix, aws_access_key=aws_access_key, aws_access_secret=aws_access_secret, options=options, ) # log cleaned sql sql_debug = sql sql_debug = sql_debug.replace(aws_access_key, '[...]') sql_debug = sql_debug.replace(aws_access_secret, '[...]') activity.logger.debug(sql_debug) # execute unload sql with psycopg2.connect( host=redshift_host, port=redshift_port, user=redshift_user, password=redshift_password, database=redshift_db, ) as conn: with conn.cursor() as cur: cur.execute(sql) activity.logger.info( 'Unload data from redshift redshift unload was successful.') @task.decorate(timeout=7200) def describe_table( activity, schema, table, redshift_host, redshift_port, redshift_db, redshift_user, redshift_password, ): """Describes columns of a Redshift table. Args: activity (ActivityWorker): The swf activity worker schema (str): redshift schema name table (str): redshift table name redshift_host (str): redshift host redshift_port (int): redshift port redshift_db (str): redshift database redshift_user (str): redshift user login redshift_password (str): redshift user password Returns: dict: context with describe json string """ activity.logger.debug('Describing redshift table: %s.%s', schema, table) # prepare sql sql = util.describe_table_sql(schema=schema, table=table) # execute unload sql with psycopg2.connect( host=redshift_host, port=redshift_port, user=redshift_user, password=redshift_password, database=redshift_db, ) as conn: with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur: cur.execute(sql) results = cur.fetchall() # prepare results json _json = json.dumps(dict(schema=schema, table=table, columns=results)) return dict(describe_json=_json)