【研发类-框架和库Skills】azure-eventhub-java 技能 📅 发布时间:2026/9/1 20:28:39 👁 浏览次数: 使用Azure Event Hubs SDK for Java构建实时流式应用程序。适用于实现事件流式传输、高吞吐量数据摄取或构建事件驱动架构。 下载地址https://github.com/sickn33/antigravity-awesome-skills/tree/main/skills/azure-eventhub-java技能概述azure-eventhub-java 技能是一个专门用于Java开发的Azure Event Hubs SDK技能包。它提供了构建实时流式应用程序所需的所有功能包括事件发送、接收、分区管理和检查点机制。该技能包适用于需要处理大规模事件流、构建事件驱动微服务或实现数据管道的Java开发者。主要功能事件发送支持单个事件和批量事件发送可指定分区或使用分区键事件接收提供EventProcessorClient用于生产环境的可靠事件处理异步客户端支持响应式编程模型的异步生产者和消费者检查点存储集成Blob存储检查点支持分布式消费者负载均衡分区管理支持分区操作和分区键路由事件属性支持自定义事件属性和元数据触发条件在以下情况下应该调用此技能用户需要在Java应用中集成Azure Event Hubs需要实现事件流式传输或高吞吐量数据摄取需要构建事件驱动架构需要使用EventProcessorClient处理事件需要了解分区处理和检查点机制使用场景场景1实时数据流处理构建高吞吐量的实时数据处理管道处理来自IoT设备、应用日志或用户行为的事件流。场景2事件驱动微服务在微服务架构中使用Event Hubs作为事件总线实现服务间的异步通信和解耦。场景3大数据摄取将大量事件数据摄取到数据湖或数据仓库中用于后续分析和处理。处理过程1. 添加Maven依赖在pom.xml中添加必要的依赖dependencygroupIdcom.azure/groupIdartifactIdazure-messaging-eventhubs/artifactIdversion5.19.0/version/dependency!-- 用于检查点存储生产环境 --dependencygroupIdcom.azure/groupIdartifactIdazure-messaging-eventhubs-checkpointstore-blob/artifactIdversion1.20.0/version/dependency2. 创建生产者客户端使用连接字符串或DefaultAzureCredential创建EventHubProducerClientimport com.azure.messaging.eventhubs.EventHubProducerClient;import com.azure.messaging.eventhubs.EventHubClientBuilder;// 使用连接字符串EventHubProducerClient producer new EventHubClientBuilder().connectionString(connection-string, event-hub-name).buildProducerClient();// 使用DefaultAzureCredentialEventHubProducerClient producer new EventHubClientBuilder().fullyQualifiedNamespace(namespace.servicebus.windows.net).eventHubName(event-hub-name).credential(new DefaultAzureCredentialBuilder().build()).buildProducerClient();3. 发送事件批次创建批次并发送事件import com.azure.messaging.eventhubs.EventDataBatch;import com.azure.messaging.eventhubs.models.CreateBatchOptions;EventDataBatch batch producer.createBatch();for (int i 0; i 100; i) {EventData event new EventData(Event i);if (!batch.tryAdd(event)) {producer.send(batch);batch producer.createBatch();batch.tryAdd(event);}}if (batch.getCount() 0) {producer.send(batch);}4. 创建事件处理器使用EventProcessorClient处理事件import com.azure.messaging.eventhubs.EventProcessorClient;import com.azure.messaging.eventhubs.EventProcessorClientBuilder;import com.azure.messaging.eventhubs.checkpointstore.blob.BlobCheckpointStore;BlobContainerAsyncClient blobClient new BlobContainerClientBuilder().connectionString(storage-connection-string).containerName(checkpoints).buildAsyncClient();EventProcessorClient 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());eventContext.updateCheckpoint();}).processError(errorContext - {System.err.println(Error: errorContext.getThrowable().getMessage());}).buildEventProcessorClient();processor.start();输入要求使用此技能时用户需要提供Event Hubs连接信息命名空间、Event Hub名称或完整连接字符串身份认证凭据连接字符串或DefaultAzureCredential消费者组指定消费者组名称默认为$Default存储账户信息用于检查点的Blob存储连接字符串和容器名称事件处理逻辑定义事件处理和错误处理的回调函数输出说明技能将提供完整的Java代码示例包含生产者、消费者和处理器的实现异步客户端示例响应式编程模型的代码示例配置指南环境变量和连接配置的详细说明最佳实践建议关于批次处理、检查点和错误处理的建议客户端类型客户端类型说明EventHubProducerClient同步生产者客户端用于发送事件EventHubConsumerClient同步消费者客户端用于接收事件EventHubProducerAsyncClient异步生产者客户端支持响应式编程EventHubConsumerAsyncClient异步消费者客户端支持响应式编程EventProcessorClient生产环境事件处理器支持检查点和负载均衡事件位置选项方法说明EventPosition.earliest()从最早的事件开始EventPosition.latest()从最新的事件开始仅新事件EventPosition.fromOffset(12345L)从特定偏移量开始EventPosition.fromSequenceNumber(100L)从特定序列号开始EventPosition.fromEnqueuedTime(instant)从特定入队时间开始使用示例示例1发送带分区键的事件CreateBatchOptions options new CreateBatchOptions().setPartitionKey(customer-123);EventDataBatch batch producer.createBatch(options);batch.tryAdd(new EventData(Customer event));producer.send(batch);示例2发送带属性的事件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));示例3批量处理事件EventProcessorClient processor new EventProcessorClientBuilder().connectionString(connection-string, event-hub-name).consumerGroup($Default).checkpointStore(new BlobCheckpointStore(blobClient)).processEventBatch(eventBatchContext - {ListEventData events eventBatchContext.getEvents();System.out.printf(Received %d events%n, events.size());for (EventData event : events) {System.out.println(event.getBodyAsString());}eventBatchContext.updateCheckpoint();}, 50) // maxBatchSize.processError(errorContext - {System.err.println(Error: errorContext.getThrowable());}).buildEventProcessorClient();最佳实践使用EventProcessorClient生产环境中使用EventProcessorClient它提供负载均衡和检查点功能批量发送事件使用EventDataBatch进行高效发送使用分区键在分区内保证事件顺序性检查点处理处理完成后进行检查点避免重复处理错误处理处理瞬态错误并进行重试关闭客户端使用try-with-resources或手动关闭生产者/消费者触发短语Event Hubs Javaevent streaming Azurereal-time data ingestionEventProcessorClientevent hub producer consumerpartition processing注意事项此技能仅适用于任务明确匹配上述范围的情况输出不应替代环境特定的验证、测试或专家审查如果缺少必需的输入、权限、安全边界或成功标准请停止并请求澄清