Apache Airflow Common Messaging Provider:用 MessageQueueTrigger 统一抽象各类消息队列触发
Apache Airflow Common Messaging Provider用 MessageQueueTrigger 统一抽象各类消息队列触发【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文围绕apache-airflow-providers-common-messaging包展开它如何以MessageQueueTrigger这一统一触发器屏蔽 Kafka、SQS、Redis Pub/Sub 等具体消息队列的差异以及如何在 DAG 中通过Asset/AssetWatcher实现“消息到达即触发”的事件驱动调度。读完本文你可以完成该 Provider 的安装与依赖配置、用scheme参数编写跨队列触发的 DAG并理解其底层的 Provider 匹配与触发器委派机制。包概览与定位Common Messaging Provider 是 Airflow 3 中面向“消息队列触发”场景的横切 Provider。其核心目标是提供一个统一抽象层不同队列 ProviderAWS SQS、Apache Kafka、Redis、IBM MQ、Azure Service Bus 等各自实现的触发器都可以通过同一个MessageQueueTrigger类来调用用户可以在不修改 DAG 代码的前提下切换底层队列实现参见 providers.rst。当前发布版本为2.1.0完整版本历史见 provider.yamlProvider 状态为ready、生命周期为production。其 Python 代码全部位于airflow.providers.common.messaging包下目录结构如下triggers/msg_queue.pyMessageQueueTrigger实现providers/base_provider.pyBaseMessageQueueProvider抽象基类tests/system/common/messaging/example_message_queue_trigger.pySQS 系统测试示例 DAG安装与依赖要求在已有 Airflow 安装之上通过 pip 安装即可pip install apache-airflow-providers-common-messaging根据文档声明的Requirements该 Provider 要求的最低 Apache Airflow 版本为3.0.1PIP packageVersion requiredapache-airflow3.0.1由于本 Provider 本身只是抽象层真正监听消息的触发器由具体队列 Provider 提供因此需要按需安装对应 extras。文档中声明的可选依赖Optional dependencies为Extra安装的第三方依赖amazonapache-airflow-providers-amazon9.7.0apache.kafkaapache-airflow-providers-apache-kafka1.9.0以启用 Kafka 支持为例pip install apache-airflow-providers-common-messaging[apache.kafka]生产环境建议从官方 Apache 下载站点获取发布包并核对附带的sha512校验和与asc签名避免安装被篡改的构件。若需从源码安装可参考 installing-providers-from-sources。核心抽象BaseMessageQueueProvider统一抽象的基石是 base_provider.py 中的BaseMessageQueueProvider类。它定义了任何想接入该抽象层的队列 Provider 必须实现的契约源码 L26-L68class BaseMessageQueueProvider: scheme: str | None None def scheme_matches(self, scheme: str) - bool: # 默认实现精确字符串比较子类可覆盖 return self.scheme scheme abstractmethod def queue_matches(self, queue: str) - bool: ... # 判断某个 queue URI 是否归属本 Provider abstractmethod def trigger_class(self) - type[BaseEventTrigger]: ... # 返回本 Provider 真正执行监听的具体触发器类 abstractmethod def trigger_kwargs(self, queue: str, **kwargs) - dict: ... # 把 queue URI 换算成触发器构造参数四个成员的职责scheme/scheme_matches基于 scheme 字符串如redispubsub的匹配默认做精确比较这是 2.x 版本引入的新匹配方式queue_matches基于 queue URI 的正则匹配用于向后兼容旧版按 URI 路由的用法trigger_class返回具体 Provider 中的触发器类如 Redis 的AwaitMessageTriggertrigger_kwargs把抽象层收到的参数转换成具体触发器需要的构造参数。以 Redis Provider 为例redis/queues/redis.py 中的实现非常简短直观展示了契约的规模QUEUE_REGEXP r^redis\pubsub:// class RedisPubSubMessageQueueProvider(BaseMessageQueueProvider): scheme redispubsub def trigger_class(self) - type[BaseEventTrigger]: return AwaitMessageTrigger其余匹配与参数转换均复用基类默认实现scheme 精确匹配、URI 正则匹配与AwaitMessageTrigger的同名 kwargs 透传。Provider 发现机制谁被注册进来MessageQueueTrigger并不是硬编码了一张队列清单而是在模块导入时动态发现所有已安装队列 Provider。查看 msg_queue.pyproviders_manager ProvidersManager() providers_manager.initialize_providers_queues() def create_class_by_name(name: str): module_name, class_name name.rsplit(., 1) module importlib.import_module(module_name) return getattr(module, class_name) MESSAGE_QUEUE_PROVIDERS [create_class_by_name(name)() for name in providers_manager.queue_class_names]调用链为ProvidersManager.initialize_providers_queues()扫描各 Provider 的 provider.yaml 中声明的queues属性一组全限定类名再逐一importlib动态加载并实例化。这意味着装了哪个 Provider就有哪条队列可用。例如安装了 amazon 即可用 SQS安装了 apache/kafka 即可用 Kafka安装了 ibm/mq、microsoft/azure 则分别可用 IBM MQ 与 Azure Service Bus若一个队列 Provider 都没有安装trigger属性访问时会直接抛出ValueError(No message queue providers are available.)msg_queue.py L109-L114这是一个清晰的安装诊断信号。目前支持的消息队列完整清单可在 message-queues 文档 中查阅仓库内确认存在queues声明的 Provider 包括amazonSQS、apache.kafka、redisPub/Sub、ibm.mq与microsoft.azureService Bus。MessageQueueTrigger 参数详解MessageQueueTrigger是BaseEventTrigger的子类构造签名msg_queue.py L68-L97为MessageQueueTrigger( *, queue: str | None None, # 已废弃队列 URI如 redispubsub://host/db scheme: str | None None, # 推荐队列 scheme如 kafka、redispubsub、sqs trigger_queue: str | None None, # 指定触发器所在 Triggerer 队列triggerer.queues_enabled **kwargs: Any, # 透传给具体 Provider 触发器的参数 )三个显式参数的要点scheme必填推荐用法队列 scheme如kafka、redispubsub、sqs。源码中若queue与scheme均未提供会抛出ValueError(Eitherqueueorschemeparameter must be provided.)queue已废弃旧版以 URI 形式定位队列的方式。仍传入时会触发AirflowProviderDeprecationWarning且在两者同时提供时queue优先保持向后兼容trigger_queue将本触发器分配到指定 Triggerer 队列对应配置项triggerer.queues_enabled与airflow triggerer的--queues选项。源码中它与废弃的 brokerqueue参数刻意分名通过self._trigger_queue私有属性保存避免与旧参数语义冲突L76-L78 注释明确说明了这一考虑。连接配置要求具体队列的认证信息由各 Provider 的默认 Connection 提供。例如监听 AWS SQS 时需配置名为aws_default的 ConnectionRedis 示例中则通过redis_conn_idredis_default显式指定。匹配与委派逻辑MessageQueueTrigger的trigger属性cached_propertyL107-L162是整套机制的中枢逻辑分三步选择匹配方式queue_uri存在时走provider.queue_matches(uri)URI 匹配否则走provider.scheme_matches(scheme)scheme 匹配唯一性校验若没有 Provider 认领该队列抛出ValueError并列出当前已注册的全部 Provider 名L132-L142若被多于一个Provider 认领scheme 冲突同样抛出ValueErrorL144-L153。这与BaseMessageQueueProvider文档字符串中“匹配规则必须尽可能特异、不得互相重叠”的设计约束相呼应构造并委派命中唯一 Provider 后实例化其trigger_class()。URI 模式下会把参数经trigger_kwargs(queue, **kwargs)换算后传入scheme 模式下直接透传**kwargs。随后serialize()与run()全部委派给内部触发器L164-L169即持久化与事件流完全由具体 Provider 的触发器负责。实战消息到达触发 DAG将触发器挂到Asset的AssetWatcher上并以该 Asset 作为 DAG 的 schedule即可实现事件驱动运行。官方系统测试示例example_message_queue_trigger.py展示了监听 Amazon SQS 的完整 DAG# [START howto_trigger_message_queue] 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 # Define a trigger that listens to an external message queue (AWS SQS in this case) trigger MessageQueueTrigger(schemesqs, sqs_queuehttps://sqs.us-east-1.amazonaws.com/0123456789/Test) # Define an asset that watches for messages on the queue asset Asset(sqs_queue_asset, watchers[AssetWatcher(namesqs_watcher, triggertrigger)]) with DAG(dag_idexample_msgq_watcher, schedule[asset]) as dag: EmptyOperator(task_idtask) # [END howto_trigger_message_queue]注意sqs_queue并非MessageQueueTrigger自己的参数而是经由**kwargs透传给 Amazon Provider 的 SQS 触发器——这正是统一抽象“参数透明传递”设计的体现。工作原理消息队列触发器MessageQueueTrigger监听外部队列AWS SQS、Kafka 或其他消息系统中的消息Asset 与 WatcherAsset抽象外部实体本例中是 SQS 队列AssetWatcher把一个命名触发器与该 Asset 关联名字用于识别“哪个触发器对应哪个 Asset”事件驱动 DAGDAG 不再按固定 schedule 运行而是在 Asset 收到更新队列新消息时执行。在任务中使用消息体触发器会把消息负载payload放入触发事件。DAG 任务可通过triggering_asset_events参数访问它该参数按 Asset 索引触发本次运行的事件每个事件的extra字典中以payload键存放消息体triggers.rstfrom airflow.decorators import task from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger from airflow.sdk import DAG, Asset, AssetWatcher, chain # Define the asset with trigger trigger MessageQueueTrigger(schemekafka, topics[my-kafka-topic]) asset Asset(kafka_queue_asset, watchers[AssetWatcher(namekafka_watcher, triggertrigger)]) with DAG(dag_idexample_msgq_payload, schedule[asset]) as dag: task def process_message(triggering_asset_events): for event in triggering_asset_events[asset]: # Access the message payload payload event.extra[payload] # Process the payload as needed print(fReceived message: {payload}) chain(process_message())这里topics[my-kafka-topic]同样经 kwargs 透传给 Kafka Provider 的触发器与 SQS 示例中透传sqs_queue是同一机制。扩展为新的队列 Provider 接入抽象层providers.rst 给出了接入新 Provider 的两步流程创建 Provider 类在新 Provider 中继承BaseMessageQueueProvider实现全部抽象方法queue_matches、trigger_class、trigger_kwargs并视需要覆盖scheme_matches或设置scheme类属性。参考现成实现Redisredis.py最简形态scheme 精确匹配 kwargs 透传Kafkakafka.pySQSsqs.py注册暴露在新增类所在 Provider 的provider.yaml中通过queues属性暴露其值为队列类的全限定名列表例如queues: - airflow.providers.redis.queues.redis.RedisPubSubMessageQueueProvider注册后无需修改 common-messaging 的任何代码——ProvidersManager会在下次初始化时自动发现该队列 Provider 并加入MESSAGE_QUEUE_PROVIDERS从源码结构看这是发现机制直接带来的扩展性收益。测试验证该 Provider 的单元与系统测试覆盖了上述关键行为可作为行为契约的参考test_msg_queue.pyMessageQueueTrigger的参数校验、scheme/queue 匹配、多 Provider 冲突报错等场景test_base_provider.pyBaseMessageQueueProvider默认 scheme 匹配行为与抽象方法约束example_message_queue_trigger.py端到端系统测试 DAG验证 SQS 消息真实到达后 DAG 被触发。小结Common Messaging Provider 的价值在于一张薄薄的契约BaseMessageQueueProvider的四个方法 provider.yaml中的queues声明。MessageQueueTrigger在模块加载期完成 Provider 动态发现运行期按scheme推荐或queueURI已废弃完成唯一匹配后把序列化、事件流与参数全部委派给具体 Provider 的触发器。对 DAG 作者而言这意味着一套Asset/AssetWatcher写法即可在 SQS、Kafka、Redis Pub/Sub 等队列之间迁移对 Provider 开发者而言接入抽象层只需两个文件改动。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考