Airflow任务间传递列表时元素被加双引号的问题求助
问题背景
从以下JSON结构生成嵌套列表:
op = {'specific': { 'run': { 'name': 'jon', 'location': 'NY' } } }
生成列表的代码(原代码存在语法错误):
lst = [] jsn = op['specific'] for task in jsn: for attrs in jsn[task]: # 原代码误用[]索引,应为append()方法调用 lst.append([[task, attrs], jsn[task][attrs]]) return lst
该列表在Airflow中从一个Operator传递到另一个Operator后,接收端的格式异常变为:
("[[['run','name'],'jon'],", "[[['run','location'],'NY'],")
原Operator输出的列表结构正常,但传递后元素被添加双引号,整体转为元组。
调试步骤
验证源数据结构
在第一个Operator的execute方法中添加日志,打印生成列表的类型和原始内容,确认输出是否符合预期:import logging logger = logging.getLogger(__name__) # 生成lst后添加 logger.info(f"列表类型: {type(lst)}") logger.info(f"列表内容: {repr(lst)}")同时查看Airflow UI的XCom面板,检查该任务输出的XCom数据,确认序列化前的结构是否正常。
检查XCom序列化配置
Airflow默认用JSON序列化XCom数据,若数据包含非JSON兼容类型会出现异常:- 查看
airflow.cfg中的core.xcom_serializer配置,默认值为json; - 若使用自定义Operator,确认是否正确设置了
do_xcom_push=True,且返回的是标准Python可序列化对象。
- 查看
排查接收端获取逻辑
在第二个Operator中打印获取到的XCom数据的类型和内容:def execute(self, context): ti = context['ti'] received_data = ti.xcom_pull(task_ids='上游任务ID') logger.info(f"接收数据类型: {type(received_data)}") logger.info(f"接收数据内容: {repr(received_data)}")确认是序列化环节异常,还是接收时的解析错误。
解决办法
修复代码语法错误
原代码中lst.append[xxx]是语法错误,append是列表方法,必须用括号调用:lst.append([[task, attrs], jsn[task][attrs]])错误的调用方式可能导致生成非预期的数据结构,甚至触发异常,先修复此问题。
显式控制序列化/反序列化
手动用JSON序列化和反序列化来规避Airflow自动处理的异常:- 上游Operator返回时序列化:
import json return json.dumps(lst) - 下游Operator接收时反序列化:
import json received_data = ti.xcom_pull(task_ids='上游任务ID') lst = json.loads(received_data)
- 上游Operator返回时序列化:
调整XCom序列化器
若数据包含JSON不兼容的类型(此场景暂不涉及),可在可信环境下切换为pickle序列化器:
修改airflow.cfg:[core] xcom_serializer = pickle重启Airflow服务后测试数据传递。
确认Operator类型兼容性
若使用的是BashOperator等非PythonOperator,其返回值是命令输出的字符串,需确保上游用PythonOperator返回标准Python对象,而非字符串形式的列表。
内容的提问来源于stack exchange,提问作者sammy morgan

