如何在Azure Event Hub中创建多个接收器以避免重复消费?
解决方案
问题根因
你当前出现重复消费的问题是由以下几个代码错误导致的:
- 你已经创建了
BlobCheckpointStore检查点存储实例,但在初始化EventHubConsumerClient时没有传入该实例,导致两个消费者无法共享消费进度,也无法触发同消费组的分区自动分配逻辑。 - 调用
receive_batch方法时硬编码了PARTITION="0"参数,强制两个消费者都监听0号分区,同个分区的消息会被同消费组下的所有消费者接收,自然就出现重复消费。 - 现有代码逻辑中收到空事件批次就关闭消费者客户端,会导致消费者提前退出,无法持续正常工作。
修改步骤
- 初始化
EventHubConsumerClient时传入checkpoint_store参数,让两个消费者共享同一个Blob存储的消费进度,实现同消费组下的分区自动负载分配。 - 删除
receive_batch中的PARTITION="0"参数,让客户端自动分配分区,同消费组下的多个消费者会自动分摊不同分区的消息,不会重复消费。 - 移除收到空事件就关闭客户端的逻辑,保证消费者可以持续监听事件。
- 额外注意:如果要两个消费者都同时工作分摊负载,请确保你的Event Hub实例的分区数≥2,单分区下同消费组只能有一个消费者消费消息,另一个会处于闲置状态。
修正后的代码
from azure.storage.blob import BlobServiceClient, ContainerClient from azure.core.exceptions import ResourceExistsError from azure.eventhub import EventHubConsumerClient from azure.eventhub.extensions.checkpointstoreblob import BlobCheckpointStore # Eventhub访问凭据 connection_str = **** consumer_group = '$Default' eventhub_name = **** # Blob存储凭据 storage_connection_str = **** container_name = **** # 初始化Blob检查点存储,用于同消费组消费者共享进度和分区分配 checkpoint_store = BlobCheckpointStore.from_connection_string(storage_connection_str, container_name) # 初始化Blob存储客户端,用于后续写入结果 blob_service_client = BlobServiceClient.from_connection_string(storage_connection_str) container_client = blob_service_client.get_container_client(container_name) try: container_client.create_container() except ResourceExistsError: print("Container already exists.") def get_messages(): # 初始化消费者客户端时传入checkpoint_store client = EventHubConsumerClient.from_connection_string( connection_str, consumer_group, eventhub_name=eventhub_name, checkpoint_store=checkpoint_store # 新增这行,传入检查点存储 ) def on_event_batch(partition_context, events): print(f"Received event from partition {partition_context.partition_id}, 事件数量: {len(events)}") if len(events) > 0: for event in events: list_ = event.body_as_json() # 这里添加你处理事件的业务逻辑 pass # 处理完一批事件后更新检查点 partition_context.update_checkpoint() try: with client: client.receive_batch( on_event_batch=on_event_batch, starting_position="-1", # 首次启动从分区最早位置消费,后续启动会从检查点位置消费 max_wait_time=5 # 批次最大等待时间,可根据业务调整 # 删掉了PARTITION="0"参数,让客户端自动分配分区 ) except KeyboardInterrupt: print('Stopped receiving.') if __name__ == "__main__": get_messages()
修改完成后分别运行两个消费者实例,同消费组下的两个消费者会自动分配Event Hub的分区,100条事件会被两个消费者分摊处理,不会出现重复消费的情况。
内容的提问来源于stack exchange,提问作者Nagesh Singh Chauhan
相关产品推荐
相关产品推荐

