使用Python设置Azure Event Hub Partition_ID后事件存入全部分区的问题排查
Azure Event Hub 指定分区发送事件问题解决
问题描述
我是Azure Event Hub的新手,已在Event Hub中初始化3个分区。我使用Python代码指定Partition_ID,通过EventHub的Capture功能将事件存储到Azure Blob Storage。代码可以正常执行,但事件被存入了所有3个分区文件夹,而非仅第2个;控制台打印event_data_batch显示Partition_ID已设为2。
原代码如下:
import asyncio from azure.eventhub import EventData from azure.eventhub.aio import EventHubProducerClient EVENT_HUB_CONNECTION_STR = "..." EVENT_HUB_NAME = "....." async def run(): # Create a producer client to send messages to the event hub. # Specify a connection string to your event hubs namespace and # the event hub name. producer = EventHubProducerClient.from_connection_string( conn_str=EVENT_HUB_CONNECTION_STR, eventhub_name=EVENT_HUB_NAME ) async with producer: # Create a batch. event_data_batch = await producer.create_batch(**partition_id="2"**) # Add events to the batch. event_data_batch.add(EventData("First event---")) event_data_batch.add(EventData("Second event---")) event_data_batch.add(EventData("Third event---")) event_data_batch.add(EventData("Fourth event---")) event_data_batch.add(EventData("Fifth event---")) event_data_batch.add(EventData("Sixth event---")) # Send the batch of events to the event hub. await producer.send_batch(event_data_batch) print(event_data_batch) print('Batch Sent') asyncio.run(run())
问题原因与解决步骤
- 语法错误导致参数失效:原代码中
create_batch(**partition_id="2"**)的写法有误,双星号是用于解包字典参数的语法,这里直接传递关键字参数不需要双星号,错误语法导致partition_id参数未被正确解析,Producer使用了默认的轮询分区策略,将事件分发到了所有分区。 - 修正代码参数写法:将
create_batch(**partition_id="2"**)改为create_batch(partition_id="2"),确保Producer明确将事件批次发送到指定的分区2。
修正后的完整代码:
import asyncio from azure.eventhub import EventData from azure.eventhub.aio import EventHubProducerClient EVENT_HUB_CONNECTION_STR = "..." EVENT_HUB_NAME = "....." async def run(): # 创建生产者客户端以向事件中心发送消息 # 指定事件中心命名空间的连接字符串和事件中心名称 producer = EventHubProducerClient.from_connection_string( conn_str=EVENT_HUB_CONNECTION_STR, eventhub_name=EVENT_HUB_NAME ) async with producer: # 创建事件批次并指定目标分区ID event_data_batch = await producer.create_batch(partition_id="2") # 向批次中添加事件 event_data_batch.add(EventData("First event---")) event_data_batch.add(EventData("Second event---")) event_data_batch.add(EventData("Third event---")) event_data_batch.add(EventData("Fourth event---")) event_data_batch.add(EventData("Fifth event---")) event_data_batch.add(EventData("Sixth event---")) # 将事件批次发送到事件中心 await producer.send_batch(event_data_batch) print(event_data_batch) print('Batch Sent') asyncio.run(run())
验证方式
- 执行修正后的代码,查看Event Hub的分区监控指标,确认仅分区2的入站消息数增加
- 检查Azure Blob Storage中的Capture文件,确认只有
partition=2的文件夹下生成了新的事件存储文件
内容的提问来源于stack exchange,提问作者Raheel
相关产品推荐
相关产品推荐

