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

使用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())

验证方式

  1. 执行修正后的代码,查看Event Hub的分区监控指标,确认仅分区2的入站消息数增加
  2. 检查Azure Blob Storage中的Capture文件,确认只有partition=2的文件夹下生成了新的事件存储文件

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 09:33:34