You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 15:47:30