Azure Event Hubs SDK for .NET - 用于.NET的高吞吐量事件流式传输SDK,用于通过Azure Event Hubs发送和接收事件。
📥 下载地址:
https://github.com/sickn33/antigravity-awesome-skills/tree/main/skills/azure-eventhub-dotnet
技能概述
azure-eventhub-dotnet 技能是一个专门用于.NET开发的Azure Event Hubs SDK技能包。它提供了高吞吐量的事件流式传输功能,支持事件的发送和接收操作。该技能包适用于需要构建实时数据管道、事件驱动架构或大数据流处理应用的.NET开发者。
主要功能
- 事件发送:支持批量发送和缓冲发送两种模式,可高效处理大量事件数据
- 事件接收:提供EventProcessorClient用于生产环境的事件处理,支持检查点和负载均衡
- 分区管理:支持分区操作,可通过分区键保证事件顺序性
- 身份认证:集成Azure Identity,支持DefaultAzureCredential进行安全认证
- 检查点策略:提供多种检查点策略,平衡吞吐量和可靠性
- ASP.NET Core集成:支持依赖注入,可轻松集成到ASP.NET Core应用中
触发条件
在以下情况下应该调用此技能:
- 用户需要在.NET应用中集成Azure Event Hubs
- 需要实现高吞吐量的事件发送功能
- 需要构建可靠的事件消费者应用
- 需要处理实时数据流或事件驱动架构
- 需要了解Event Hubs的分区和检查点机制
使用场景
场景1:实时数据管道
构建高吞吐量的数据管道,将事件从生产者发送到Event Hubs,再由消费者处理并存储到数据库或数据湖。
场景2:事件驱动架构
在微服务架构中使用Event Hubs作为事件总线,实现服务间的解耦和异步通信。
场景3:日志和遥测收集
收集应用程序日志、指标和遥测数据,进行实时分析和监控。
处理过程
1. 安装依赖
通过NuGet安装必要的包:
# 核心包(发送和简单接收)
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
2. 创建生产者客户端
使用DefaultAzureCredential创建EventHubProducerClient:
using Azure.Identity;
using Azure.Messaging.EventHubs;
using Azure.Messaging.EventHubs.Producer;var credential = new DefaultAzureCredential();
var producer = new EventHubProducerClient(
"<namespace>.servicebus.windows.net",
"<event-hub-name>",
credential);
3. 发送事件
创建批次并发送事件:
using EventDataBatch batch = await producer.CreateBatchAsync();
foreach (var eventData in events)
{
if (!batch.TryAdd(eventData))
{
await producer.SendAsync(batch);
batch = await producer.CreateBatchAsync();
batch.TryAdd(eventData);
}
}if (batch.Count > 0)
{
await producer.SendAsync(batch);
}
4. 创建消费者处理器
使用EventProcessorClient处理事件:
var processor = new EventProcessorClient(
blobClient,
EventHubConsumerClient.DefaultConsumerGroup,
fullyQualifiedNamespace,
eventHubName,
new DefaultAzureCredential());processor.ProcessEventAsync += async args =>
{
Console.WriteLine($"Data: {args.Data.EventBody}");
await args.UpdateCheckpointAsync();
};processor.ProcessErrorAsync += args =>
{
Console.WriteLine($"Error: {args.Exception.Message}");
return Task.CompletedTask;
};await processor.StartProcessingAsync();
输入要求
使用此技能时,用户需要提供:
- Event Hubs命名空间:完全限定的命名空间名称(如:mynamespace.servicebus.windows.net)
- Event Hub名称:事件中心的名称
- 身份认证凭据:DefaultAzureCredential或连接字符串
- 存储账户信息:用于检查点的Blob存储连接字符串和容器名称(消费者需要)
- RBAC角色:需要分配适当的Azure角色(Sender、Receiver或Owner)
输出说明
技能将提供:
- 完整的代码示例:包含生产者和消费者的完整实现
- 最佳实践指南:关于批次处理、检查点策略和错误处理的建议
- 配置说明:环境变量和连接配置的详细说明
- 故障排除建议:常见问题的解决方案
客户端类型对比
| 客户端 | 用途 | 使用场景 |
|---|---|---|
| EventHubProducerClient | 立即批量发送事件 | 实时发送,完全控制批处理 |
| EventHubBufferedProducerClient | 自动批处理和后台发送 | 高容量、即发即弃场景 |
| EventHubConsumerClient | 简单事件读取 | 仅用于原型开发,不适用于生产 |
| EventProcessorClient | 生产环境事件处理 | 生产环境接收事件的首选 |
使用示例
示例1:批量发送事件
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)
{
batch.TryAdd(eventData);
}await producer.SendAsync(batch);
示例2:生产环境事件处理
var processor = new EventProcessorClient(
blobClient,
EventHubConsumerClient.DefaultConsumerGroup,
fullyQualifiedNamespace,
eventHubName,
new DefaultAzureCredential());processor.ProcessEventAsync += async args =>
{
Console.WriteLine($"Partition: {args.Partition.PartitionId}");
Console.WriteLine($"Data: {args.Data.EventBody}");
await args.UpdateCheckpointAsync();
};processor.ProcessErrorAsync += args =>
{
Console.WriteLine($"Error: {args.Exception.Message}");
return Task.CompletedTask;
};await processor.StartProcessingAsync();
最佳实践
- 使用EventProcessorClient接收事件:不要在生产环境中使用EventHubConsumerClient
- 策略性检查点:在处理N个事件或时间间隔后检查点,而不是每个事件
- 使用分区键:在分区内保证顺序性
- 重用客户端:创建一次,作为单例使用(线程安全)
- 使用await using:确保正确释放资源
- 处理ProcessErrorAsync:始终注册错误处理程序
- 批量事件:使用CreateBatchAsync()遵守大小限制
- 使用缓冲生产者:对于高容量场景,使用自动批处理
检查点策略
| 策略 | 使用场景 |
|---|---|
| 每个事件 | 低容量、关键数据 |
| 每N个事件 | 平衡吞吐量/可靠性 |
| 基于时间 | 一致的检查点间隔 |
| 批次完成 | 处理逻辑批次后 |
相关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 |
注意事项
- 此技能仅适用于任务明确匹配上述范围的情况
- 输出不应替代环境特定的验证、测试或专家审查
- 如果缺少必需的输入、权限、安全边界或成功标准,请停止并请求澄清