"""Sales data ETL flow.""" from garcon_contrib.aws import garcon_sns from flows import runner from flows.flow import DatabaseParam from flows.flow import FlowBase from flows.sales_data import tasks class Flow(FlowBase): """Sales data ETL flow class.""" def decider(self, schedule): """Schedule the next appropriate activity for the SWF execution.""" bootstrap = schedule( 'bootstrap', self.bootstrap_activity) if bootstrap.result.get('bootstrap.stop'): return unload_sales_data = schedule( 'unload_sales_data', self.unload_sales_data_activity, requires=[bootstrap]) load_raw_table = schedule( 'load_raw_table', self.load_raw_table_activity, requires=[unload_sales_data]) insert_daily_revenue = schedule( 'insert_daily_revenue', self.insert_daily_revenue_activity, requires=[load_raw_table]) send_notification = schedule( 'send_notification', self.send_notification, requires=[insert_daily_revenue]) schedule('set_status', self.set_status, requires=[send_notification]) @property def bootstrap_activity(self): """Bootstrap all required params. Returns: garcon.activity.Activity instance. """ return self.create( name='bootstrap', retry=4, tasks=runner.Sync( tasks.bootstrap.fill( namespace='bootstrap', correlation_id='correlation_id', accounting_period_id='accounting_period_id'))) @property def unload_sales_data_activity(self): """Unload sales data to S3. Returns: garcon.activity.Activity instance. """ return self.create( name='unload_sales_data', retry=4, tasks=runner.Sync( tasks.unload_sales_data.fill( namespace='unload_sales_data', correlation_id='correlation_id', upcs=DatabaseParam('upcs'), accounting_period_id='bootstrap.accounting_period_id'))) @property def load_raw_table_activity(self): """Load raw table with S3 data. Returns: garcon.activity.Activity instance. """ return self.create( name='load_raw_table', retry=4, tasks=runner.Sync( tasks.create_temp_table.fill( correlation_id='correlation_id', namespace='create_temp_table'), tasks.load_temp_table.fill( correlation_id='correlation_id', batch_count='unload_sales_data.batch_count', temp_table_name='create_temp_table.table_name'), tasks.load_raw_from_temp.fill( correlation_id='correlation_id', temp_table_name='create_temp_table.table_name', upcs=DatabaseParam('upcs'), accounting_period_id='bootstrap.accounting_period_id'))) @property def insert_daily_revenue_activity(self): """Insert daily sales data into accounting_revenue table. Returns: garcon.activity.Activity instance. """ return self.create( name='insert_daily_revenue', retry=4, tasks=runner.Sync( tasks.insert_daily_revenue.fill( correlation_id='correlation_id', namespace='insert_daily_revenue', upcs=DatabaseParam('upcs'), accounting_period_id='bootstrap.accounting_period_id'))) @property def send_notification(self): """Tell the world the ETL has finished. Returns: Activity: garcon.activity.Activity instance. """ return self.create( name='send_notification', retry=4, tasks=runner.Sync( tasks.prepare_sns_message.fill( namespace='sns_message', correlation_id='correlation_id', upcs='upcs', accounting_period_id='bootstrap.accounting_period_id'), garcon_sns.sns_publish_message.fill( topic='sns_message.topic', message='sns_message.message', subject='sns_message.subject'))) @property def set_status(self): """Set DynamoDB statuses and mark log to completed. Returns: Activity: garcon.activity.Activity instance. """ return self.create( name='set_status', retry=4, tasks=runner.Sync( tasks.set_dynamo_status.fill( correlation_id='correlation_id', date_start='bootstrap.date_start'), tasks.set_final_status.fill( correlation_id='correlation_id')))