如何正确使用@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,也未安装外部依赖。请问该场景的最佳处理方式是什么?是否需要创建包含依赖的自定义镜像并在函数内导入模块?
解决方案
最佳实践:构建自定义镜像
这是生产环境下最可靠、易维护的方案,具体步骤如下:
- 编写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/
- 构建并推送镜像
将镜像推送到你的容器仓库(私有仓库或公共仓库均可):
docker build -t your-registry/model-trainer:v1 . docker push your-registry/model-trainer:v1
- 更新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
相关产品推荐
相关产品推荐

