"""Snowflake executor to operate backload_tasks table.""" from snowflake_connector.etl_connector import SnowflakeSQLExecutor, SQLLoader sql_loader = SQLLoader(__file__) class BackloadTasksExecutor(SnowflakeSQLExecutor): """Snowflake executor to operate backload_tasks table.""" def insert_backload_tasks(self, licensor, flow, date, context): """Insert one record to backload_tasks table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], licensor=licensor, date=date, flow=flow, context=context, ) self.execute_query(sql_loader, 'insert_backload_tasks', params) def delete_backload_tasks(self, licensor, flow, period_start, period_end): """Delete multiple records from backload_tasks table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], licensor=licensor, flow=flow, period_start=period_start, period_end=period_end, ) self.execute_query(sql_loader, 'delete_backload_tasks', params) def select_flows_to_process(self): """Get flows from tasks which has not yet been processed.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], ) return self.fetchall_query(sql_loader, 'select_flows_to_process', params) def select_backload_tasks_by_flow(self, flow, limit): """Select up to `limit` ready to process tasks by `flow`.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], flow=flow, limit=limit, ) return self.fetchall_query( sql_loader, 'select_backload_tasks_by_flow', params) def update_backload_tasks_as_processed(self, licensor, flow, date, context): """Delete multiple records from backload_tasks table.""" params = dict( db=self.sf_config['db'], schema=self.sf_config['schema'], licensor=licensor, flow=flow, date=date, context=context, ) self.execute_query(sql_loader, 'update_backload_tasks_as_processed', params)