Airflow升级至2.5.0后无法向XCom推送数据问题排查
Airflow 2.5.0推送bytes类型到XCom报错的解决方法
问题原因
Airflow 2.5.0对XCom的序列化逻辑进行了调整:2.3.4版本允许直接序列化bytes类型,但升级后默认的XComEncoder在处理非JSON原生类型时,会尝试读取对象的__module__属性做序列化标记,而bytes作为Python内置基础类型没有该属性,因此触发AttributeError。
解决方法
方法1:将bytes转换为Base64字符串(推荐)
这是最安全通用的方案,把bytes编码为JSON支持的字符串格式,拉取时再解码回bytes:
推送XCom修改代码:
import base64 from airflow.operators.python import get_current_context context = get_current_context() ti = context['ti'] # 将bytes转为Base64编码的字符串 encoded_doc = base64.b64encode(doc).decode('utf-8') ti.xcom_push(key="file", value=encoded_doc)
拉取XCom代码:
import base64 from airflow.operators.python import get_current_context context = get_current_context() ti = context['ti'] encoded_doc = ti.xcom_pull(key="file") # 将Base64字符串解码回bytes doc = base64.b64decode(encoded_doc)
方法2:自定义XCom编码器
若需在多任务中批量处理bytes类型,可自定义编码器专门处理该类型:
import base64 from airflow.utils.json import XComEncoder class BytesAwareXComEncoder(XComEncoder): def default(self, o): if isinstance(o, bytes): # 用自定义标记+Base64编码存储bytes return {"__type__": "bytes", "data": base64.b64encode(o).decode('utf-8')} return super().default(o)
在DAG默认参数中指定该编码器:
default_args = { 'xcom_encoder': BytesAwareXComEncoder, } with DAG( dag_id="your_dag_id", default_args=default_args, # 其他DAG参数 ) as dag: # 任务定义
注:拉取XCom时需对应解码,也可额外实现自定义解码器。
方法3:使用Pickle序列化(不推荐生产环境)
Airflow提供基于Pickle的XCom后端,支持直接序列化bytes,但Pickle存在安全风险(可能执行恶意代码),仅在完全信任所有任务的场景下使用:
全局配置(airflow.cfg):
xcom_backend = airflow.models.xcom.PickleXComBackend
DAG级别配置:
default_args = { 'xcom_backend': 'airflow.models.xcom.PickleXComBackend', }
内容的提问来源于stack exchange,提问作者Rishabhg
相关产品推荐
相关产品推荐

