Azure Eventhub Java

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

Build real-time streaming applications with Azure Event Hubs SDK for Java. Use when implementing event streaming, high-throughput data ingestion, or building event-driven architectures.

精选仅含说明Software Development
AI 生成的概览

指导 Java 开发者构建 Azure 事件中心流式生产者、消费者和事件处理器。

功能
该技能提供使用 Azure 事件中心 Java SDK 的参考说明和代码模式。内容涵盖使用连接字符串或 Entra ID 凭据创建客户端、发送单个事件和批次、分区键、接收事件、带 Blob 检查点的 EventProcessorClient、事件位置、错误处理和资源清理。它产出的是指导和代码片段,而非可直接运行的产物。
适用场景
适用于在 Java 中实现 Azure 事件中心上的实时事件流、高吞吐量数据摄取或事件驱动架构。适合需要生产者、消费者或生产级事件处理器具体 SDK 用法模式的开发者。
运行要求
需要 Java 以及 com.azure:azure-messaging-eventhubs Maven 依赖,检查点功能可选配 azure-messaging-eventhubs-checkpointstore-blob 和 Azure 存储 Blob。需要 Azure 事件中心命名空间以及连接字符串或 Entra ID 凭据,并需要访问 Azure 的网络。不包含脚本,仅为说明文档。

Azure Event Hubs SDK for Java

Build real-time streaming applications using the Azure Event Hubs SDK for Java.

Installation

xml
<dependency>    <groupId>com.azure</groupId>    <artifactId>azure-messaging-eventhubs</artifactId>    <version>5.19.0</version></dependency>
<!-- For checkpoint store (production) --><dependency>    <groupId>com.azure</groupId>    <artifactId>azure-messaging-eventhubs-checkpointstore-blob</artifactId>    <version>1.20.0</version></dependency>

Client Creation

EventHubProducerClient

java
import com.azure.messaging.eventhubs.EventHubProducerClient;import com.azure.messaging.eventhubs.EventHubClientBuilder;
// With connection stringEventHubProducerClient producer = new EventHubClientBuilder()    .connectionString("<connection-string>", "<event-hub-name>")    .buildProducerClient();
// Full connection string with EntityPathEventHubProducerClient producer = new EventHubClientBuilder()    .connectionString("<connection-string-with-entity-path>")    .buildProducerClient();

With DefaultAzureCredential

java
import com.azure.core.credential.TokenCredential;import com.azure.identity.AzureIdentityEnvVars;import com.azure.identity.DefaultAzureCredentialBuilder;import com.azure.identity.ManagedIdentityCredentialBuilder;
// Local dev: DefaultAzureCredential. Production: set AZURE_TOKEN_CREDENTIALS=prod or AZURE_TOKEN_CREDENTIALS=<specific_credential>TokenCredential credential = new DefaultAzureCredentialBuilder()    .requireEnvVars(AzureIdentityEnvVars.AZURE_TOKEN_CREDENTIALS)    .build();// Or use a specific credential directly in production:// See https://learn.microsoft.com/java/api/overview/azure/identity-readme?view=azure-java-stable#credential-classes// TokenCredential credential = new ManagedIdentityCredentialBuilder().build();
EventHubProducerClient producer = new EventHubClientBuilder()    .fullyQualifiedNamespace("<namespace>.servicebus.windows.net")    .eventHubName("<event-hub-name>")    .credential(credential)    .buildProducerClient();

EventHubConsumerClient

java
import com.azure.messaging.eventhubs.EventHubConsumerClient;
EventHubConsumerClient consumer = new EventHubClientBuilder()    .connectionString("<connection-string>", "<event-hub-name>")    .consumerGroup(EventHubClientBuilder.DEFAULT_CONSUMER_GROUP_NAME)    .buildConsumerClient();

Async Clients

java
import com.azure.messaging.eventhubs.EventHubProducerAsyncClient;import com.azure.messaging.eventhubs.EventHubConsumerAsyncClient;
EventHubProducerAsyncClient asyncProducer = new EventHubClientBuilder()    .connectionString("<connection-string>", "<event-hub-name>")    .buildAsyncProducerClient();
EventHubConsumerAsyncClient asyncConsumer = new EventHubClientBuilder()    .connectionString("<connection-string>", "<event-hub-name>")    .consumerGroup("$Default")    .buildAsyncConsumerClient();

Core Patterns

Send Single Event

java
import com.azure.messaging.eventhubs.EventData;
EventData eventData = new EventData("Hello, Event Hubs!");producer.send(Collections.singletonList(eventData));

Send Event Batch

java
import com.azure.messaging.eventhubs.EventDataBatch;import com.azure.messaging.eventhubs.models.CreateBatchOptions;
// Create batchEventDataBatch batch = producer.createBatch();
// Add events (returns false if batch is full)for (int i = 0; i < 100; i++) {    EventData event = new EventData("Event " + i);    if (!batch.tryAdd(event)) {        // Batch is full, send and create new batch        producer.send(batch);        batch = producer.createBatch();        batch.tryAdd(event);    }}
// Send remaining eventsif (batch.getCount() > 0) {    producer.send(batch);}

Send to Specific Partition

java
CreateBatchOptions options = new CreateBatchOptions()    .setPartitionId("0");
EventDataBatch batch = producer.createBatch(options);batch.tryAdd(new EventData("Partition 0 event"));producer.send(batch);

Send with Partition Key

java
CreateBatchOptions options = new CreateBatchOptions()    .setPartitionKey("customer-123");
EventDataBatch batch = producer.createBatch(options);batch.tryAdd(new EventData("Customer event"));producer.send(batch);

Event with Properties

java
EventData event = new EventData("Order created");event.getProperties().put("orderId", "ORD-123");event.getProperties().put("customerId", "CUST-456");event.getProperties().put("priority", 1);
producer.send(Collections.singletonList(event));

Receive Events (Simple)

java
import com.azure.messaging.eventhubs.models.EventPosition;import com.azure.messaging.eventhubs.models.PartitionEvent;
// Receive from specific partitionIterable<PartitionEvent> events = consumer.receiveFromPartition(    "0",                           // partitionId    10,                            // maxEvents    EventPosition.earliest(),      // startingPosition    Duration.ofSeconds(30)         // timeout);
for (PartitionEvent partitionEvent : events) {    EventData event = partitionEvent.getData();    System.out.println("Body: " + event.getBodyAsString());    System.out.println("Sequence: " + event.getSequenceNumber());    System.out.println("Offset: " + event.getOffset());}

EventProcessorClient (Production)

java
import com.azure.messaging.eventhubs.EventProcessorClient;import com.azure.messaging.eventhubs.EventProcessorClientBuilder;import com.azure.messaging.eventhubs.checkpointstore.blob.BlobCheckpointStore;import com.azure.storage.blob.BlobContainerAsyncClient;import com.azure.storage.blob.BlobContainerClientBuilder;
// Create checkpoint storeBlobContainerAsyncClient blobClient = new BlobContainerClientBuilder()    .connectionString("<storage-connection-string>")    .containerName("checkpoints")    .buildAsyncClient();
// Create processorEventProcessorClient processor = new EventProcessorClientBuilder()    .connectionString("<eventhub-connection-string>", "<event-hub-name>")    .consumerGroup("$Default")    .checkpointStore(new BlobCheckpointStore(blobClient))    .processEvent(eventContext -> {        EventData event = eventContext.getEventData();        System.out.println("Processing: " + event.getBodyAsString());                // Checkpoint after processing        eventContext.updateCheckpoint();    })    .processError(errorContext -> {        System.err.println("Error: " + errorContext.getThrowable().getMessage());        System.err.println("Partition: " + errorContext.getPartitionContext().getPartitionId());    })    .buildEventProcessorClient();
// Start processingprocessor.start();
// Keep running...Thread.sleep(Duration.ofMinutes(5).toMillis());
// Stop gracefullyprocessor.stop();

Batch Processing

java
EventProcessorClient processor = new EventProcessorClientBuilder()    .connectionString("<connection-string>", "<event-hub-name>")    .consumerGroup("$Default")    .checkpointStore(new BlobCheckpointStore(blobClient))    .processEventBatch(eventBatchContext -> {        List<EventData> events = eventBatchContext.getEvents();        System.out.printf("Received %d events%n", events.size());                for (EventData event : events) {            // Process each event            System.out.println(event.getBodyAsString());        }                // Checkpoint after batch        eventBatchContext.updateCheckpoint();    }, 50) // maxBatchSize    .processError(errorContext -> {        System.err.println("Error: " + errorContext.getThrowable());    })    .buildEventProcessorClient();

Async Receiving

java
asyncConsumer.receiveFromPartition("0", EventPosition.latest())    .subscribe(        partitionEvent -> {            EventData event = partitionEvent.getData();            System.out.println("Received: " + event.getBodyAsString());        },        error -> System.err.println("Error: " + error),        () -> System.out.println("Complete")    );

Get Event Hub Properties

java
// Get hub infoEventHubProperties hubProps = producer.getEventHubProperties();System.out.println("Hub: " + hubProps.getName());System.out.println("Partitions: " + hubProps.getPartitionIds());
// Get partition infoPartitionProperties partitionProps = producer.getPartitionProperties("0");System.out.println("Begin sequence: " + partitionProps.getBeginningSequenceNumber());System.out.println("Last sequence: " + partitionProps.getLastEnqueuedSequenceNumber());System.out.println("Last offset: " + partitionProps.getLastEnqueuedOffset());

Event Positions

java
// Start from beginningEventPosition.earliest()
// Start from end (new events only)EventPosition.latest()
// From specific offsetEventPosition.fromOffset(12345L)
// From specific sequence numberEventPosition.fromSequenceNumber(100L)
// From specific timeEventPosition.fromEnqueuedTime(Instant.now().minus(Duration.ofHours(1)))

Error Handling

java
import com.azure.messaging.eventhubs.models.ErrorContext;
.processError(errorContext -> {    Throwable error = errorContext.getThrowable();    String partitionId = errorContext.getPartitionContext().getPartitionId();        if (error instanceof AmqpException) {        AmqpException amqpError = (AmqpException) error;        if (amqpError.isTransient()) {            System.out.println("Transient error, will retry");        }    }        System.err.printf("Error on partition %s: %s%n", partitionId, error.getMessage());})

Resource Cleanup

java
// Always close clientstry {    producer.send(batch);} finally {    producer.close();}
// Or use try-with-resourcestry (EventHubProducerClient producer = new EventHubClientBuilder()        .connectionString(connectionString, eventHubName)        .buildProducerClient()) {    producer.send(events);}

Environment Variables

bash
EVENT_HUBS_CONNECTION_STRING=Endpoint=sb://<namespace>.servicebus.windows.net/;SharedAccessKeyName=...  # Alternative to Entra ID authEVENT_HUBS_NAME=<event-hub-name>  # Required for event hub nameSTORAGE_CONNECTION_STRING=<for-checkpointing>  # Alternative to Entra ID auth for checkpointingAZURE_TOKEN_CREDENTIALS=prod  # Required only if DefaultAzureCredential is used in production

Best Practices

  1. Use EventProcessorClient: For production, provides load balancing and checkpointing
  2. Batch Events: Use EventDataBatch for efficient sending
  3. Partition Keys: Use for ordering guarantees within a partition
  4. Checkpointing: Checkpoint after processing to avoid reprocessing
  5. Error Handling: Handle transient errors with retries
  6. Close Clients: Always close producer/consumer when done

Trigger Phrases

  • "Event Hubs Java"
  • "event streaming Azure"
  • "real-time data ingestion"
  • "EventProcessorClient"
  • "event hub producer consumer"
  • "partition processing"

来源与署名

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

许可证: MIT

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

举报或申请下架