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

如何正确使用@task.kubernetes()装饰器导入第三方依赖与自定义模块?

问题描述

我有一个要在K8s Pod中运行的Airflow DAG任务,示例代码如下:

from mymodule import process_data # 自定义模块导入
from decouple import AutoConfig # 外部安装的依赖
from airflow.decorators import dag, task
from datetime import datetime

@dag(
    "model_trainer",
    start_date=datetime(2023, 1, 1),
    catchup=False,
    schedule=None,
)
def pipeline():
    @task.kubernetes(image="python")
    def fetch_data():
        return process_data()

运行pipeline().test()时出现错误:NameError: name 'process_data' is not defined。原因是使用的python镜像中不存在自定义模块process_data,也未安装外部依赖。请问该场景的最佳处理方式是什么?是否需要创建包含依赖的自定义镜像并在函数内导入模块?

解决方案

最佳实践:构建自定义镜像

这是生产环境下最可靠、易维护的方案,具体步骤如下:

  1. 编写Dockerfile
    基于官方Python镜像,安装所需依赖并复制自定义模块:
# 选择匹配你DAG环境的Python版本
FROM python:3.10-slim

# 设置工作目录
WORKDIR /app

# 复制依赖清单(提前整理好requirements.txt,包含python-decouple等依赖)
COPY requirements.txt .

# 安装依赖,--no-cache-dir减少镜像体积
RUN pip install --no-cache-dir -r requirements.txt

# 复制自定义模块到容器的Python可识别路径下
COPY mymodule/ /app/mymodule/
  1. 构建并推送镜像
    将镜像推送到你的容器仓库(私有仓库或公共仓库均可):
docker build -t your-registry/model-trainer:v1 .
docker push your-registry/model-trainer:v1
  1. 更新DAG代码
    在@task.kubernetes中指定自定义镜像:
from mymodule import process_data
from decouple import AutoConfig
from airflow.decorators import dag, task
from datetime import datetime

@dag(
    "model_trainer",
    start_date=datetime(2023, 1, 1),
    catchup=False,
    schedule=None,
)
def pipeline():
    @task.kubernetes(image="your-registry/model-trainer:v1")
    def fetch_data():
        return process_data()

关于模块导入的说明

不需要把导入语句放到fetch_data函数内部。只要自定义镜像中包含了mymodule和依赖包,DAG顶部的导入语句就能在K8s Pod执行任务时正常解析——Pod会使用镜像内置的Python环境运行代码。

临时调试替代方案

如果只是临时测试,不想构建镜像,可以通过挂载ConfigMap传递自定义模块,同时在任务启动命令中安装依赖,但这种方式不适合生产环境:

from airflow.decorators import dag, task
from airflow.kubernetes.volume import Volume
from airflow.kubernetes.volume_mount import VolumeMount
from datetime import datetime

@dag(
    "model_trainer",
    start_date=datetime(2023, 1, 1),
    catchup=False,
    schedule=None,
)
def pipeline():
    # 先将mymodule打包为ConfigMap(kubectl create configmap mymodule-cm --from-file=mymodule/)
    volume = Volume(
        name="mymodule-volume",
        configs=[{
            "name": "mymodule-cm",
            "mountPath": "/app/mymodule"
        }]
    )
    volume_mount = VolumeMount(
        "mymodule-volume",
        mount_path="/app/mymodule",
        sub_path=None,
        read_only=True
    )

    @task.kubernetes(
        image="python",
        volume_mounts=[volume_mount],
        volumes=[volume],
        env_vars={"PYTHONPATH": "/app"},
        cmds=["pip", "install", "python-decouple", "&&", "python"]
    )
    def fetch_data():
        from mymodule import process_data
        return process_data()

该方案需要手动维护ConfigMap,且每次任务运行都要重新安装依赖,效率低,仅适合短期调试。

内容的提问来源于stack exchange,提问作者Kaushil Kundalia

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 22:50:28