如何在单个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
相关产品推荐
相关产品推荐

