Azure Service Bus 消息队列¶
Azure Service Bus 队列提供程序¶
实现于 AzureServiceBusMessageQueueProvider
Azure Service Bus 队列提供程序是一个 BaseMessageQueueProvider,它使用 Azure Service Bus 作为底层消息队列系统。它让您能够在 Airflow 工作流中使用 Azure Service Bus 队列发送和接收消息,并通过 MessageQueueTrigger 这一通用消息队列接口进行交互。
它使用
azure+servicebus作为标识提供程序的方案(scheme)。有关参数定义,请参阅
AzureServiceBusQueueTrigger。from airflow.providers.common.messaging.triggers.msg_queue import MessageQueueTrigger from airflow.sdk import Asset, AssetWatcher trigger = MessageQueueTrigger( scheme="azure+servicebus", # AzureServiceBusQueueTrigger parameters queues=["my-queue"], azure_service_bus_conn_id="azure_service_bus_default", poll_interval=60, ) asset = Asset( "asb_queue_asset", watchers=[AssetWatcher(name="asb_watcher", trigger=trigger)], )
Azure Service Bus 消息队列触发器¶
实现于 AzureServiceBusQueueTrigger
等待队列中的消息¶
以下示例展示了如何配置 Airflow DAG,使其在 Azure Service Bus 中收到消息时被触发。
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 Azure Service Bus queue
trigger = MessageQueueTrigger(
scheme="azure+servicebus",
queues=["my-queue"],
azure_service_bus_conn_id="azure_service_bus_default",
poll_interval=60,
)
# Define an asset that watches for messages on the Azure Service Bus queue
asset = Asset(
"event_schedule_asb_asset_1",
watchers=[AssetWatcher(name="event_schedule_asb_watcher_1", trigger=trigger)],
)
with DAG(
dag_id="example_event_schedule_asb",
schedule=[asset],
) as dag:
process_message_task = EmptyOperator(task_id="process_asb_message")
它是如何工作的¶
Azure Service Bus 消息队列触发器:`AzureServiceBusQueueTrigger` 监听来自 Azure Service Bus 队列的消息。
Asset 与 Watcher:`Asset` 抽象了外部实体,在本例中即 Azure Service Bus 队列。`AssetWatcher` 将触发器与名称关联。该名称帮助您识别哪个触发器对应哪个资产。
事件驱动的 DAG:与其按固定计划运行,DAG 会在资产收到更新时执行(例如,队列中出现新消息)。
有关如何使用该触发器,请参考 Messaging Trigger 的文档。