Azure Event Hubs .NET SDK
Azure.Messaging.EventHubs (.NET)
用于通过 Azure Event Hubs 发送和接收事件的高吞吐量事件流 SDK。
安装
# 核心包(发送和简单接收)
dotnet add package Azure.Messaging.EventHubs
处理器包(带检查点的生产级接收)
dotnet add package Azure.Messaging.EventHubs.Processor
身份验证
dotnet add package Azure.Identity
用于检查点(EventProcessorClient 必需)
dotnet add package Azure.Storage.Blobs当前版本: Azure.Messaging.EventHubs v5.12.2, Azure.Messaging.EventHubs.Processor v5.12.2
环境变量
EVENTHUB_FULLY_QUALIFIED_NAMESPACE=<namespace>.servicebus.windows.net
EVENTHUB_NAME=<event-hub-name>
用于检查点 (EventProcessorClient)
BLOB_STORAGE_CONNECTION_STRING=<storage-connection-string>
BLOB_CONTAINER_NAME=<checkpoint-container>
备选方案:连接字符串认证(不推荐用于生产环境)
EVENTHUB_CONNECTION_STRING=Endpoint=sb://<namespace>.servicebus.windows.net/;SharedAccessKeyName=...身份验证
using Azure.Identity;
using Azure.Messaging.EventHubs;
using Azure.Messaging.EventHubs.Producer;
// 生产环境请始终使用 DefaultAzureCredential
var credential = new DefaultAzureCredential();
var fullyQualifiedNamespace = Environment.GetEnvironmentVariable("EVENTHUB_FULLY_QUALIFIED_NAMESPACE");
var eventHubName = Environment.GetEnvironmentVariable("EVENTHUB_NAME");
var producer = new EventHubProducerClient(
fullyQualifiedNamespace,
eventHubName,
credential);
必需的 RBAC 角色:
- 发送:
Azure Event Hubs Data Sender
- 接收:
Azure Event Hubs Data Receiver
- 两者兼有:
Azure Event Hubs Data Owner
客户端类型
| 客户端 | 用途 | 使用场景 |
|--------|---------|-------------|
| EventHubProducerClient | 立即以批次发送事件 | 实时发送,需完全控制批处理 |
| EventHubBufferedProducerClient | 自动批处理并后台发送 | 高吞吐量、发后即忘 (fire-and-forget) 场景 |
| EventHubConsumerClient | 简单的事件读取 | 仅用于原型开发,不适用于生产环境 |
| EventProcessorClient | 生产级事件处理 | 生产环境接收事件请始终使用此客户端 |
核心工作流
1. 发送事件(批处理)
using Azure.Identity;
using Azure.Messaging.EventHubs;
using Azure.Messaging.EventHubs.Producer;
await using var producer = new EventHubProducerClient(
fullyQualifiedNamespace,
eventHubName,
new DefaultAzureCredential());
// 创建批次(自动遵守大小限制)
using EventDataBatch batch = await producer.CreateBatchAsync();
// 向批次添加事件
var events = new[]
{
new EventData(BinaryData.FromString("{\"id\": 1, \"message\": \"Hello\"}")),
new EventData(BinaryData.FromString("{\"id\": 2, \"message\": \"World\"}"))
};
foreach (var eventData in events)
{
if (!batch.TryAdd(eventData))
{
// 批次已满 - 发送当前批次并创建新批次
await producer.SendAsync(batch);
batch = await producer.CreateBatchAsync();
if (!batch.TryAdd(eventData))
{
throw new Exception("事件过大,无法放入空批次");
}
}
}
// 发送剩余事件
if (batch.Count > 0)
{
await producer.SendAsync(batch);
}
2. 发送事件(缓冲模式 - 高吞吐量)
using Azure.Messaging.EventHubs.Producer;
var options = new EventHubBufferedProducerClientOptions
{
MaximumWaitTime = TimeSpan.FromSeconds(1)
};
await using var producer = new EventHubBufferedProducerClient(
fullyQualifiedNamespace,
eventHubName,
new DefaultAzureCredential(),
options);
// 处理发送成功/失败
producer.SendEventBatchSucceededAsync += args =>
{
Console.WriteLine($"批次发送成功: {args.EventBatch.Count} 个事件");
return Task.CompletedTask;
};
producer.SendEventBatchFailedAsync += args =>
{
Console.WriteLine($"批次发送失败: {args.Exception.Message}");
return Task.CompletedTask;
};
// 将事件入队(在后台自动发送)
for (int i = 0; i < 1000; i++)
{
await producer.EnqueueEventAsync(new EventData($"Event {i}"));
}
// 在释放前刷新剩余事件
await producer.FlushAsync();
3. 接收事件(生产环境 - EventProcessorClient)
using Azure.Identity;
using Azure.Messaging.EventHubs;
using Azure.Messaging.EventHubs.Consumer;
using Azure.Messaging.EventHubs.Processor;
using Azure.Storage.Blobs;
// 用于检查点(checkpointing)的 Blob 容器
var blobClient = new BlobContainerClient(
Environment.GetEnvironmentVariable("BLOB_STORAGE_CONNECTION_STRING"),
Environment.GetEnvironmentVariable("BLOB_CONTAINER_NAME"));
await blobClient.CreateIfNotExistsAsync();
// 创建处理器
var processor = new EventProcessorClient(
blobClient,
EventHubConsumerClient.DefaultConsumerGroup,
fullyQualifiedNamespace,
eventHubName,
new DefaultAzureCredential());
// 处理事件
processor.ProcessEventAsync += async args =>
{
Console.WriteLine($"分区: {args.Partition.PartitionId}");
Console.WriteLine($"数据: {args.Data.EventBody}");
// 处理后更新检查点(或进行批次检查点更新)
await args.UpdateCheckpointAsync();
};
// 处理错误
processor.ProcessErrorAsync += args =>
{
Console.WriteLine($"错误: {args.Exception.Message}");
Console.WriteLine($"分区: {args.PartitionId}");
return Task.CompletedTask;
};
// 开始处理
await processor.StartProcessingAsync();
// 运行直到被取消
await Task.Delay(Timeout.Infinite, cancellationToken);
// 优雅停止
await processor.StopProcessingAsync();
4. 分区操作
// 获取分区 ID
string[] partitionIds = await producer.GetPartitionIdsAsync();
// 发送到指定分区(谨慎使用)
var options = new SendEventOptions
{
PartitionId = "0"
};
await producer.SendAsync(events, options);
// 使用分区键(推荐用于保证顺序)
var batchOptions = new CreateBatchOptions
{
PartitionKey = "customer-123" // 具有相同键的事件将进入同一分区
};
using var batch = await producer.CreateBatchAsync(batchOptions);
EventPosition 选项
控制读取的起始位置:
// 从最早的消息开始
EventPosition.Earliest
// 从最新消息开始(仅接收新事件)
EventPosition.Latest
// 从特定偏移量开始
EventPosition.FromOffset(12345)
// 从特定序列号开始
EventPosition.FromSequenceNumber(100)
// 从特定时间开始
EventPosition.FromEnqueuedTime(DateTimeOffset.UtcNow.AddHours(-1))
ASP.NET Core 集成
// Program.cs
using Azure.Identity;
using Azure.Messaging.EventHubs.Producer;
using Microsoft.Extensions.Azure;
builder.Services.AddAzureClients(clientBuilder =>
{
clientBuil
der.AddEventHubProducerClient(
builder.Configuration["EventHub:FullyQualifiedNamespace"],
builder.Configuration["EventHub:Name"]);
clientBuilder.UseCredential(new DefaultAzureCredential());
});
// 在控制器/服务中注入
public class EventService
{
private readonly EventHubProducerClient _producer;
public EventService(EventHubProducerClient producer)
{
_producer = producer;
}
public async Task SendAsync(string message)
{
using var batch = await _producer.CreateBatchAsync();
batch.TryAdd(new EventData(message));
await _producer.SendAsync(batch);
}
}
最佳实践
1. 接收端使用 EventProcessorClient —— 生产环境请勿使用 EventHubConsumerClient
2. 策略性设置检查点 (Checkpoint) —— 在处理 N 个事件或达到一定时间间隔后设置,而非每个事件都设置
3. 使用分区键 (Partition Keys) —— 以确保分区内的顺序保证
4. 复用客户端 —— 创建一次,作为单例使用(线程安全)
5. 使用 await using —— 确保正确释放资源
6. 处理 ProcessErrorAsync —— 务必注册错误处理器
7. 批量发送事件 —— 使用 CreateBatchAsync() 以遵守大小限制
8. 使用缓冲生产者 —— 适用于具有自动批处理功能的高吞吐量场景
错误处理
using Azure.Messaging.EventHubs;
try
{
await producer.SendAsync(batch);
}
catch (EventHubsException ex) when (ex.Reason == EventHubsException.FailureReason.ServiceBusy)
{
// 使用指数退避重试
await Task.Delay(TimeSpan.FromSeconds(5));
}
catch (EventHubsException ex) when (ex.IsTransient)
{
// 瞬时错误 - 可以安全重试
Console.WriteLine($"Transient error: {ex.Message}");
}
catch (EventHubsException ex)
{
// 非瞬时错误
Console.WriteLine($"Error: {ex.Reason} - {ex.Message}");
}
检查点策略
| 策略 | 使用场景 |
|----------|-------------|
| 每个事件 | 低吞吐量,关键数据 |
| 每 N 个事件 | 吞吐量与可靠性的平衡 |
| 基于时间 | 固定的检查点间隔 |
| 批处理完成 | 处理完一个逻辑批次后 |
// 每 100 个事件设置一次检查点
private int _eventCount = 0;
processor.ProcessEventAsync += async args =>
{
// 处理事件...
_eventCount++;
if (_eventCount >= 100)
{
await args.UpdateCheckpointAsync();
_eventCount = 0;
}
};
相关 SDK
| SDK | 用途 | 安装命令 |
|-----|---------|---------|
| Azure.Messaging.EventHubs | 核心发送/接收 | dotnet add package Azure.Messaging.EventHubs |
| Azure.Messaging.EventHubs.Processor | 生产级处理 | dotnet add package Azure.Messaging.EventHubs.Processor |
| Azure.ResourceManager.EventHubs | 管理平面(创建 Hub) | dotnet add package Azure.ResourceManager.EventHubs |
| Microsoft.Azure.WebJobs.Extensions.EventHubs | Azure Functions 绑定 | dotnet add package Microsoft.Azure.WebJobs.Extensions.EventHubs |
适用场景
本技能适用于执行概述中描述的工作流或操作。局限性
- 仅在任务明确符合上述范围时使用此技能。
- 不要将输出结果视为针对特定环境的验证、测试或专家评审的替代方案。
- 如果缺少必要的输入、权限、安全边界或成功标准,请停止并请求澄清。