"""Transform data from staging_raw table and load into marketshare table.""" 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.decorate(timeout=7200) def load_marketshare_data( activity, feed_name, date, sfdb_params, secrets_path=None, kwargs=None): """Load data on main_market_share. 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. """ task_id = 'load_marketshare_data' 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) ExecutorMS = registered_executors.get(feed_name) with ExecutorMS(sf_config_custom) as sf_executor: activity.logger.info( 'Deleting rows for {date} from main_market_share for feed ' '{feed_name}'.format(date=date, feed_name=feed_name)) sf_executor.delete_from_marketshare_table(date, **kwargs) activity.logger.info( 'Loading main_market_share for feed {}'.format(feed_name)) sf_executor.load_marketshare_data(date, **kwargs) activity.logger.info( 'Data was loaded into main_market_share')