如何在两个DockerOperator间使用Airflow XCom传递数据?
Airflow DockerOperator XCom数据传递故障排查与修复
问题背景
已花费3天解决Airflow中DockerOperator间的XCom数据传递问题,当前XCom值已成功存储,但other/main.py无法获取目标值——要么返回Retrieved XCom value: None,要么提示XCom.get_one需要run_id但容器内无法获取该参数。
DAG目录结构
./dags/ ├── custom_operator.py ├── docker-xcom │ ├── other │ │ ├── Dockerfile │ │ ├── main.py │ │ └── requirements.txt │ └── scraper │ ├── Dockerfile │ ├── main.py │ └── requirements.txt ├── jobs.py
现有代码细节
1. scraper/main.py(任务1容器内代码)
my_dict = {"name":"Luna","age":31} print(my_dict)
2. 自定义CustomDockerOperator(custom_operator.py)
from airflow.providers.docker.operators.docker import DockerOperator class CustomDockerOperator(DockerOperator): def __init__(self, *args, **kwargs): xcom_pull = kwargs.pop('xcom_pull', None) super().__init__(*args, **kwargs) self.xcom_pull = xcom_pull def execute(self, context): if self.xcom_pull: ti = context['ti'] value = ti.xcom_pull(key=self.xcom_pull) print(f'Retrieved XCom value: {value}') return super().execute(context)
3. DAG配置(jobs.py)
from airflow import DAG from airflow.operators.empty import EmptyOperator from custom_operator import CustomDockerOperator from datetime import datetime, timedelta # 原代码缺失该导入,需补充 default_args = { 'owner' : 'airflow', 'description' : 'docker xcom test', 'depend_on_past' : False, 'start_date' : datetime(2021, 5, 1), 'email_on_failure' : False, 'email_on_retry' : False, 'retries' : 1, 'retry_delay' : timedelta(minutes=5) } def push_xcom_value(**kwargs): # 原代码存在语法错误,修正后如下 return_value = kwargs['ti'].xcom_pull(task_ids='task1') kwargs["ti"].xcom_push(value=return_value) with DAG( dag_id="XCOM_DOCKER", default_args=default_args, description="Docker + XCOM", start_date=datetime(2022, 12, 12), schedule_interval="30 * * * *", catchup=False ) as dag: start_dag = EmptyOperator(task_id='start_dag') end_dag = EmptyOperator(task_id='end_dag') task1 = CustomDockerOperator( task_id='task1', image='xcom_scraper:latest', container_name='xcom_scrape', do_xcom_push=True, api_version='auto', auto_remove=True, docker_url="unix://var/run/docker.sock", network_mode="bridge", on_success_callback=push_xcom_value, ) task2 = CustomDockerOperator( task_id="task2", image='xcom_other:latest', container_name='xcom_other', api_version='auto', auto_remove=True, xcom_pull=task1.task_id, docker_url="unix://var/run/docker.sock", network_mode="bridge", ) start_dag >> task1 >> task2 >> end_dag
核心问题分析
- XCom拉取逻辑错误:当前
CustomDockerOperator中ti.xcom_pull(key=self.xcom_pull)是按key拉取XCom,但task1的XCom默认key为return_value,而非task_id,导致拉取空值。 - 容器无法直接访问Airflow元数据:
other/main.py直接调用XCom API会失败,因为容器默认无Airflow元数据库访问权限,且缺少run_id等上下文参数。 - 冗余回调与语法错误:
push_xcom_value函数存在语法错误,且完全冗余——do_xcom_push=True已自动将容器stdout输出推为XCom。
修复方案
1. 修正CustomDockerOperator,传递XCom到容器环境变量
修改Operator逻辑,拉取正确的XCom值并通过环境变量注入容器:
from airflow.providers.docker.operators.docker import DockerOperator class CustomDockerOperator(DockerOperator): def __init__(self, *args, **kwargs): self.xcom_pull_task_ids = kwargs.pop('xcom_pull_task_ids', None) super().__init__(*args, **kwargs) def execute(self, context): if self.xcom_pull_task_ids: ti = context['ti'] # 拉取指定task_id的默认XCom值(key为return_value) xcom_value = ti.xcom_pull(task_ids=self.xcom_pull_task_ids, key='return_value') # 将XCom值转为字符串,注入容器环境变量 self.environment = self.environment or {} self.environment['XCOM_VALUE'] = str(xcom_value) print(f'Passed XCom value to container: {xcom_value}') return super().execute(context)
2. 更新task2的DAG配置
使用修正后的Operator,指定拉取task1的XCom:
task2 = CustomDockerOperator( task_id="task2", image='xcom_other:latest', container_name='xcom_other', api_version='auto', auto_remove=True, xcom_pull_task_ids='task1', # 改为指定task_id docker_url="unix://var/run/docker.sock", network_mode="bridge", )
3. 修改other/main.py,从环境变量读取XCom值
容器内无需访问Airflow API,直接读取环境变量并解析:
import os import ast # 读取环境变量并转为字典 xcom_str = os.getenv('XCOM_VALUE') if xcom_str: my_dict = ast.literal_eval(xcom_str) print(f'Received XCom value: {my_dict}') else: print('No XCom value received')
4. 移除冗余回调
删除push_xcom_value函数及task1中的on_success_callback配置,保留do_xcom_push=True即可自动推送XCom。
额外注意事项
- 若XCom为复杂JSON结构,建议使用
json.dumps和json.loads替代ast.literal_eval,避免解析错误。 - 确保jobs.py中补充
datetime和timedelta的导入,否则DAG无法正常运行。
内容的提问来源于stack exchange,提问作者Roku
相关产品推荐
相关产品推荐

