Airflow参数化Timetable报错:__init__缺hour和minute参数
参数化Airflow Timetable报错解决方案
问题现象
实现参数化Timetable并在DAG中传入hour和minute参数后,出现如下错误:
TypeError: __init__() missing 2 required positional arguments: 'hour' and 'minute'
非参数化版本运行正常,参数化后失效。
相关代码
Timetable代码
class EveryFiscalPeriod(Timetable): def __init__(self, hour: int, minute: int) -> None: self._hour = hour self._minute = minute def next_dagrun_info( self, *, last_automated_data_interval: Optional[DataInterval], restriction: TimeRestriction, ) -> Optional[DagRunInfo]: delta = timedelta(days=28) if last_automated_data_interval is not None: # There was a previous run on the regular schedule. next_start = last_automated_data_interval.end next_end = last_automated_data_interval.end + delta else: # This is the first ever run on the regular schedule. restriction_earliest = restriction.earliest next_start = restriction_earliest - delta if next_start is None: # No start_date. Don't schedule. return None next_end = restriction_earliest return DagRunInfo( data_interval=DataInterval(start=next_start, end=next_end), run_after=DateTime.combine(next_end.date(), Time(self.hour), Time(self.minute)).replace(tzinfo=UTC), )
DAG代码片段
with DAG( catchup=False, # 其他配置项 max_active_runs=1, schedule=EveryFiscalPeriod(hour=15, minute=30), ) as dag: # 任务定义
堆栈跟踪信息
[2024-02-03T06:29:28.294+0000] {app.py:1744} ERROR - Exception on /dags/accrual_repot_missing_orders/grid [GET] Traceback (most recent call last): File "/home/airflow/.local/lib/python3.8/site-packages/flask/app.py", line 2529, in wsgi_app response = self.full_dispatch_request() File "/home/airflow/.local/lib/python3.8/site-packages/flask/app.py", line 1825, in full_dispatch_request rv = self.handle_user_exception(e) File "/home/airflow/.local/lib/python3.8/site-packages/flask/app.py", line 1823, in full_dispatch_request rv = self.dispatch_request() File "/home/airflow/.local/lib/python3.8/site-packages/flask/app.py", line 1799, in dispatch_request return self.ensure_sync(self.view_functions[rule.endpoint])(**view_args) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/www/auth.py", line 53, in decorated return func(*args, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/www/decorators.py", line 168, in view_func return f(*args, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/www/decorators.py", line 127, in wrapper return f(*args, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/utils/session.py", line 79, in wrapper return func(*args, session=session, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/www/views.py", line 2936, in grid dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/utils/session.py", line 76, in wrapper return func(*args, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/dagbag.py", line 189, in get_dag self._add_dag_from_db(dag_id=dag_id, session=session) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/dagbag.py", line 271, in _add_dag_from_db dag = row.dag File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/serialized_dag.py", line 221, in dag return SerializedDAG.from_dict(data) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/serialization/serialized_objects.py", line 1413, in from_dict return cls.deserialize_dag(serialized_obj["dag"]) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/serialization/serialized_objects.py", line 1341, in deserialize_dag v = _decode_timetable(v) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/serialization/serialized_objects.py", line 211, in _decode_timetable return timetable_class.deserialize(var[Encoding.VAR]) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/timetables/base.py", line 168, in deserialize return cls() TypeError: __init__() missing 2 required positional arguments: 'hour' and 'minute'
错误原因
从堆栈信息可以看出,Airflow在反序列化Timetable实例时,调用了默认的deserialize方法(直接返回cls()),但你的Timetable类需要hour和minute参数才能实例化,因此报错。另外代码还存在两个小问题:
next_dagrun_info中访问属性时用了self.hour和self.minute,但初始化时定义的是self._hour和self._minute,会引发属性错误。DateTime.combine方法参数错误,该方法只接受一个date和一个time对象,不能传两个Time实例。
解决方案
需要为参数化Timetable实现序列化、反序列化方法,同时修正代码中的属性访问和DateTime.combine的错误:
修改后的Timetable代码:
from airflow.timetables.base import Timetable, DataInterval, DagRunInfo, TimeRestriction from datetime import timedelta, DateTime, Time, UTC from typing import Optional, Dict, Any class EveryFiscalPeriod(Timetable): def __init__(self, hour: int, minute: int) -> None: self._hour = hour self._minute = minute def next_dagrun_info( self, *, last_automated_data_interval: Optional[DataInterval], restriction: TimeRestriction, ) -> Optional[DagRunInfo]: delta = timedelta(days=28) if last_automated_data_interval is not None: next_start = last_automated_data_interval.end next_end = last_automated_data_interval.end + delta else: restriction_earliest = restriction.earliest if restriction_earliest is None: return None next_start = restriction_earliest - delta next_end = restriction_earliest # 修正DateTime.combine的用法,以及属性访问 run_time = Time(self._hour, self._minute) run_after = DateTime.combine(next_end.date(), run_time).replace(tzinfo=UTC) return DagRunInfo( data_interval=DataInterval(start=next_start, end=next_end), run_after=run_after, ) # 实现序列化方法,保存参数 def serialize(self) -> Dict[str, Any]: return {"hour": self._hour, "minute": self._minute} # 实现反序列化方法,从保存的参数中实例化 @classmethod def deserialize(cls, data: Dict[str, Any]) -> "EveryFiscalPeriod": return cls(hour=data["hour"], minute=data["minute"]) # 实现__eq__和__hash__,确保Airflow能正确比较实例 def __eq__(self, other: Any) -> bool: if not isinstance(other, EveryFiscalPeriod): return False return self._hour == other._hour and self._minute == other._minute def __hash__(self) -> int: return hash((self._hour, self._minute))
说明
serialize方法将Timetable的参数转为字典,用于Airflow序列化存储。deserialize类方法从序列化的字典中恢复Timetable实例,传入必要的参数。__eq__和__hash__确保Airflow在比较DAG配置时能正确识别相同的Timetable实例,避免不必要的DAG更新。- 修正了属性访问错误和
DateTime.combine的参数错误,确保运行时间计算正确。
内容的提问来源于stack exchange,提问作者Luis Lema
相关产品推荐
相关产品推荐

