如何在Airflow PostgresHook游标中给PostgreSQL查询传参并解决序列化报错
错误原因
- 核心是返回结果序列化失败:你定义的Python函数最终返回了查询结果
result,如果结果集中包含datetime类型的字段(比如你的created_at字段),Airflow默认会将任务返回值存入XCom,而XCom默认的JSON序列化方案不支持Python原生datetime对象,直接触发该报错。 - 次要问题是SQL参数拼接不规范:你用
str.format()直接拼接SQL字符串,既存在SQL注入风险,也可能因为日期格式不匹配引发额外的SQL执行错误。
修复方案
第一步:修改SQL文件的占位符
PostgreSQL的Python驱动psycopg2支持参数化查询,把pg_query.sql中的占位符改成%s,不需要额外加单引号:
SELECT * FROM public.airtimes WHERE created_at > %s;
第二步:调整Python代码
有两种可选方案解决序列化报错,根据你的使用场景选即可:
方案A:不需要将查询结果传递给下游任务
直接关闭当前任务的XCom推送,不需要对结果做额外处理,修改任务定义即可:
# 原有函数不用大幅修改,只需要在定义PythonOperator的时候加do_xcom_push=False read_pg_task = PythonOperator( task_id='read_latest_data_from_pg', python_callable=read_latest_data_from_pg, do_xcom_push=False, dag=dag )
同时把原函数里的cursor.execute(pg_export_data_query.format(pg_date))改成参数化调用,避免注入风险:
# 原来的cursor.execute行替换成下面的 cursor.execute(pg_export_data_query, (pg_date,))
方案B:需要将查询结果传递给下游任务
先把结果中的datetime类型转成可序列化的字符串格式再返回:
from datetime import datetime def read_latest_data_from_pg(**kwargs): with open('dags/scripts/sql_scripts/pg_query.sql','r') as sqlfile: pg_export_data_query = sqlfile.read() pg_date = '2021-05-01' pg_hook = PostgresHook(postgres_conn_id='pg_conn', delegate_to=None, use_legacy_sql=False) # 用hook自带的get_records方法,省略手动处理连接、游标的冗余代码 result = pg_hook.get_records(pg_export_data_query, parameters=(pg_date,)) # 序列化datetime字段 processed_result = [] for row in result: new_row = [] for col in row: if isinstance(col, datetime): new_row.append(col.isoformat()) else: new_row.append(col) processed_result.append(new_row) print('result', processed_result) return processed_result
内容的提问来源于stack exchange,提问作者Shadow Walker
相关产品推荐
相关产品推荐

