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

如何使用Azure EventProcessor从EventHubs模拟器消费数据?

解决方案

针对你遇到的Blob容器不存在、Event Hub实体找不到/端点异常问题,可通过以下步骤解决:

1. 确保Event Hub在模拟器中已存在

Docker版Event Hubs模拟器不会自动创建Event Hub,必须手动创建:

  • 使用Azure CLI连接本地模拟器执行创建命令:
    az eventhubs eventhub create --name <你的EventHub名称> --namespace-name emulatorns1 --resource-group dummy
    
    注意:模拟器默认命名空间为emulatorns1,代码中指定的eventHubName必须与创建的名称完全一致。

2. 修复Blob容器不存在的问题

BlobCheckpointStore不会自动创建Blob容器,需手动初始化:
在创建BlobContainerClient后添加容器创建逻辑:

var storageClient = new BlobContainerClient(blobConnection, "localhost-test");
// 确保容器存在,不存在则自动创建
await storageClient.CreateIfNotExistsAsync();

3. 修正Event Hubs客户端配置,解决端点异常

SDK使用UseDevelopmentEmulator=true时,指定WebSockets传输类型可避免端口映射问题,确保端点正确指向localhost:

  • 修改BasicConsumer构造函数,支持传入EventHubClientOptions:
    internal class BasicConsumer : PluggableCheckpointStoreEventProcessor<EventProcessorPartition>
    {
        public BasicConsumer(
            BlobContainerClient storageClient,
            string connectionString,
            string consumerGroupName,
            string eventHubName,
            int batchSize,
            EventHubClientOptions clientOptions = null)
                : base(
                      checkpointStore: new BlobCheckpointStore(storageClient),
                      eventBatchMaximumCount: batchSize,
                      consumerGroup: consumerGroupName,
                      connectionString: connectionString,
                      eventHubName: eventHubName,
                      clientOptions: clientOptions ?? new EventHubClientOptions()
                      )
        {
        }
    
        // 保留原有重写方法...
    }
    
  • 初始化客户端时指定WebSockets传输:
    var eventHubOptions = new EventHubClientOptions
    {
        TransportType = EventHubsTransportType.AmqpWebSockets
    };
    
    var basicConsumer = new BasicConsumer(storageClient, eventHubConnection, consumerGroupName, eventHubName, batchSize, eventHubOptions);
    

4. 启动Event Processor(关键步骤)

实例化BasicConsumer后必须手动启动处理流程,否则不会触发任何回调:

var cancellationTokenSource = new CancellationTokenSource();
// 启动事件处理
await basicConsumer.StartProcessingAsync(cancellationTokenSource.Token);

// 保持程序运行(示例为控制台程序)
Console.WriteLine("处理器已启动,按回车停止...");
Console.ReadLine();

// 停止处理
await basicConsumer.StopProcessingAsync(cancellationTokenSource.Token);

5. 避免Azurite冲突

不要同时运行独立的Azurite和模拟器内置的Azurite(两者默认端口均为10000,会导致冲突),仅使用模拟器自带的Azurite即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 13:03:21