""" Transform data from staging_raw table and load into fact tables. Used by FlowLoadFactMixinSF mixin. """ from garcon import task from feed_ingestion.conf.config import merge_configs from feed_ingestion.flows import registered_executors from feed_ingestion.flows.helpers import get_sf_config from feed_ingestion.util import task_status TASK_ID = 'load_fact_data' @task.decorate(timeout=7200) def create_staging_fact( activity, feed_name, date, sfdb_params, secrets_path=None, kwargs=None): """Create temp staging fact table for unloading fact and fact error data. Args: activity (ActivityWorker): The activity worker. feed_name (str): Name of the feed to get executor class. date (str): Date of the data being process (YYYY-MM-DD). sfdb_params (dict): Dict with params to optionally override default ones (Snowflake db and schema name). secrets_path (str): Secrets manager path of the flow. """ # keyword arguments for the task kwargs = kwargs or {} # completion check if task_status.is_completed_task(feed_name, date, TASK_ID): activity.logger.info( 'Task {task_id} of {feed_name} for {date} already complete, and ' '"reload" flag was not passed, skipping...'.format( task_id=TASK_ID, feed_name=feed_name, date=date)) return sf_config = get_sf_config(secrets_path) sf_config_custom = merge_configs(sf_config, sfdb_params) ExecutorFA = registered_executors.get(feed_name) with ExecutorFA(sf_config_custom) as sf_executor: sf_executor.drop_table(sf_executor.staging_fact_table(date)) activity.logger.info('Staging fact table was dropped, if existed') sf_executor.create_staging_fact_table(date, **kwargs) activity.logger.info('Staging fact table was created') @task.decorate(timeout=60*60*4) def load_staging_fact( activity, feed_name, date, sfdb_params, secrets_path=None, kwargs=None): """Load staging fact_analytics table from staging_raw table. Args: activity (ActivityWorker): The activity worker. feed_name (str): Name of the feed to get executor class. date (str): Date of the data being process (YYYY-MM-DD). sfdb_params (dict): Dict with params to optionally override default ones (Snowflake db and schema name). secrets_path (str): Secrets manager path of the flow. """ # keyword arguments for the task kwargs = kwargs or {} if task_status.is_completed_task(feed_name, date, TASK_ID): activity.logger.info( 'Task {task_id} of {feed_name} for {date} already complete, and ' '"reload" flag was not passed, skipping...'.format( task_id=TASK_ID, feed_name=feed_name, date=date)) return sf_config = get_sf_config(secrets_path) sf_config_custom = merge_configs(sf_config, sfdb_params) ExecutorFA = registered_executors.get(feed_name) with ExecutorFA(sf_config_custom) as sf_executor: sf_executor.load_staging_fact_table(date, **kwargs) activity.logger.info(f'Data was loaded into staging fact table ' f'feed_name {feed_name} date {date}') @task.decorate(timeout=14000) def load_aggregated_skips_and_saves(activity, feed_name, date, sfdb_params, secrets_path=None, kwargs=None): """Load aggregated_skips_and_saves from staging fact_analytics table. Args: activity (ActivityWorker): The activity worker. feed_name (str): Name of the feed to get executor class. date (str): Date of the data being process (YYYY-MM-DD). sfdb_params (dict): Dict with params to optionally override default ones (Snowflake db and schema name). secrets_path (str): Secrets manager path of the flow. kwargs (dict): Custom activity params. Passed to snowflake_executor """ task_id = 'load_aggregated_skips_and_saves' kwargs = kwargs or {} if task_status.is_completed_task(feed_name, date, task_id): activity.logger.info( 'Task {task_id} of {feed_name} for {date} already complete, and ' '"reload" flag was not passed, skipping...'.format( task_id=task_id, feed_name=feed_name, date=date)) return sf_config = get_sf_config(secrets_path) sf_config_custom = merge_configs(sf_config, sfdb_params) ExecutorFA = registered_executors.get(feed_name) with ExecutorFA(sf_config_custom) as sf_executor: sf_executor.delete_from_aggregated_skips_and_saves(date, **kwargs) activity.logger.info( 'Data was deleted from load_aggregated_skips_and_saves table') sf_executor.load_aggregated_skips_and_saves(date, **kwargs) activity.logger.info( 'Data was loaded into load_aggregated_skips_and_saves table') task_status.mark_completed_task(feed_name, date, task_id) activity.logger.info( 'Task {task_id} of {feed_name} for {date} completed'.format( task_id=task_id, feed_name=feed_name, date=date)) @task.decorate(timeout=60*60*5) def load_fact_data( activity, feed_name, date, sfdb_params, secrets_path=None, kwargs=None): """Load data on fact_analytics & fact_analytics_error. Args: activity (ActivityWorker): The activity worker. feed_name (str): Name of the feed to get executor class. date (str): Date of the data being process (YYYY-MM-DD). sfdb_params (dict): Dict with params to optionally override default ones (Snowflake db and schema name). kwargs (dict): Custom activity params. secrets_path (str): Secrets manager path of the flow. """ kwargs = kwargs or {} completed = task_status.is_completed_report(feed_name, date) if task_status.is_completed_task(feed_name, date, TASK_ID) and completed: activity.logger.info( 'Task {task_id} of {feed_name} for {date} already complete, and ' '"reload" flag was not passed, skipping...'.format( task_id=TASK_ID, feed_name=feed_name, date=date)) return sf_config = get_sf_config(secrets_path) sf_config_custom = merge_configs(sf_config, sfdb_params) ExecutorFA = registered_executors.get(feed_name) with ExecutorFA(sf_config_custom) as sf_executor: message = 'Deleting rows for {date} from fact_analytics for feed ' \ '{feed_name} and kwargs {kwargs}'.format(date=date, feed_name=feed_name, kwargs=kwargs) activity.logger.info(message) sf_executor.delete_from_fact_table(date, **kwargs) activity.logger.info( 'Loading fact_analytics for feed {}'.format(feed_name)) sf_executor.load_fact_data(date, **kwargs) with ExecutorFA(sf_config_custom) as sf_executor: activity.logger.info( 'Deleting rows for {date} from fact_analytics_error for feed ' '{feed_name}'.format(date=date, feed_name=feed_name)) sf_executor.delete_from_fact_error_table(date, **kwargs) activity.logger.info( 'Loading fact_analytics_error for feed {}'.format(feed_name)) sf_executor.load_fact_error_data(date, **kwargs) sf_executor.drop_table(sf_executor.staging_fact_table(date)) activity.logger.info('Staging fact table was dropped') activity.logger.info( 'Data was loaded into fact_analytics and fact_analytics_error') task_status.mark_completed_task(feed_name, date, TASK_ID) activity.logger.info( 'Task {task_id} of {feed_name} for {date} completed'.format( task_id=TASK_ID, feed_name=feed_name, date=date)) @task.decorate(timeout=21600) def update_dim_tables( activity, feed_name, date, sfdb_params, secrets_path=None, kwargs=None): """Update dimension tables. Args: activity (ActivityWorker): The activity worker. feed_name (str): Feed name of workflow execution for status updates. date (str): Reporting date (YYYY-MM-DD). sfdb_params (dict): Dict with params to optionally override default ones (Snowflake db and schema name). kwargs (dict): Custom activity params. secrets_path (str): Secrets manager path of the flow. """ task_id = 'update_dim_tables' kwargs = kwargs or {} completed = task_status.is_completed_report(feed_name, date) if task_status.is_completed_task(feed_name, date, task_id) and completed: activity.logger.info( 'Task {task_id} of {feed_name} for {date} already complete, and ' '"reload" flag was not passed, skipping...'.format( task_id=task_id, feed_name=feed_name, date=date)) return sf_config = get_sf_config(secrets_path) sf_config_custom = merge_configs(sf_config, sfdb_params) ExecutorFA = registered_executors.get(feed_name) with ExecutorFA(sf_config_custom) as sf_executor: sns_report = {} for table_name in kwargs['tables_to_update']: activity.logger.info('Updating {}...'.format(table_name)) result = sf_executor.update_dimension_table(date, table_name) sns_report.update({table_name: result}) activity.logger.info('All the dimension tables were updated.') changed_tables = [ table_name for table_name in sns_report if sns_report[table_name].get( 'number of rows inserted', 0) + sns_report[table_name].get('number of rows updated', 0) and table_name in kwargs['include_to_report'] ] report_subject = '{feed_name} ETL for {date} update' \ ' following tables {tables}'.format( date=date, tables=', '.join(changed_tables), feed_name=feed_name) if len(report_subject) > 100: report_subject = report_subject[0:96] + '...' context = { 'sns_report_subject': report_subject, 'sns_report_message': '\n'.join( '{table_name}: {inserted} rows inserted, ' '{updated} rows updated'.format( table_name=table_name, inserted=sns_report[table_name].get( 'number of rows inserted', '0'), updated=sns_report[table_name].get( 'number of rows updated', '0')) for table_name in changed_tables) } task_status.mark_completed_task(feed_name, date, task_id) activity.logger.info( 'Task {task_id} of {feed_name} for {date} completed'.format( task_id=task_id, feed_name=feed_name, date=date)) return context if context.get('sns_report_message') else {}