如何用Airflow存储SQL查询结果并用于条件判断?XCom返回None求助
看起来你在尝试用自定义Operator执行SQL并通过XCom传递结果时遇到了返回None的问题,我来帮你分析一下可能的原因和解决办法,同时也会告诉你如何用这些结果做条件判断。
首先分析XCom返回None的可能原因
你的自定义DWPostgresReturn Operator继承自DWPostgresOperator,虽然你在execute方法中返回了hook.get_records()的结果,但还是拿到None,大概率是这几个问题:
1. 基类可能禁用了XCom自动推送
很多自定义的Postgres Operator基类会默认设置do_xcom_push=False,这样即使你在execute中返回了结果,Airflow也不会自动把它推送到XCom里。
解决办法:在实例化t1的时候显式开启do_xcom_push=True:
t1 = DWPostgresReturn( task_id='t1', postgres_conn_id='db_conn', sql="select output from table", config=config, start_date = dt.datetime(2020,3,9), do_xcom_push=True # 加上这个参数开启自动推送 )
如果这个参数在基类里没有暴露,那你可以在自定义Operator的__init__方法中默认设置:
class DWPostgresReturn(DWPostgresOperator): def __init__(self, *args, do_xcom_push=True, **kwargs): super().__init__(*args, do_xcom_push=do_xcom_push, **kwargs) def execute(self, context): self.log.info('Executing: %s', self.sql) hook = DWPostgresHook(postgres_conn_id=self.postgres_conn_id, schema=self.database) return hook.get_records(self.sql, parameters=self.parameters)
2. 手动推送XCom更稳妥
如果自动推送还是有问题,你可以在execute方法里手动把结果推送到XCom,这样更可控:
class DWPostgresReturn(DWPostgresOperator): def execute(self, context): self.log.info('Executing: %s', self.sql) hook = DWPostgresHook(postgres_conn_id=self.postgres_conn_id, schema=self.database) records = hook.get_records(self.sql, parameters=self.parameters) # 手动推送XCom,指定key方便后续精准获取 context['ti'].xcom_push(key='query_result', value=records) return records
然后在t2中这样获取:
def get_records(**kwargs): ti = kwargs['ti'] xcom = ti.xcom_pull(task_ids='t1', key='query_result') string_to_print = 'Value in xcom is: {}'.format(xcom) print(xcom)
3. 确认SQL查询确实返回了数据
如果你的SQL查询select output from table本身没有返回任何记录,那hook.get_records()会返回空列表[],而不是None。如果拿到的是None,那肯定是XCom没有被推送成功,重点检查前面两点。
如何用查询结果做if-else条件判断
当你成功拿到XCom中的查询结果后,就可以根据结果做条件判断了。比如你的查询返回的是单列数据,示例代码如下:
def get_records(**kwargs): ti = kwargs['ti'] records = ti.xcom_pull(task_ids='t1') # 先判断是否有返回记录 if not records: print("查询没有返回任何数据") return # 取出第一条记录的第一个值(因为get_records返回的是二维列表) output_value = records[0][0] # 进行条件判断 if output_value == 'success': print("任务执行成功,进入后续流程") # 这里可以调用其他函数或者触发后续任务逻辑 elif output_value == 'failed': print("任务执行失败,触发告警") else: print("返回值不符合预期,进行异常处理")
如果你的查询返回多列数据,只需要调整索引即可,比如records[0][1]取第一条记录的第二列值。
内容的提问来源于stack exchange,提问作者Undesignated_Survivor

