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

继承自 MessageQueueTrigger

等待队列中的消息

以下示例展示了如何配置 Airflow DAG,使其在 Azure Service Bus 中收到消息时被触发。

tests/system/microsoft/azure/example_event_schedule_asb.py[source]

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")

它是如何工作的

  1. Azure Service Bus 消息队列触发器:`AzureServiceBusQueueTrigger` 监听来自 Azure Service Bus 队列的消息。

  2. Asset 与 Watcher:`Asset` 抽象了外部实体,在本例中即 Azure Service Bus 队列。`AssetWatcher` 将触发器与名称关联。该名称帮助您识别哪个触发器对应哪个资产。

  3. 事件驱动的 DAG:与其按固定计划运行,DAG 会在资产收到更新时执行(例如,队列中出现新消息)。

有关如何使用该触发器,请参考 Messaging Trigger 的文档。

此条目是否有帮助?