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 從公開儲存庫中收錄這些內容。

檢舉或申請下架