Airflow自定义算子移至独立文件后运行异常求助
检查Python路径与模块导入
确认my_custom_operator.py所在目录在Airflow的PYTHONPATH里。Astronomer环境下默认dags目录会被加入路径,但如果是子目录得额外配置。可以先在DAG文件开头加临时路径测试:import sys from pathlib import Path sys.path.append(str(Path(__file__).parent))要是测试有效,要么改
airflow.cfg的python_path参数,要么在Astronomer的Dockerfile里设ENV PYTHONPATH=/usr/local/airflow/dags来永久配置。验证自定义算子的依赖与兼容性
确保CustomTriggerDagRunOperator继承的TriggerDagRunOperator在Airflow 2.1.3里的导入路径没错,原算子得是from airflow.operators.trigger_dagrun import TriggerDagRunOperator,自定义算子里别写错。同时检查有没有漏加依赖导入,比如from airflow.utils.decorators import apply_defaults这类Airflow必需的装饰器,迁移时很容易漏。排查Kubernetes执行环境的文件同步问题
在Astronomer的K8s部署里,worker节点得同步dags目录下的文件。先确认my_custom_operator.py已经同步到worker的dags目录:- 进worker pod:
kubectl exec -it <worker-pod-name> -- bash - 查文件是否存在:
ls /usr/local/airflow/dags/my_custom_operator.py
要是文件不存在,得确认Astronomer的DAG同步机制(比如git-sync或volume挂载)配置对不对,保证新添的文件能同步到所有worker节点。
- 进worker pod:
查看Airflow核心日志定位错误
既然任务没日志输出,就得看Airflow scheduler和worker的核心日志:- 看scheduler日志:
kubectl logs <scheduler-pod-name> - 看worker日志:
kubectl logs <worker-pod-name>
重点搜和my_custom_operator相关的报错,比如导入错误、类初始化错误之类的,这些错误可能导致任务启动阶段直接失败,没法生成任务日志。
- 看scheduler日志:
测试自定义算子的最小可复现案例
整个极简DAG只放自定义算子,排除其他任务干扰:from airflow import DAG from datetime import datetime from my_custom_operator import CustomTriggerDagRunOperator with DAG( 'test_custom_operator', start_date=datetime(2023, 1, 1), schedule_interval=None ) as dag: trigger_task = CustomTriggerDagRunOperator( task_id='test_trigger', trigger_dag_id='target_dag_id' )运行这个测试DAG,要是还失败,就把排查范围缩小到自定义算子本身的代码;要是成功,那原DAG里肯定有其他依赖冲突。
内容的提问来源于stack exchange,提问作者Vanshaj Bhatia

