ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Apache Airflow `@deadline_reference` 装饰器修复:支持无括号用法,杜绝自定义 Deadline 引用注册静默失败

Apache Airflow `@deadline_reference` 装饰器修复:支持无括号用法,杜绝自定义 Deadline 引用注册静默失败 Apache Airflowdeadline_reference装饰器修复支持无括号用法杜绝自定义 Deadline 引用注册静默失败【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读本文围绕 Apache Airflow 的一则 bugfix 展开对应变更说明 70708.bugfix.rst自定义 Deadline 引用Deadline Reference的注册装饰器deadline_reference现在可以不带括号直接使用。此前无括号写法会静默跳过注册并把被装饰的类重新绑定为装饰器内部函数最终在运行期以一条与装饰器无关的TypeError形式暴露出来。读完本文你将理解该缺陷的根源、修复后的装饰器分发逻辑以及如何正确编写、注册和反序列化自定义 Deadline 引用并能在 DAG 中安全使用两种带/不带括号写法。一、变更说明原文与问题现象1.1 变更说明内容仓库中的 70708.bugfix.rst 对本次修复的描述如下Thedeadline_referencedecorator can now be used without parentheses. Previously, using it that way silently skipped registration and rebound the decorated class to the decorators inner function, which surfaced later as an unrelatedTypeError.翻译为中文其含义包含三个关键事实修复目标deadline_reference允许不带括号使用即裸装饰器语法deadline_reference旧行为缺陷无括号使用时注册被静默跳过不报错、不注册缺陷后果被装饰的类被重新绑定为装饰器的内部函数后续以类的身份使用它时会抛出与装饰器本身无关的TypeError排查成本极高。1.2 什么是 deadline_referencedeadline_reference是 Airflow 中用于注册自定义 Deadline 引用类的装饰器属于 Deadline Alerts截止时间告警Airflow 3.1 引入的实验特性能力的一部分。它负责把用户自定义的、继承自BaseDeadlineReference的引用类挂载到DeadlineReference命名空间下形如DeadlineReference.ClassName并决定该引用在何时被求值DAG run 创建时或排队时。其完整使用背景可参考官方指南 Deadline Alerts。二、缺陷根源装饰器对裸用法的分发逻辑缺失要理解这个 bug需要先看装饰器的实现。deadline_reference定义在 task-sdk/src/airflow/sdk/definitions/deadline.py 中它需要同时支持三种调用形态deadline_reference # 裸用法无括号 deadline_reference() # 空括号 deadline_reference(DeadlineReference.TYPES.DAGRUN_QUEUED) # 带参数Python 装饰器语法决定了当使用裸用法时被装饰的类会作为第一个位置参数直接传给装饰器函数本身而当使用带括号用法时传入装饰器函数的是一个配置参数或None返回值才是一个真正的装饰器函数。修复后的实现见 deadline.py通过typing.overload声明了两种签名并在运行时用isinstance(deadline_reference_type, type)精确区分两种形态overload def deadline_reference( deadline_reference_type: type[BaseDeadlineReference], ) - type[BaseDeadlineReference]: ... overload def deadline_reference( deadline_reference_type: DeadlineReferenceTypes | None None, ) - Callable[[type[BaseDeadlineReference]], type[BaseDeadlineReference]]: ... def deadline_reference(deadline_reference_typeNone): # 裸用法无括号传入的就是被装饰的类本身 if isinstance(deadline_reference_type, type): return DeadlineReference.register_custom_reference(deadline_reference_type) # 带括号用法返回真正的装饰器 def decorator( reference_class: type[BaseDeadlineReference], ) - type[BaseDeadlineReference]: DeadlineReference.register_custom_reference(reference_class, deadline_reference_type) return reference_class return decorator从修复后的代码可以推断旧版本的缺陷机制旧的实现没有裸用法分支无论传入的是类还是配置参数都会直接返回内部的decorator函数。于是deadline_reference这种裸用法执行时传入的类被当作配置参数接收但它并不是合法的DeadlineReference.TYPES选项函数没有走注册逻辑注册被静默跳过因为旧逻辑对未知参数既不报错也不注册函数返回的是内部decorator函数而装饰器语法会把返回值重新绑定到原类名上——即MyReference这个名称指向的不再是类而是一个普通函数之后无论是DeadlineAlert(referenceDeadlineReference.MyReference)这样的实例化调用还是序列化过程中对类的属性访问都会以一条与装饰器毫无关联的TypeError失败。这正是变更说明中silently skipped registration and rebound the decorated class to the decorators inner function的完整含义错误被推迟、且错误信息具有误导性。三、修复后的正确行为与三种合法用法修复的核心是第 393 行的isinstance(deadline_reference_type, type)分支当检测到第一个参数本身就是类对象时立即执行注册并把类本身作为装饰器结果返回register_custom_reference在 deadline.py 末尾return reference_class从而保证类名始终指向真正的类。3.1 裸用法deadline_referencedeadline_reference class MyBareReference(BaseDeadlineReference): # 等价于 deadline_reference()默认在 DAG run 创建时求值 def _evaluate_with(self, *, session: Session, **kwargs) - datetime: return some_datetime3.2 空括号用法deadline_reference()deadline_reference() class MyCustomReference(BaseDeadlineReference): # 默认情况下 evaluate_with 在 DAG run 创建时被调用 def _evaluate_with(self, *, session: Session, **kwargs) - datetime: # 在这里编写业务逻辑对 Core 类型使用延迟导入 from airflow.models import DagRun return some_datetime def serialize_reference(self) - dict: return {reference_type: self.reference_name}3.3 带参数用法指定求值时机deadline_reference(DeadlineReference.TYPES.DAGRUN_QUEUED) class MyQueuedRef(BaseDeadlineReference): # 在 DAG run 排队时求值 def _evaluate_with(self, *, session: Session, **kwargs) - datetime: return some_datetime def serialize_reference(self) - dict: return {reference_type: self.reference_name}三种用法最终都汇入同一个注册入口DeadlineReference.register_custom_reference(reference_class, deadline_reference_type)其内部依次完成默认时机未指定类型时回退到DeadlineReference.TYPES.DAGRUN_CREATED基类校验类必须继承BaseDeadlineReference兼容 Core 侧的ReferenceModels.BaseDeadlineReference否则抛出ValueError: xxx must inherit from BaseDeadlineReference无参构造校验注册时会对类执行一次reference_class()实例化若类需要必填构造参数会抛出带指引信息的TypeError提示使用dataclass并给字段默认值挂载命名空间setattr(cls, reference_class.__name__, reference_instance)使 DAG 作者可通过DeadlineReference.ClassName访问登记求值时机根据deadline_reference_type把类追加到TYPES.DAGRUN_CREATED或TYPES.DAGRUN_QUEUED元组并刷新合并的TYPES.DAGRUN。四、测试验证修复行为有据可查仓库中的单元测试直接覆盖了本次修复的行为位于 airflow-core/tests/unit/models/test_deadline.pydef test_deadline_reference_decorator_without_parentheses(self): deadline_reference class BareDecoratedRef(BaseDeadlineReference): def _evaluate_with(self, *, session: Session, **kwargs) - datetime: return timezone.datetime(DEFAULT_DATE) # 修复后的关键断言名字必须仍是类而不是内部装饰器函数 assert isinstance(BareDecoratedRef, type) assert issubclass(BareDecoratedRef, BaseDeadlineReference) assert hasattr(DeadlineReference, BareDecoratedRef.__name__) assert getattr(DeadlineReference, BareDecoratedRef.__name__).__class__ is BareDecoratedRef assert_correct_timing(BareDecoratedRef, DeadlineReference.TYPES.DAGRUN_CREATED) assert_builtin_types_unchanged( DeadlineReference.TYPES.DAGRUN_QUEUED, DeadlineReference.TYPES.DAGRUN_CREATED )该测试的三条断言与 bugfix 的语义一一对应isinstance(BareDecoratedRef, type)直接验证类没有被重绑为内部函数这一旧缺陷已被消除getattr(DeadlineReference, ...).__class__ is BareDecoratedRef验证注册确实发生类已挂载到DeadlineReference命名空间assert_correct_timing(...)验证裸用法默认登记为DAGRUN_CREATED时机且内置引用类型不被破坏。此外同文件还提供了配套的边界用例test_deadline_reference_decorator_calls_register_method断言带参用法会且仅会调用一次register_custom_reference(DecoratedCustomRef, timing)test_deadline_reference_decorator_without_parentheses_invalid_class裸用法下非法基类同样会被拒绝ValueErrortest_deadline_reference_requiring_arguments_raises_helpful_error无参构造失败的类会得到可读的错误提示。五、注册之外的另一半插件注册与反序列化需要特别强调的是装饰器注册 ≠ 插件注册。register_custom_reference的 docstring 明确警告见 deadline.py装饰器只影响解析 DAG 文件的那个进程让类以DeadlineReference.ClassName形式可用而调度器反序列化 DAG 时能否重新解析该类取决于它是否被列在某个AirflowPlugin的deadline_references属性中。这一插件注册要求由另一则 significant 变更引入见 66737.significant.rst自定义 Deadline 引用必须像自定义 Timetable、自定义 Partition Mapper 一样通过AirflowPlugin.deadline_references列表注册未注册的引用在反序列化时会抛出DeadlineReferenceNotRegistered。5.1 插件注册的收集逻辑调度器进程通过 plugins_manager.py 的get_deadline_references_plugins()收集所有插件声明的引用类并以类的限定名qualname为键建立查找表cache def get_deadline_references_plugins() - dict[str, type[DeadlineReferenceType]]: Collect and get deadline reference classes registered by plugins. return { qualname(deadline_ref_cls): deadline_ref_cls for plugin in _get_plugins()[0] for deadline_ref_cls in plugin.deadline_references }5.2 反序列化时的解析与错误路径反序列化侧由 serialization/helpers.py 的find_registered_custom_deadline_reference()负责按__class_path查找注册类未命中时抛出DeadlineReferenceNotRegistered提示语会明确告知必须通过AirflowPlugin的deadline_references属性注册。序列化包装类SerializedCustomReference见 serialization/definitions/deadline.py在deserialize_reference中依次处理三类情形缺少__class_path提示存储的引用损坏、来自更新版本或插件未安装类未注册抛出DeadlineReferenceNotRegistered注册成功动态委托给包装的内部引用执行_evaluate_with求值逻辑并校验required_kwargs声明。上述查表解析 未注册报错的行为有独立测试覆盖见 airflow-core/tests/unit/serialization/test_deadline_reference_registry.py其中test_serialized_custom_reference_uses_registry、test_serialized_custom_reference_rejects_unregistered分别验证了注册命中与未注册拒绝两条路径。六、完整的自定义引用示例推荐写法综合以上机制一个生产可用的自定义 Deadline 引用应同时完成装饰器注册与插件注册。官方指南 Deadline Alerts 给出的完整模式如下文件置于插件目录如$AIRFLOW_HOME/plugins/deadline_references.pyfrom sqlalchemy.orm import Session from airflow.plugins_manager import AirflowPlugin from airflow.sdk import BaseDeadlineReference, DeadlineReference, deadline_reference from airflow.sdk.timezone import datetime # 默认在 DAG run 创建时求值等价于 deadline_reference() deadline_reference() class MyCustomDecoratedReference(BaseDeadlineReference): A custom reference evaluated when Dag runs are created. def _evaluate_with(self, *, session: Session, **kwargs) - datetime: # 在这里编写业务逻辑 return your_datetime # 指定在 DAG run 排队时求值并声明需要的 DAG run 上下文 deadline_reference(DeadlineReference.TYPES.DAGRUN_QUEUED) class MyQueuedReference(BaseDeadlineReference): A custom reference evaluated when Dag runs are queued. required_kwargs {dag_id, run_id} def _evaluate_with(self, *, session: Session, **kwargs) - datetime: dag_id kwargs[dag_id] run_id kwargs[run_id] return your_datetime # 关键注册到插件调度器反序列化时才能解析 class MyDeadlineReferencePlugin(AirflowPlugin): name my_deadline_reference_plugin deadline_references [MyCustomDecoratedReference, MyQueuedReference]然后在 DAG 文件中按内置引用的方式使用with DAG( dag_idcustom_reference_example, deadlineDeadlineAlert( referenceDeadlineReference.MyCustomDecoratedReference, intervaltimedelta(hours2), callbackAsyncCallback(my_callback), ), ): ...七、最佳实践与注意事项结合修复本身与官方指南的Important Notes见 deadline-alerts.rst给出如下实操建议裸用法与带括号用法语义一致deadline_reference与deadline_reference()完全等价均默认在DAGRUN_CREATED时机求值。但注意装饰器必须有括号才是可配置形式——需要指定DAGRUN_QUEUED等时机时必须带参数。时区感知_evaluate_with必须返回时区感知timezone-aware的 datetime 对象。无参构造约束自定义引用在注册时即被实例化因此必须可无参构造若需要构造参数应装饰dataclass并给每个字段默认值这也与 test_deadline_reference_requiring_arguments_raises_helpful_error 验证的错误提示一致。插件注册不可或缺仅装饰不注册会导致调度器反序列化时抛出DeadlineReferenceNotRegistered。装饰器与插件注册是DAG 文件解析与调度器反序列化两个阶段各自的需求缺一不可。required_kwargs仅支持dag_id与run_id声明其他上下文键会在求值时抛出ValueError引用自身的配置应通过构造字段或读取 Airflow Variable 完成。重启生效新增或修改自定义引用后需要重启 Airflow API Server异步回调在 Triggerer 中执行变更后同样需要重启 Triggerer 以重新加载文件。升级注意如果此前因旧缺陷被迫写成deadline_reference()形式升级到包含本次修复的版本后可保持原写法不变完全兼容而旧的裸写法若代码中恰有从静默失效变为正确注册属于行为修复而非破坏性变更。八、相关文件索引变更说明70708.bugfix.rst装饰器与注册实现task-sdk/src/airflow/sdk/definitions/deadline.py插件注册收集airflow-core/src/airflow/plugins_manager.py反序列化解析与异常airflow-core/src/airflow/serialization/helpers.py序列化包装类airflow-core/src/airflow/serialization/definitions/deadline.py装饰器行为测试airflow-core/tests/unit/models/test_deadline.py注册表解析测试airflow-core/tests/unit/serialization/test_deadline_reference_registry.py功能使用指南airflow-core/docs/howto/deadline-alerts.rst插件注册要求变更说明66737.significant.rst【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表