"""snowflake_views_etl workflow.""" from garcon import runner from garcon.param import StaticParam from snowflake_views import config as flow_config from snowflake_views import tasks from snowflake_views.base_flow import BaseFlow class Flow(BaseFlow): """Class representing the workflow.""" timeout = 32400 # nine hours def __init__(self): """Initialize flow object.""" super(Flow, self).__init__( flow_name=flow_config.SWF_FLOW_NAME, version=flow_config.SWF_FLOW_VERSION) def decider(self, schedule, context): """Activity decider. Args: schedule (callable): The scheduler method. """ check_dynamo_status = schedule( 'check_dynamo_status', self.check_dynamo_status) if check_dynamo_status.result.get( 'check_dynamo_status.rebuild_view') is False: return create_view = schedule( 'create_view', self.create_view, requires=[check_dynamo_status]) schedule( 'cache_view', self.cache_view, requires=[create_view]) schedule( 'save_last_processed_timestamp', self.save_last_processed_timestamp, requires=[create_view]) def workflow_id(self, context): """Create identifier based on view being processed.""" return '{flow_name}-{view_name}'.format( flow_name=self.name, view_name=context['name']) @property def check_dynamo_status(self): """Check Dynamo status.""" return self.create( name='check_dynamo_status', tasks=runner.Sync( tasks.check_dynamo_status.fill( namespace='check_dynamo_status', view_name='name'))) @property def create_view(self): """Create view.""" return self.create( name='create_view', tasks=runner.Sync( tasks.create_view.fill( namespace='create_view', view_name='name'))) @property def cache_view(self): """Cache view.""" return self.create( name='cache_view', tasks=runner.Sync( tasks.cache_view.fill( namespace='cache_view', view_name='name', attempt_limit=StaticParam( flow_config.SF_CACHE_WARMUP_ATTEMPT_LIMIT)))) @property def save_last_processed_timestamp(self): """Save last prococessed timestamp in DynamoDB.""" return self.create( name='save_last_processed_timestamp', tasks=runner.Sync( tasks.save_last_processed_timestamp.fill( namespace='save_last_processed_timestamp', view_name='name', latest_timestamp=('check_dynamo_status.' 'current_last_processed_timestamp'))))