Azure Eventhub Dotnet

作者 microsoft354361d83247MIT收录于 2026年10月8日更新于 2026年10月8日

Azure Event Hubs SDK for .NET. Use for high-throughput event streaming: sending events (EventHubProducerClient, EventHubBufferedProducerClient), receiving events (EventProcessorClient with checkpointing), partition management, and real-time data ingestion. Triggers: "Event Hubs", "event streaming", "EventHubProducerClient", "EventProcessorClient", "send events", "receive events", "checkpointing", "partition".

AI 生成的概览

使用 Azure Event Hubs .NET SDK 发送、接收事件流并设置检查点的参考指南。

功能
该技能提供使用 Azure.Messaging.EventHubs 系列包进行 Azure Event Hubs 事件流处理的 .NET 代码指导。内容涵盖使用 EventHubProducerClient 和 EventHubBufferedProducerClient 发送事件、使用 EventProcessorClient 配合 Blob 检查点接收事件、分区操作、事件位置选项、ASP.NET Core 集成、错误处理以及检查点策略。产出为代码片段和配置指导,而非可执行产物。
适用场景
适用于编写或审查通过 Azure Event Hubs 发送或接收事件的 .NET 代码。适合高吞吐量数据摄取、带检查点的生产环境事件处理,以及分区或顺序相关决策。
运行要求
需要 .NET SDK 以及 Azure.Messaging.EventHubs、Azure.Messaging.EventHubs.Processor、Azure.Identity 和 Azure.Storage.Blobs 包。需要 Event Hubs 命名空间和事件中心名称,以及具有相应 RBAC 角色的 Entra ID 凭据或连接字符串,另需用于检查点的 Blob 容器。不包含脚本,仅为说明文档。

Azure.Messaging.EventHubs (.NET)

High-throughput event streaming SDK for sending and receiving events via Azure Event Hubs.

Installation

bash
# Core package (sending and simple receiving)dotnet add package Azure.Messaging.EventHubs
# Processor package (production receiving with checkpointing)dotnet add package Azure.Messaging.EventHubs.Processor
# Authenticationdotnet add package Azure.Identity
# For checkpointing (required by EventProcessorClient)dotnet add package Azure.Storage.Blobs

Current Versions: Azure.Messaging.EventHubs v5.12.2, Azure.Messaging.EventHubs.Processor v5.12.2

Environment Variables

bash
EVENTHUB_FULLY_QUALIFIED_NAMESPACE=<namespace>.servicebus.windows.net  # Required: Event Hubs fully qualified namespaceEVENTHUB_NAME=<event-hub-name>  # Required: Event Hub nameBLOB_STORAGE_CONNECTION_STRING=<storage-connection-string>  # Alternative to Entra ID authBLOB_CONTAINER_NAME=<checkpoint-container>  # Required: checkpoint container nameEVENTHUB_CONNECTION_STRING=Endpoint=sb://<namespace>.servicebus.windows.net/;SharedAccessKeyName=...  # Alternative to Entra ID authAZURE_TOKEN_CREDENTIALS=prod  # Required only if DefaultAzureCredential is used in production

Authentication

csharp
using Azure.Identity;using Azure.Messaging.EventHubs;using Azure.Messaging.EventHubs.Producer;
// Local dev: DefaultAzureCredential. Production: set AZURE_TOKEN_CREDENTIALS=prod or AZURE_TOKEN_CREDENTIALS=<specific_credential>var credential = new DefaultAzureCredential(    DefaultAzureCredential.DefaultEnvironmentVariableName);// Or use a specific credential directly in production:// See https://learn.microsoft.com/dotnet/api/overview/azure/identity-readme?view=azure-dotnet#credential-classes// var credential = new ManagedIdentityCredential();
var fullyQualifiedNamespace = Environment.GetEnvironmentVariable("EVENTHUB_FULLY_QUALIFIED_NAMESPACE");var eventHubName = Environment.GetEnvironmentVariable("EVENTHUB_NAME");
var producer = new EventHubProducerClient(    fullyQualifiedNamespace,    eventHubName,    credential);

Required RBAC Roles:

  • Sending: Azure Event Hubs Data Sender
  • Receiving: Azure Event Hubs Data Receiver
  • Both: Azure Event Hubs Data Owner

Client Types

ClientPurposeWhen to Use
EventHubProducerClientSend events immediately in batchesReal-time sending, full control over batching
EventHubBufferedProducerClientAutomatic batching with background sendingHigh-volume, fire-and-forget scenarios
EventHubConsumerClientSimple event readingPrototyping only, NOT for production
EventProcessorClientProduction event processingAlways use this for receiving in production

Core Workflow

1. Send Events (Batch)

csharp
using Azure.Identity;using Azure.Messaging.EventHubs;using Azure.Messaging.EventHubs.Producer;
await using var producer = new EventHubProducerClient(    fullyQualifiedNamespace,    eventHubName,    new DefaultAzureCredential());
// Create a batch (respects size limits automatically)using EventDataBatch batch = await producer.CreateBatchAsync();
// Add events to batchvar 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))    {        // Batch is full - send it and create a new one        await producer.SendAsync(batch);        batch = await producer.CreateBatchAsync();                if (!batch.TryAdd(eventData))        {            throw new Exception("Event too large for empty batch");        }    }}
// Send remaining eventsif (batch.Count > 0){    await producer.SendAsync(batch);}

2. Send Events (Buffered - High Volume)

csharp
using Azure.Messaging.EventHubs.Producer;
var options = new EventHubBufferedProducerClientOptions{    MaximumWaitTime = TimeSpan.FromSeconds(1)};
await using var producer = new EventHubBufferedProducerClient(    fullyQualifiedNamespace,    eventHubName,    new DefaultAzureCredential(),    options);
// Handle send success/failureproducer.SendEventBatchSucceededAsync += args =>{    Console.WriteLine($"Batch sent: {args.EventBatch.Count} events");    return Task.CompletedTask;};
producer.SendEventBatchFailedAsync += args =>{    Console.WriteLine($"Batch failed: {args.Exception.Message}");    return Task.CompletedTask;};
// Enqueue events (sent automatically in background)for (int i = 0; i < 1000; i++){    await producer.EnqueueEventAsync(new EventData($"Event {i}"));}
// Flush remaining events before disposingawait producer.FlushAsync();

3. Receive Events (Production - EventProcessorClient)

csharp
using Azure.Identity;using Azure.Messaging.EventHubs;using Azure.Messaging.EventHubs.Consumer;using Azure.Messaging.EventHubs.Processor;using Azure.Storage.Blobs;
// Blob container for checkpointingvar blobClient = new BlobContainerClient(    Environment.GetEnvironmentVariable("BLOB_STORAGE_CONNECTION_STRING"),    Environment.GetEnvironmentVariable("BLOB_CONTAINER_NAME"));
await blobClient.CreateIfNotExistsAsync();
// Create processorvar processor = new EventProcessorClient(    blobClient,    EventHubConsumerClient.DefaultConsumerGroup,    fullyQualifiedNamespace,    eventHubName,    new DefaultAzureCredential());
// Handle eventsprocessor.ProcessEventAsync += async args =>{    Console.WriteLine($"Partition: {args.Partition.PartitionId}");    Console.WriteLine($"Data: {args.Data.EventBody}");        // Checkpoint after processing (or batch checkpoints)    await args.UpdateCheckpointAsync();};
// Handle errorsprocessor.ProcessErrorAsync += args =>{    Console.WriteLine($"Error: {args.Exception.Message}");    Console.WriteLine($"Partition: {args.PartitionId}");    return Task.CompletedTask;};
// Start processingawait processor.StartProcessingAsync();
// Run until cancelledawait Task.Delay(Timeout.Infinite, cancellationToken);
// Stop gracefullyawait processor.StopProcessingAsync();

4. Partition Operations

csharp
// Get partition IDsstring[] partitionIds = await producer.GetPartitionIdsAsync();
// Send to specific partition (use sparingly)var options = new SendEventOptions{    PartitionId = "0"};await producer.SendAsync(events, options);
// Use partition key (recommended for ordering)var batchOptions = new CreateBatchOptions{    PartitionKey = "customer-123"  // Events with same key go to same partition};using var batch = await producer.CreateBatchAsync(batchOptions);

EventPosition Options

Control where to start reading:

csharp
// Start from beginningEventPosition.Earliest
// Start from end (new events only)EventPosition.Latest
// Start from specific offsetEventPosition.FromOffset(12345)
// Start from specific sequence numberEventPosition.FromSequenceNumber(100)
// Start from specific timeEventPosition.FromEnqueuedTime(DateTimeOffset.UtcNow.AddHours(-1))

ASP.NET Core Integration

csharp
// Program.csusing Azure.Identity;using Azure.Messaging.EventHubs.Producer;using Microsoft.Extensions.Azure;
builder.Services.AddAzureClients(clientBuilder =>{    clientBuilder.AddEventHubProducerClient(        builder.Configuration["EventHub:FullyQualifiedNamespace"],        builder.Configuration["EventHub:Name"]);        clientBuilder.UseCredential(new DefaultAzureCredential());});
// Inject in controller/servicepublic 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);    }}

Best Practices

  1. Use EventProcessorClient for receiving — Never use EventHubConsumerClient in production
  2. Checkpoint strategically — After N events or time interval, not every event
  3. Use partition keys — For ordering guarantees within a partition
  4. Reuse clients — Create once, use as singleton (thread-safe)
  5. Use await using — Ensures proper disposal
  6. Handle ProcessErrorAsync — Always register error handler
  7. Batch events — Use CreateBatchAsync() to respect size limits
  8. Use buffered producer — For high-volume scenarios with automatic batching

Error Handling

csharp
using Azure.Messaging.EventHubs;
try{    await producer.SendAsync(batch);}catch (EventHubsException ex) when (ex.Reason == EventHubsException.FailureReason.ServiceBusy){    // Retry with backoff    await Task.Delay(TimeSpan.FromSeconds(5));}catch (EventHubsException ex) when (ex.IsTransient){    // Transient error - safe to retry    Console.WriteLine($"Transient error: {ex.Message}");}catch (EventHubsException ex){    // Non-transient error    Console.WriteLine($"Error: {ex.Reason} - {ex.Message}");}

Checkpointing Strategies

StrategyWhen to Use
Every eventLow volume, critical data
Every N eventsBalanced throughput/reliability
Time-basedConsistent checkpoint intervals
Batch completionAfter processing a logical batch
csharp
// Checkpoint every 100 eventsprivate int _eventCount = 0;
processor.ProcessEventAsync += async args =>{    // Process event...        _eventCount++;    if (_eventCount >= 100)    {        await args.UpdateCheckpointAsync();        _eventCount = 0;    }};

Related SDKs

SDKPurposeInstall
Azure.Messaging.EventHubsCore sending/receivingdotnet add package Azure.Messaging.EventHubs
Azure.Messaging.EventHubs.ProcessorProduction processingdotnet add package Azure.Messaging.EventHubs.Processor
Azure.ResourceManager.EventHubsManagement plane (create hubs)dotnet add package Azure.ResourceManager.EventHubs
Microsoft.Azure.WebJobs.Extensions.EventHubsAzure Functions bindingdotnet add package Microsoft.Azure.WebJobs.Extensions.EventHubs

来源与署名

来源:microsoft/skills位于.github/plugins/azure-sdk-dotnet/skills/azure-eventhub-dotnet提交354361d

许可证: MIT

内容归原作者所有。SourceWeft 从公开仓库中收录这些内容。

举报或申请下架