Airflow使用dill反序列化DAG Run时遇timedelta未定义错误如何解决?
问题解决:PythonVirtualenvOperator传递dag_run触发NameError: timedelta未定义
环境信息
- Python 3.11
- airflow==2.7.3
- dill==0.3.7
问题现象
启用render_template_as_native_obj=True的DAG中,尝试将{{ dag_run }}作为参数传递给PythonVirtualenvOperator时,执行触发NameError: name 'timedelta' is not defined错误,错误根源来自pendulum时区对象反序列化失败。
相关DAG代码
import datetime from pathlib import Path import airflow from airflow import DAG from airflow.operators.python import PythonVirtualenvOperator import dill dag = DAG( dag_id='strange_pickling_error_dag', schedule_interval='0 5 * * 1', start_date=datetime.datetime(2020, 1, 1), catchup=False, render_template_as_native_obj=True, ) context = {"ts": "{{ ts }}", "dag_run": "{{ dag_run }}"} op_args = [context, Path(__file__).parent.absolute()] def make_foo(*args, **kwargs): print("---> making foo!") print("make foo(...): args") print(args) print("make foo(...): kwargs") print(kwargs) make_foo_task = PythonVirtualenvOperator( task_id='make_foo', python_callable=make_foo, use_dill=True, system_site_packages=False, op_args=op_args, requirements=[f"dill=={dill.__version__}", f"apache-airflow=={airflow.__version__}"], dag=dag)
错误日志
[2023-11-06, 18:23:21 UTC] {process_utils.py:182} INFO - Executing cmd: /tmp/venvqse65m1b/bin/python /tmp/venvqse65m1b/script.py /tmp/venvqse65m1b/script.in /tmp/venvqse65m1b/script.out /tmp/venvqse65m1b/string_args.txt /tmp/venvqse65m1b/termination.log [2023-11-06, 18:23:21 UTC] {process_utils.py:186} INFO - Output: [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - Traceback (most recent call last): [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - File "/tmp/venvqse65m1b/script.py", line 17, in <module> [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - arg_dict = dill.load(file) [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - ^^^^^^^^^^^^^^^ [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - File "/tmp/venvqse65m1b/lib/python3.11/site-packages/dill/_dill.py", line 287, in load [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - return Unpickler(file, ignore=ignore, **kwds).load() [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - File "/tmp/venvqse65m1b/lib/python3.11/site-packages/dill/_dill.py", line 442, in load [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - obj = StockUnpickler.load(self) [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - ^^^^^^^^^^^^^^^^^^^^^^^^^ [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - File "/home/felix/Projects/alfabank/fs_etl/venv_py3.11/lib/python3.11/site-packages/pendulum/tz/timezone.py", line 312, in __init__ [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - self._utcoffset = timedelta(seconds=offset) [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - ^^^^^^^^^ [2023-11-06, 18:23:21 UTC] {process_utils.py:190} INFO - NameError: name 'timedelta' is not defined [2023-11-06, 18:23:22 UTC] {taskinstance.py:1937} ERROR - Task failed with exception Traceback (most recent call last): File "/home/felix/Projects/alfabank/fs_etl/venv_py3.11/lib/python3.11/site-packages/airflow/operators/python.py", line 395, in execute return super().execute(context=serializable_context) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/felix/Projects/alfabank/fs_etl/venv_py3.11/lib/python3.11/site-packages/airflow/operators/python.py", line 192, in execute return_value = self.execute_callable() ^^^^^^^^^^^^^^^^^^^^^^^ File "/home/felix/Projects/alfabank/fs_etl/venv_py3.11/lib/python3.11/site-packages/airflow/operators/python.py", line 609, in execute_callable result = self._execute_python_callable_in_subprocess(python_path, tmp_path) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/felix/Projects/alfabank/fs_etl/venv_py3.11/lib/python3.11/site-packages/airflow/operators/python.py", line 446, in _execute_python_callable_in_subprocess execute_in_subprocess( File "/home/felix/Projects/alfabank/fs_etl/venv_py3.11/lib/python3.11/site-packages/airflow/utils/process_utils.py", line 171, in execute_in_subprocess execute_in_subprocess_with_kwargs(cmd, cwd=cwd) File "/home/felix/Projects/alfabank/fs_etl/venv_py3.11/lib/python3.11/site-packages/airflow/utils/process_utils.py", line 194, in execute_in_subprocess_with_kwargs raise subprocess.CalledProcessError(exit_code, cmd) subprocess.CalledProcessError: Command '['/tmp/venvqse65m1b/bin/python', '/tmp/venvqse65m1b/script.py', '/tmp/venvqse65m1b/script.in', '/tmp/venvqse65m1b/script.out', '/tmp/venvqse65m1b/string_args.txt', '/tmp/venvqse65m1b/termination.log']' returned non-zero exit status 1.
原因分析
当启用render_template_as_native_obj=True时,{{ dag_run }}会被渲染成Airflow的DagRun原生对象,该对象内部包含pendulum库的时区对象。dill在序列化这个对象后,虚拟环境中反序列化时,pendulum的timezone.py中__init__方法直接使用了timedelta但未在局部作用域导入(依赖全局导入),导致反序列化过程中找不到timedelta定义。
解决方法
方案1:只传递dag_run的必要属性
避免传递整个DagRun对象,只提取需要的属性(如run_id、execution_date等),减少序列化复杂度:
# 修改context定义 context = { "ts": "{{ ts }}", "dag_run_id": "{{ dag_run.run_id }}", "execution_date": "{{ dag_run.execution_date }}" }
方案2:在requirements中添加pendulum依赖
明确指定pendulum版本(与Airflow依赖版本一致),确保虚拟环境中pendulum能正确导入依赖:
# 修改PythonVirtualenvOperator的requirements参数 requirements=[ f"dill=={dill.__version__}", f"apache-airflow=={airflow.__version__}", "pendulum==2.1.2" # 替换为你的Airflow实际依赖的pendulum版本 ]
方案3:启用system_site_packages
允许虚拟环境继承主环境的系统包,包括pendulum及其依赖:
# 修改PythonVirtualenvOperator的system_site_packages参数 make_foo_task = PythonVirtualenvOperator( task_id='make_foo', python_callable=make_foo, use_dill=True, system_site_packages=True, # 改为True op_args=op_args, requirements=[f"dill=={dill.__version__}", f"apache-airflow=={airflow.__version__}"], dag=dag)
方案4:禁用render_template_as_native_obj,转为JSON传递
关闭原生对象渲染,将dag_run序列化为JSON字符串,在callable中再解析:
# 修改DAG的render_template_as_native_obj为False dag = DAG( dag_id='strange_pickling_error_dag', schedule_interval='0 5 * * 1', start_date=datetime.datetime(2020, 1, 1), catchup=False, render_template_as_native_obj=False, # 改为False ) # 修改context定义,使用tojson过滤器 context = {"ts": "{{ ts }}", "dag_run_json": "{{ dag_run | tojson }}"} # 修改make_foo函数解析JSON import json def make_foo(*args, **kwargs): context = args[0] dag_run_data = json.loads(context["dag_run_json"]) print("---> making foo!") print("dag_run data:", dag_run_data) print("make foo(...): args") print(args) print("make foo(...): kwargs") print(kwargs)
内容的提问来源于stack exchange,提问作者Felix
相关产品推荐
相关产品推荐

