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

如何在Azure Event Hub中创建多个接收器以避免重复消费?

解决方案

问题根因

你当前出现重复消费的问题是由以下几个代码错误导致的:

  1. 你已经创建了BlobCheckpointStore检查点存储实例,但在初始化EventHubConsumerClient时没有传入该实例,导致两个消费者无法共享消费进度,也无法触发同消费组的分区自动分配逻辑。
  2. 调用receive_batch方法时硬编码了PARTITION="0"参数,强制两个消费者都监听0号分区,同个分区的消息会被同消费组下的所有消费者接收,自然就出现重复消费。
  3. 现有代码逻辑中收到空事件批次就关闭消费者客户端,会导致消费者提前退出,无法持续正常工作。

修改步骤

  • 初始化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 11:54:09