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

Airflow中DatabricksSqlOperator报datetime无法JSON序列化错误求助

问题:Airflow DatabricksSqlOperator 读取含日期/时间戳列的Delta表时XCom序列化失败

我尝试在Airflow中使用DatabricksSqlOperator从Databricks Delta表获取数据,代码如下:

select = DatabricksSqlOperator(
    databricks_conn_id=databricks_id,
    http_path=http_path,
    task_id="select_data",
    sql="select * from schema.table_name",
)

连接已正常建立,但获取数据时出现错误:

[2023-05-08, 18:32:43 UTC] {xcom.py:599} ERROR - Could not serialize the XCom value into JSON. If you are using pickle instead of JSON for XCom, then you need to enable pickle support for XCom in your *** config.
[2023-05-08, 12:44:02 UTC] {taskinstance.py:1851} ERROR - Task failed with exception
Traceback (most recent call last):
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/utils/session.py", line 72, in wrapper
    return func(*args, **kwargs)
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 2378, in xcom_push
    XCom.set(
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/utils/session.py", line 72, in wrapper
    return func(*args, **kwargs)
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/xcom.py", line 206, in set
    value = cls.serialize_value(
  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/xcom.py", line 597, in serialize_value
    return json.dumps(value).encode('UTF-8')
  File "/usr/local/lib/python3.8/json/__init__.py", line 231, in dumps
    return _default_encoder.encode(obj)
  File "/usr/local/lib/python3.8/json/encoder.py", line 199, in encode
    chunks = self.iterencode(o, _one_shot=True)
  File "/usr/local/lib/python3.8/json/encoder.py", line 257, in iterencode
    return _iterencode(o, 0)
  File "/usr/local/lib/python3.8/json/encoder.py", line 179, in default
    raise TypeError(f'Object of type {o.__class__.__name__} '
TypeError: Object of type datetime is not JSON serializable

该表包含1个date格式列和2个timestamp格式列,在无日期列的表上测试时可正常运行。


解决方案

错误原因

DatabricksSqlOperator默认会将查询结果推送到XCom,但查询返回的datetime/timestamp类型无法被Airflow默认的JSON序列化器处理,导致报错。

解决方法

1. 禁用XCom推送(推荐,若无需传递结果)

如果不需要把查询结果传递给后续任务,直接关闭XCom推送即可:

select = DatabricksSqlOperator(
    databricks_conn_id=databricks_id,
    http_path=http_path,
    task_id="select_data",
    sql="select * from schema.table_name",
    do_xcom_push=False  # 添加此参数禁用XCom推送
)

2. 将日期/时间戳转换为字符串

修改SQL查询,把date和timestamp列转换为字符串格式,确保返回结果都是可JSON序列化的类型:

select 
  cast(date_column as string) as date_column,
  cast(timestamp_column1 as string) as timestamp_column1,
  cast(timestamp_column2 as string) as timestamp_column2,
  -- 其他列保持不变
  col1, col2
from schema.table_name

3. 启用Pickle序列化XCom(不推荐生产环境)

如果必须传递原始datetime类型数据,可以修改Airflow配置,允许使用Pickle序列化XCom。在airflow.cfg中设置:

[core]
enable_xcom_pickling = True

注意:Pickle存在安全风险,若XCom数据来自不可信来源,可能导致代码注入,不建议生产环境使用。


内容的提问来源于stack exchange,提问作者vaibhav

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 09:15:31