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

如何在两个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

核心问题分析

  1. XCom拉取逻辑错误:当前CustomDockerOperator中ti.xcom_pull(key=self.xcom_pull)是按key拉取XCom,但task1的XCom默认key为return_value,而非task_id,导致拉取空值。
  2. 容器无法直接访问Airflow元数据:other/main.py直接调用XCom API会失败,因为容器默认无Airflow元数据库访问权限,且缺少run_id等上下文参数。
  3. 冗余回调与语法错误: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 08:25:20