import os import traceback import sys from garcon import activity from garcon import runner from garcon import param from ingest_training import tasks from ingest_training import config class IngestTrainingFlow(object): def __init__(self, domain=None, version=None): """Create a WorkFlow flow. Args: domain (str): domain workflow runs under. version (str): version of the workflow (ex 1.0). """ # Name of the SWF WorkFlow self.name = config.flow_name self.domain = domain or config.env self.version = version or config.version self.create = activity.create( self.domain, self.name, version=self.version, on_exception=self.on_exception) def decider(self, schedule, context): """Flow decider. Arg: schedule (callable): Call the scheduler. """ bootstrap = schedule( 'bootstrap', self.bootstrap_activity) # if no files in drop bucket, exit if bootstrap.result.get('bootstrap.stop'): return archive_raw = schedule( 'archive_raw', self.archive_raw_activity, requires=[bootstrap]) process_excel = schedule( 'process_excel', self.process_excel_activity, requires=[archive_raw]) upload_stg = schedule( 'upload_stg', self.upload_stg_activity, requires=[process_excel]) evaluate_new_data = schedule( 'evaluate_new_data', self.evaluate_new_data_activity, requires=[upload_stg]) # update_probs = schedule( # 'update_probs', # self.update_probs_activity, # requires=[evaluate_new_data ]) refit_model = schedule( 'refit_model', self.refit_model_activity, requires=[evaluate_new_data]) predict_catalog = schedule( 'predict_catalog', self.predict_catalog_activity, requires=[refit_model]) def on_exception(self, actor, exception): """Capture an exception that has occurred in the application. Without this method exceptions would silently fail in Garcon flows. For Worker see: https://github.com/xethorn/garcon/blob/bee6bd5d5afbf2d77581d235ee1d4daa88301f42/garcon/activity.py#L274-L275 For Decider seer https://github.com/xethorn/garcon/blob/967db94c2c7f759b98e9a58f276628ea773d9b84/garcon/decider.py#L138-L139 https://github.com/xethorn/garcon/blob/967db94c2c7f759b98e9a58f276628ea773d9b84/garcon/decider.py#L178-L179 Args: actor (ActivityWorker, DeciderWorker): the actor that has received the exception. exception (Exception): the exception to capture. """ traceback.print_exc() @property def bootstrap_activity(self): return self.create( name='bootstrap', tasks=runner.Sync( tasks.bootstrap.fill( namespace='bootstrap'))) @property def archive_raw_activity(self): return self.create( name='archive_raw', tasks=runner.Sync( tasks.archive_raw.fill( namespace='archive_raw', task_name=param.StaticParam('archive_raw'), s3_archive='bootstrap.s3_archive', s3_filedrop='bootstrap.s3_filedrop'))) @property def process_excel_activity(self): return self.create( name='process_excel', tasks=runner.Sync( tasks.process_excel.fill( namespace='process_excel', task_name=param.StaticParam('process_excel'), s3_archive='bootstrap.s3_archive', s3_processed='bootstrap.s3_processed', run_date='bootstrap.run_date'))) @property def upload_stg_activity(self): return self.create( name='upload_stg', tasks=runner.Sync( tasks.upload_stg.fill( namespace='upload_stg', task_name=param.StaticParam('upload_stg'), run_date='bootstrap.run_date', s3_processed='bootstrap.s3_processed'))) @property def evaluate_new_data_activity(self): return self.create( name='evaluate_new_data', tasks=runner.Sync( tasks.evaluate_new_data.fill( namespace='evaluate_new_data', latest_model='bootstrap.latest_model', s3_machine_v_human='bootstrap.s3_machine_v_human'))) @property def refit_model_activity(self): return self.create( name='refit_model', tasks=runner.Sync( tasks.refit_model.fill( namespace='evaluate_new_data', run_date='bootstrap.run_date', latest_model='bootstrap.latest_model', next_model='bootstrap.next_model', model_archive='bootstrap.model_archive'))) @property def predict_catalog_activity(self): return self.create( name='predict_catalog', tasks=runner.Sync( tasks.predict_catalog.fill( namespace='predict_catalog', next_model='bootstrap.next_model', s3_predictions='bootstrap.s3_predictions')))