Airflow中Postgres Hook返回结果时datetime对象无__module__属性报错求助
解决Airflow XCom序列化datetime字段报错问题
问题场景
使用PostgresHook查询包含timestamp类型字段的表时,任务返回结果通过XCom传递触发以下报错:
AttributeError: 'datetime.datetime' object has no attribute '__module__'
测试命令为:airflow tasks test etl extract,仅查询非timestamp字段时运行正常。
解决方案
方案一:转换查询结果中的datetime字段为字符串
直接修改数据提取函数,将返回结果中的datetime对象转换为ISO格式字符串,避免XCom序列化失败:
from airflow.hooks.postgres_hook import PostgresHook from datetime import datetime from airflow.decorators import dag, task from airflow.utils.dates import days_ago default_args = { 'owner': 'airflow', } def get_data(): sql_stmt = "SELECT * FROM table" pg_hook = PostgresHook( postgres_conn_id='postgres_connection', schema='postgres' ) pg_conn = pg_hook.get_conn() cursor = pg_conn.cursor() cursor.execute(sql_stmt) results = cursor.fetchall() # 遍历转换datetime字段为ISO字符串 converted_results = [] for row in results: converted_row = [] for item in row: if isinstance(item, datetime): converted_row.append(item.isoformat()) else: converted_row.append(item) converted_results.append(tuple(converted_row)) return converted_results @dag(default_args=default_args, schedule_interval=None, start_date=days_ago(2), tags=['etl']) def etl(): @task() def extract(): return get_data() data = extract() etl_dag = etl()
方案二:自定义XCom编码器处理datetime类型
通过自定义XCom编码器,让Airflow能够正确序列化datetime对象,无需修改原始数据类型:
步骤1:定义自定义编码器
在DAG文件中添加自定义编码器类,继承默认的XComEncoder并重写处理逻辑:
from airflow.utils.json import XComEncoder from datetime import datetime import json class CustomXComEncoder(XComEncoder): def default(self, o): if isinstance(o, datetime): return o.isoformat() # 调用父类处理其他类型 return super().default(o)
步骤2:在任务中手动处理XCom序列化
修改extract任务,手动使用自定义编码器序列化数据后推送XCom:
@task() def extract(): data = get_data() from airflow.models.xcom import XCom from airflow.operators.python import get_current_context context = get_current_context() ti = context['ti'] # 使用自定义编码器序列化数据 serialized_data = json.dumps(data, cls=CustomXComEncoder).encode("UTF-8") ti.xcom_push(key='return_value', value=serialized_data)
可选:全局配置自定义编码器
如果希望所有DAG都生效,可修改Airflow配置文件airflow.cfg,添加以下配置:
xcom_encoder = your_dag_module.CustomXComEncoder
(将your_dag_module替换为实际包含自定义编码器的Python模块路径)
方案对比
- 方案一:实现简单,无需调整Airflow配置,适合临时需求;缺点是丢失
datetime对象类型,后续任务如需使用日期操作需手动转换字符串回datetime。 - 方案二:保留原始数据类型信息(或可自定义序列化规则),全局生效后一劳永逸;缺点是需要额外代码或配置。
内容的提问来源于stack exchange,提问作者SelketDaly
相关产品推荐
相关产品推荐

