资讯动态

Apache Airflow common-messaging Provider 版本演进全解:从 1.0.0 到 2.1.0 的 Changelog 深度解读

发布时间:2026/9/14 10:09:26 来源:尧图企业网站定制
Apache Airflow common-messaging Provider 版本演进全解从 1.0.0 到 2.1.0 的 Changelog 深度解读【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本篇以 Apache Airflow 仓库中apache-airflow-providers-common-messaging的官方 Changelog 文档为主体逐版本解读该 Provider 从 1.0.0 到 2.1.0 的全部变更并结合仓库源码深入剖析 2.0.0 的队列接口重构Breaking Change与 2.1.0 新增的 Trigger 队列路由参数。读完本文你将能够看懂该 Provider 的每次版本升级意味着什么、如何安全地从旧版queue参数迁移到新版scheme参数以及 Provider 发现机制在底层是如何工作的。一、文档定位与包概况本文解读的源文档是 changelog.rst它是common.messagingProvider 的官方版本变更记录文档开头给贡献者的说明指出Changelog 由发布经理半自动维护只有在存在破坏性变更breaking changes时才需要在 Changelog 标题下方手动补充解释说明告知用户应如何应对。从包元数据看该 Provider 当前的完整事实如下以仓库实际内容为准项目内容依据包名apache-airflow-providers-common-messagingprovider.yaml当前版本2.1.0生命周期状态productionprovider.yaml、pyproject.toml最低 Airflow 版本3.0.11.0.0 起仅支持 Airflow 3.0pyproject.toml、changelog.rstPython 版本3.10classifiers 覆盖 3.10–3.14pyproject.toml可选依赖extrasamazon→apache-airflow-providers-amazon9.7.0apache.kafka→apache-airflow-providers-apache-kafka1.9.0pyproject.toml构建后端flit_core4.0.2pyproject.toml[pyproject.toml](https://link.gitcode.com/i/98ee9158b824f7c6c981339cb41bc0f3#L61-L66)中的注释还特别解释了为什么最低版本锁定在 3.0.1 而非 3.0.0Airflow 3.0.1 之前的 Provider 管理器缺少对common.messagingProvider 的发现支持新 Provider 与 Airflow 3.0.0 不兼容。二、Changelog 全版本变更总览以下为 changelog.rst 中全部 13 个版本的完整内容整理。带 Below changes are excluded from the changelog 注释的条目是被自动流程过滤掉的内部变更本文在总览表中标注不展开。2.1.0当前版本FeaturesAdd trigger queue support for AsyncCallback and BaseEventTrigger (#71346)—— 为 Trigger 增加对 Triggerer 命名队列--queues的显式路由支持详见第四节的源码剖析。被排除项Adopt flit 4 构建后端#71186、准备 2026-07-22 发布#70256、为每个 Provider 文档补充可选 extras 说明#69478、修复生成文档与 pyproject.toml 不一致#68991。2.0.4MiscAdd explicit [tool.flit.sdist] sections to flit-based pyproject.tomls (#65861)。这一点可以在本包的 pyproject.toml 中直接验证文件末尾显式声明了[tool.flit.sdist]的include清单docs/、provider.yaml、各级__init__.py、tests/注释说明这是为了让构建不依赖 VCS 信息flit 4.0 默认--no-use-vcs。被排除项Fix stale system test documentation links#65071。2.0.3MiscAdd Python 3.14 Support (#63520)。对应 pyproject.toml 中出现的Programming Language :: Python :: 3.14classifier。被排除项*.iml加入 .gitignore#63636、准备 2026-03-09 发布#63198、准备下一轮文档#62495、按 AIP-95 为所有 provider.yaml 增加lifecycle字段#62190。2.0.2MiscNew year means updated Copyright notices (#60344)。被排除项Providers wave 2025-12-30#59947、准备 2025-12-09 发布#59249、Provider 发布流程更新#58316。2.0.1MiscConvert all airflow distributions to be compliant with ASF requirements (#58138)即所有发行版包按 Apache 软件基金会合规要求改造LICENSE/NOTICE 结构、ASF 元数据等。被排除项较多其中与用户相关的包括在文档中新增如何在使用中访问消息 payload章节#55438、删除不必要的 LICENSE 文件#58191等。2.0.0重大版本Breaking changesRefactor Common Queue Interface (#54651)—— 重构通用队列接口这是该 Provider 唯一一次标注 Breaking Change 的版本详见第三节。Bug Fixesfix(messaging): improve MessageQueueTrigger logging and add comprehensive tests (#54492)—— 改进MessageQueueTrigger的日志输出并补充测试。Doc-onlyMake term Dag consistent in providers docs (#55101)—— 文档术语统一为 Dag。被排除项Switch pre-commit to prek#54258、修复 README/index 中的 Airflow 2 引用#55240。1.0.5MiscAdd Python 3.13 support for Airflow (#46891)、Remove type ignore across codebase after mypy upgrade (#53243)、Remove upper-binding for python-requires (#52980)、Temporarily switch to use , pattern instead of ~ (#52967)。后两条值得注意前者移除了requires-python的上限约束当前文件即为3.10见 pyproject.toml后者说明依赖约束写法从 PEP 440 兼容操作符~临时切换到,区间写法这也是为什么现在看到的依赖是apache-airflow3.0.1这种形式。1.0.4MiscDrop support for Python 3.9 (#52072)—— 放弃 Python 3.9最低支持版本提升到 3.10。1.0.3Bug FixesMove MESSAGE_QUEUE_PROVIDERS array to where it belongs - to msq_queue (#51774)—— 将消息队列 Provider 注册表数组移到它本应所在的模块。从当前源码结构看这一职责最终落在 msg_queue.py 的MESSAGE_QUEUE_PROVIDERS构建语句上它基于ProvidersManager的发现结果动态实例化各队列 Provider。1.0.2MiscAIP-82: Add KafkaMessageQueueProvider (#49938)—— 按 AIP-82 规范引入 Kafka 消息队列 Provider 接入这也是为什么本包的可选依赖中包含apache.kafkaextra要求apache-airflow-providers-apache-kafka1.9.0。1.0.1MiscMove SQS message queue to Amazon provider (#50057)—— SQS 队列实现从 core 迁移到 Amazon Provider 中。被排除项中还包括将 SQS 队列代码示例从 core 移入本 Provider 文档#49208以及修正版本号到 1.0.1#50099。1.0.0InitialInitial version of the provider (#46694)。原始文档在此处附有一条 note该 Provider 的首个版本仅面向 Airflow 3.0原因是社区维护 Provider 对 Airflow 最低支持版本的策略文档中引用了 PROVIDERS.rst 的相关章节当前仓库对应文件为 PROVIDERS.rst。三、2.0.0 Breaking Change 深度剖析队列接口重构2.0.0 的Refactor Common Queue Interface (#54651)是整份 Changelog 中唯一一次 Breaking Change。结合当前源码重构后的最终形态可以完整还原这次接口变化的技术含义。3.1 抽象基类BaseMessageQueueProvider所有队列接入方都继承 base_provider.py 中的BaseMessageQueueProvider。它定义了四个契约方法scheme_matches(scheme)判断给定 scheme 字符串是否匹配该 Provider。基类默认实现为self.scheme scheme的精确匹配docstring 明确要求各 Provider 的匹配条件不得互相重叠以避免冲突queue_matches(queue)抽象方法判断给定 queue URI 是否匹配该 Provider 的模式trigger_class()抽象方法返回queue_matches命中时实际使用的 Trigger 类trigger_kwargs(queue, **kwargs)抽象方法返回构造上述 Trigger 实例所需的参数字典。基类的 docstring 直接给出了扩展指引要新增一个受 common-messaging 支持的 Provider请创建继承本基类的新类并将其加入MESSAGE_QUEUE_PROVIDERS在动态发现机制下实际是通过各 Provider 包的provider.yaml声明队列类并由 Provider 管理器发现见下文。3.2 核心触发器MessageQueueTrigger 的双匹配模式triggers/msg_queue.py 中的MessageQueueTrigger是统一入口它抽象掉了 Provider 细节让用户无需关心底层是 SQS 还是 Kafka 就能监控一个队列。重构后的关键行为对应 2.0.0 的接口变化与 #54492 的日志改进参数校验与弃用处理L68-L97queue与scheme二选一二者皆空则抛出ValueError传queue会触发AirflowProviderDeprecationWarning明确提示该参数将在未来版本移除请改用scheme参数并将配置作为关键字参数传入。缓存的 Provider 匹配triggercached_propertyL107-L162未安装任何队列 Provider 时记录 error 日志并抛错提示安装缺失的依赖按queueURI 匹配向后兼容路径或按scheme匹配新路径匹配到 0 个或多个 Provider 都会记录包含可用 Provider 列表的 error 日志并抛ValueError——这正是 #54492 improve logging 的落点报错信息中会列出所有已注册的 Provider 类名便于快速定位是哪个依赖没装命中唯一 Provider 后实例化其trigger_class()并把匹配过程产生的trigger_kwargs与之合并。透传执行run()只是把事件流从被包装的真实 Trigger 转发出来serialize()同样委托给底层 Trigger因此MessageQueueTrigger本质上是一层路由与适配。官方文档 triggers.rst 对重构后的用法给出的权威说明是必填参数只有scheme例如kafka、redispubsub、sqs连接信息使用对应的默认连接 ID如 SQS 用aws_default。系统测试示例 example_message_queue_trigger.py 展示了标准写法from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger from airflow.providers.standard.operators.empty import EmptyOperator from airflow.sdk import DAG, Asset, AssetWatcher # 监听外部消息队列此处以 AWS SQS 为例 trigger MessageQueueTrigger(schemesqs, sqs_queuehttps://sqs.us-east-1.amazonaws.com/0123456789/Test) # 定义一个监听该队列消息的 Asset asset Asset(sqs_queue_asset, watchers[AssetWatcher(namesqs_watcher, triggertrigger)]) with DAG(dag_idexample_msgq_watcher, schedule[asset]) as dag: EmptyOperator(task_idtask)triggers.rst 还说明了事件驱动的执行模型Asset抽象外部实体、AssetWatcher将命名 Trigger 绑定到 Asset 上DAG 不再按固定周期运行而是在队列收到新消息、Asset 更新时触发任务内可通过triggering_asset_events参数读取事件消息 payload 存放在event.extra[payload]中。3.3 升级迁移建议从 1.x 迁移到 2.x 的用户对照 Changelog 与源码可得出明确的操作路径停止传queueURI改为传schemescheme名 各 Provider 要求的关键字参数若触发AirflowProviderDeprecationWarning按提示改造即可。两种模式在源码中仍并存queue路径会调用trigger_kwargs解析 URIscheme路径直接把**self.kwargs透传给底层 Trigger因此旧写法在 2.x 依然可运行只是被标记为待移除。四、2.1.0 新特性剖析Trigger 的 queue 路由参数2.1.0 的唯一 Feature#71346为AsyncCallback与BaseEventTrigger增加了queue参数支持。在MessageQueueTrigger中这一能力体现为独立于队列 URI 的trigger_queue构造参数msg_queue.py L55-L58用途把 Trigger 分配到 Triggerer 的某个命名队列上对应 Airflow 核心配置triggerer__queues_enabled与airflow triggerer --queues选项docstring 中直接引用了这两个入口命名考量docstring 特意解释了为什么参数叫trigger_queue而不是queue——queue这个名字已被弃用的 broker 队列 URI 参数占用改名是为了不破坏现有调用方。从源码实现看L68-L105trigger_queue被存到私有属性_trigger_queue并通过自定义的queueproperty/setter 对外暴露。注释说明了这样做的原因绕过BaseEventTrigger.__init__(queue...)的签名使该行为不受所安装 airflow-core 版本影响保证向前兼容。五、Provider 发现机制与 1.x 系列变更的底层印证Changelog 中 1.0.1/1.0.2/1.0.3 的一系列变更SQS 迁入 Amazon Provider、按 AIP-82 增加 Kafka Provider、MESSAGE_QUEUE_PROVIDERS归位反映的其实是同一套底层机制的成熟过程common-messaging 本身不硬编码任何具体队列实现而是通过动态发现聚合各 Provider 包声明的队列类。当前源码印证了这一点msg_queue.py L30-L40模块加载时调用ProvidersManager().initialize_providers_queues()再通过create_class_by_name按全限定类名反射实例化providers_manager.queue_class_names中的每个队列 Provider组装出MESSAGE_QUEUE_PROVIDERS列表发现入口位于 Airflow 核心的 providers_manager.pyinitialize_providers_queues标注了provider_info_cache(queues)装饰器先initialize_providers_list()再_discover_queues()即队列 Provider 来自已安装各 Provider 包的元数据声明本包自己的声明见 provider.yaml其中triggers段落注册了airflow.providers.common.messaging.triggers.msg_queue模块这也解释了 2.0.3 版本 Changelog 中 AIP-95provider.yaml 增加lifecycle字段与本包 provider.yaml 中lifecycle: production字段的由来。由此可以推断用户侧新增一种消息队列支持不需要修改 common-messaging 本身而是由对应 Provider 包如 amazon、apache.kafka声明队列类后自动被MESSAGE_QUEUE_PROVIDERS拾取——1.0.2 引入 Kafka、1.0.1 迁入 SQS 都是这一机制下的产物这也正是 pyproject.toml 中两个 optional extra 与两类队列支持的对应关系。六、兼容性与构建体系演进1.0.4–2.0.4Changelog 的 Misc 条目虽然琐碎串起来恰好是一条完整的兼容性与工程化演进线Python 支持区间1.0.4 放弃 Python 3.9#52072→ 1.0.5 增加 Python 3.13#46891→ 2.0.3 增加 Python 3.14#63520。最终形态为 pyproject.toml 中requires-python 3.10且 classifiers 列出 3.10–3.14。依赖约束写法1.0.5 的两条 Misc#52980、#52967解释了当前文件中apache-airflow3.0.1这类仅下限、无上限约束的成因provider.yaml的 versions 列表 上方注释也写明这些版本由发布经理维护不要手动修改除非同源码树中的其他 Provider 要求更高的 Airflow 下限。ASF 合规2.0.1#58138将全部发行版改造为符合 ASF 要求当前包根目录保留有独立的 LICENSE 与 NOTICE与 2.0.1 被排除项 #58191删除各子目录多余 LICENSE相吻合。构建后端2.0.4 显式化[tool.flit.sdist]#658612.1.0 被排除项中记录了正式切换到 flit 4#71186当前 pyproject.toml 即flit_core4.0.2L128-L137 的 sdist 清单注释直接引用了 flit 4.0 的--no-use-vcs默认行为变化作为该段声明的动机。七、如何阅读与使用这份 Changelog对使用者升级前先看 Breaking changes 段落。本包历史上唯一一次出现在 2.0.0队列接口重构应对方式已在第三节给出将queueURI改为schemescheme 关键字参数。被排除的变更行Below changes are excluded...注释块不要删除——文档明确要求保留这些行需要时可手动上移为正式条目。这类条目多数是发布流程性提交版本准备、依赖对齐对用户运行时行为无直接影响但其中的依赖版本变化如 2.0.1 期间的 extras 调整值得关注。结合 provider.yaml 的versions列表与 index.rst 中的安装说明pip install apache-airflow-providers-common-messaging核对当前可用版本与 Airflow 最低版本要求3.0.1。对贡献者按文档头部的 NOTE TO CONTRIBUTORS只有破坏性变更需要手动在 Changelog 标题下方添加用户侧应对说明其余变更由发布经理半自动流程生成条目——这解释了为什么每个版本块末尾都有相同的排除区注释模板。总结common-messaging的 Changelog 呈现出一条清晰的演进主线1.0.0 作为 Airflow 3 专属的首发版本#466941.x 系列完成了消息队列生态的重组SQS 归位 Amazon Provider、Kafka 按 AIP-82 接入、注册表归位2.0.0 以一次 Breaking Change 确立了scheme 优先的队列接口#54651并同步改进了MessageQueueTrigger的错误日志与测试覆盖#544922.1.0 则补齐了 Trigger 在 Triggerer 命名队列上的路由能力#71346。整条线索与仓库中 msg_queue.py、base_provider.py、providers_manager.py 的源码实现一一对应构成了一份文档说变更、源码给证据的完整版本图谱。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价