"""Tasks of the Playlists sharing with Chartmetric Workflow.""" from datetime import datetime, timedelta from garcon import task from feed_ingestion.flows.chartmetric_playlists_share import config from feed_ingestion.flows.chartmetric_playlists_share.snowflake_executor \ import SnowflakeExecutor from feed_ingestion.flows.helpers import get_sf_config @task.decorate(timeout=600) def bootstrap(activity, date): """Bootstrap workflow by getting the correct configurations. Args: activity (ActivityWorker): The activity worker. date (str): The date to ingest from (YYYY-MM-DD). Returns: dict: Context. """ # date is the date passed in or yesterday's date if date: date_obj = datetime.strptime(date, '%Y-%m-%d') else: date_obj = datetime.utcnow().date() - timedelta(days=1) return dict( feed_name=config.feed_name, date=date_obj.strftime('%Y-%m-%d'), ) @task.decorate(timeout=7200) def load_missing_playlists(activity, feed_name, platform_name): """Get missing playlists for platform. Args: activity (ActivityWorker): The activity worker. feed_name (str): The name of the feed. platform_name (str): The name of the platform (e.g. instagram). Returns: dict: Context. """ query_name = 'load_missing_playlists_for_' + platform_name activity.logger.info('Starting ingestion for {}'.format(query_name)) sf_config = get_sf_config(feed_name) params = { 'temp_table_name': '{}_{}'.format(feed_name, platform_name)} with SnowflakeExecutor(sf_config) as sf_executor: sf_executor.execute_query(query_name, **params) activity.logger.info('Finished ingesting {}'.format(query_name)) @task.decorate(timeout=1800) def clear_shared_table(activity, feed_name): """Clear all records from shared table. Args: activity (ActivityWorker): The activity worker. Returns: dict: Context. """ activity.logger.info('Starting clearing shared table') sf_config = get_sf_config(feed_name) params = { 'shared_table_name': config.shared_table_name} with SnowflakeExecutor(sf_config) as sf_executor: sf_executor.execute_query('clear_shared_table', **params) activity.logger.info('Finished clearing shared table') @task.decorate(timeout=7200) def share_missing_playlists(activity, feed_name, platform_name): """Share missing playlists for platform. Args: activity (ActivityWorker): The activity worker. feed_name (str): The name of the feed. platform_name (str): The name of the platform (e.g. instagram). Returns: dict: Context. """ query_name = 'share_missing_playlists_for_' + platform_name activity.logger.info('Starting ingestion for {}'.format(query_name)) sf_config = get_sf_config(feed_name) params = { 'temp_table_name': '{}_{}'.format(feed_name, platform_name), 'shared_table_name': config.shared_table_name} with SnowflakeExecutor(sf_config) as sf_executor: sf_executor.execute_query(query_name, **params) activity.logger.info('Finished ingesting {}'.format(query_name))