You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在单个Python程序中用asyncio同时消费多个Azure Event Hubs事件

问题根因

原写法属于串行等待协程执行,receive_batch()是长期运行的阻塞协程,只要第一个消费任务不主动退出,永远不会执行到第二个客户端的消费逻辑,属于异步任务调度方式错误,和手动加sleep交权无关。

实现方案

通过asyncio.gather将两个消费逻辑作为独立协程并发调度即可,示例代码如下:

import asyncio
from azure.eventhub.aio import EventHubConsumerClient
from azure.eventhub.extensions.checkpointstoreblobaio import BlobCheckpointStore

# 自定义第一个Event Hub的批次事件处理逻辑
async def on_client1_batch(partition_context, events):
    # 这里写你自己的事件处理逻辑
    for event in events:
        print(f"Client1收到事件: {event.body_as_str()}")
    # 提交检查点(可选,用于断点续传)
    await partition_context.update_checkpoint()

# 自定义第二个Event Hub的批次事件处理逻辑
async def on_client2_batch(partition_context, events):
    # 这里写你自己的事件处理逻辑
    for event in events:
        print(f"Client2收到事件: {event.body_as_str()}")
    # 提交检查点(可选,用于断点续传)
    await partition_context.update_checkpoint()

# 封装第一个Event Hub的完整消费逻辑
async def consume_eventhub_1():
    # 替换为你第一个Event Hub的实际参数
    conn_str = "第一个Event Hub的连接字符串"
    consumer_group = "第一个的专属消费组"
    eventhub_name = "第一个Event Hub名称"
    
    # 如果需要持久化检查点,初始化Blob存储(不需要可以去掉这部分)
    checkpoint_store = BlobCheckpointStore.from_connection_string(
        "Blob存储连接字符串",
        "checkpoint容器名1"
    )

    client = EventHubConsumerClient.from_connection_string(
        conn_str=conn_str,
        consumer_group=consumer_group,
        eventhub_name=eventhub_name,
        checkpoint_store=checkpoint_store # 不需要检查点可以去掉这个参数
    )

    async with client:
        await client.receive_batch(
            on_event_batch=on_client1_batch,
            max_batch_size=100,
            max_wait_time=5
        )

# 封装第二个Event Hub的完整消费逻辑
async def consume_eventhub_2():
    # 替换为你第二个Event Hub的实际参数
    conn_str = "第二个Event Hub的连接字符串"
    consumer_group = "第二个的专属消费组"
    eventhub_name = "第二个Event Hub名称"
    
    # 如果需要持久化检查点,初始化Blob存储(不需要可以去掉这部分)
    checkpoint_store = BlobCheckpointStore.from_connection_string(
        "Blob存储连接字符串",
        "checkpoint容器名2"
    )

    client = EventHubConsumerClient.from_connection_string(
        conn_str=conn_str,
        consumer_group=consumer_group,
        eventhub_name=eventhub_name,
        checkpoint_store=checkpoint_store # 不需要检查点可以去掉这个参数
    )

    async with client:
        await client.receive_batch(
            on_event_batch=on_client2_batch,
            max_batch_size=100,
            max_wait_time=5
        )

# 主函数并发调度两个消费任务
async def main():
    await asyncio.gather(
        consume_eventhub_1(),
        consume_eventhub_2()
    )

if __name__ == "__main__":
    asyncio.run(main())

补充说明

  • asyncio.gather会自动把多个协程提交到事件循环调度,两个消费任务的IO等待间隙会自动切换执行,无需手动释放控制权
  • 可在每个消费函数内部加try-except捕获单个客户端的运行异常,避免单个消费任务失败导致整个程序退出
  • 两个消费逻辑使用独立的配置、消费组、检查点存储,完全互不影响

内容的提问来源于stack exchange,提问作者Mircea Stoica

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 05:15:04