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
相关产品推荐
相关产品推荐

