"""Digital ETL related queries.""" from flows import config as flows_config from flows.digital import config from flows.queries import sql_upcs_condition_in from flows.util import batch_param_calls # The select statement of a datawarehouse unload query. The columns are mapped # to the datastore's digital_revenue table without the primary key. UNLOAD_SELECT = """ SELECT fa.releaseid AS upc, dr.display_upc AS display_upc, sum(fa.royaltydollar) AS amount, dd.displaydate AS date, fa.transactiontypeid AS transaction_type_id, fa.countryid AS country_id, sum(fa.paidunits) + sum(fa.freeunits) AS units, CASE WHEN sum(fa.freeunits) > 0 THEN 0 ELSE 1 END AS paid, fa.formatid AS format_id FROM facts.prod.fact_analytics fa LEFT JOIN facts.prod.dim_day dd ON dd.dayid = fa.dayid LEFT JOIN facts.prod.dim_release dr ON dr.releaseid = fa.releaseid WHERE {upcs_where_clause} dd.displaydate >= '{date_start}' AND dd.displaydate < '{date_end}' GROUP BY upc, display_upc, date, transaction_type_id, country_id, format_id """ UNLOAD = """ COPY INTO '{destination}' FROM ({select}) CREDENTIALS=( AWS_KEY_ID='{access_key_id}' AWS_SECRET_KEY='{secret_access_key}') FILE_FORMAT=( TYPE='CSV' FIELD_DELIMITER=',' DATE_FORMAT='YYYY-MM-DD' RECORD_DELIMITER='\\n' COMPRESSION='GZIP' ESCAPE='\\134' ESCAPE_UNENCLOSED_FIELD='\\134' NULL_IF=('NULL', '__NULL__') FIELD_OPTIONALLY_ENCLOSED_BY='"' TRIM_SPACE=TRUE ) SINGLE = TRUE MAX_FILE_SIZE = {max_size} """ CREATE_TEMP_TABLE = """ CREATE TABLE IF NOT EXISTS {table_name} LIKE digital_revenue""" TEMP_TABLE_NAME = 'digital_revenue_{correlation_hex}' INSERT_TO_TEMP_TABLE = """ INSERT INTO {table_name} ( upc, display_upc, amount, date, transaction_type_id, country_id, units, paid, format_id) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)""" DELETE_FROM_DIGITAL_REVENUE = """ DELETE FROM digital_revenue WHERE {upc_in_clause} date >= %(date_start)s AND date < %(date_end)s""" INSERT_FROM_TEMP_TABLE = """ INSERT INTO digital_revenue ( upc, display_upc, amount, date, transaction_type_id, country_id, units, paid, format_id) SELECT upc, display_upc, amount, date, transaction_type_id, country_id, units, paid, format_id FROM {table_name}""" DROP_TEMP_TABLE = 'DROP TABLE {table_name}' @batch_param_calls('upcs', flows_config.ITERATION_BATCH_SIZE) def get_unload_from_snowflake_sql( bucket, correlation_id, date_end, date_start, upcs): """Generate snowflake unload query templates for upcs. This function is a generator due to the batching decorator. Args: bucket (str): bucket to unload to. correlation_id (str): activity correlation ID. date_end (str): YYYY-MM-DD format end of date range query. date_start (str): YYYY-MM-DD format start of date range query. upcs (list): upcs of query. Yields: str: sql format string with "batch" parameter. """ s3_destination = config.UNLOAD_DESTINATION.format( bucket=bucket, correlation_id=correlation_id) upcs_condition = sql_upcs_condition_in(upcs, column_name='fa.releaseid') select_sql = UNLOAD_SELECT.format( date_start=date_start, date_end=date_end, upcs_where_clause=upcs_condition) sql = UNLOAD.format( access_key_id=flows_config.AWS_CREDENTIALS['aws_access_key_id'], destination=s3_destination, options=config.UNLOAD_OPTIONS, secret_access_key=( flows_config.AWS_CREDENTIALS['aws_secret_access_key']), select=select_sql, max_size=config.S3_UPLOAD_MAX_SIZE_BYTES) return sql @batch_param_calls('upcs', flows_config.ITERATION_BATCH_SIZE) def get_delete_from_digital_revenue_sql(upcs): """Generate delete digital_revenue rows sql with placeholders. Args: upcs (list): upcs of query. Yields: str: sql with db api placeholders for date ranges. """ upc_in_clause = sql_upcs_condition_in(upcs, column_name='upc') sql = DELETE_FROM_DIGITAL_REVENUE.format(upc_in_clause=upc_in_clause) return sql