Azure Event Grid Java SDK

azure-eventgrid-java
分类编程
作者Agentic Awesome Skills 社区
许可MIT
评分4.70/5
使用16.5K

Azure Event Grid Java SDK

使用 Azure Event Grid Java SDK 构建事件驱动应用程序。

安装

xml
<dependency>
    <groupId>com.azure</groupId>
    <artifactId>azure-messaging-eventgrid</artifactId>
    <version>4.27.0</version>
</dependency>

客户端创建

EventGridPublisherClient

java
import com.azure.messaging.eventgrid.EventGridPublisherClient;
import com.azure.messaging.eventgrid.EventGridPublisherClientBuilder;
import com.azure.core.credential.AzureKeyCredential;

// 使用 API 密钥
EventGridPublisherClient<EventGridEvent> client = new EventGridPublisherClientBuilder()
.endpoint("<topic-endpoint>")
.credential(new AzureKeyCredential("<access-key>"))
.buildEventGridEventPublisherClient();

// 用于 CloudEvents
EventGridPublisherClient<CloudEvent> cloudClient = new EventGridPublisherClientBuilder()
.endpoint("<topic-endpoint>")
.credential(new AzureKeyCredential("<access-key>"))
.buildCloudEventPublisherClient();

使用 DefaultAzureCredential

java
import com.azure.identity.DefaultAzureCredentialBuilder;

EventGridPublisherClient<EventGridEvent> client = new EventGridPublisherClientBuilder()
.endpoint("<topic-endpoint>")
.credential(new DefaultAzureCredentialBuilder().build())
.buildEventGridEventPublisherClient();

异步客户端

java
import com.azure.messaging.eventgrid.EventGridPublisherAsyncClient;

EventGridPublisherAsyncClient<EventGridEvent> asyncClient = new EventGridPublisherClientBuilder()
.endpoint("<topic-endpoint>")
.credential(new AzureKeyCredential("<access-key>"))
.buildEventGridEventPublisherAsyncClient();

事件类型

| 类型 | 描述 |
|------|-------------|
| EventGridEvent | Azure Event Grid 原生架构 |
| CloudEvent | CNCF CloudEvents 1.0 规范 |
| BinaryData | 自定义架构事件 |

核心模式

发布 EventGridEvent

java
import com.azure.messaging.eventgrid.EventGridEvent;
import com.azure.core.util.BinaryData;

EventGridEvent event = new EventGridEvent(
"resource/path", // subject
"MyApp.Events.OrderCreated", // eventType
BinaryData.fromObject(new OrderData("order-123", 99.99)), // data
"1.0" // dataVersion
);

client.sendEvent(event);

发布多个事件

java
List<EventGridEvent> events = Arrays.asList(
    new EventGridEvent("orders/1", "Order.Created", 
        BinaryData.fromObject(order1), "1.0"),
    new EventGridEvent("orders/2", "Order.Created", 
        BinaryData.fromObject(order2), "1.0")
);

client.sendEvents(events);

发布 CloudEvent

java
import com.azure.core.models.CloudEvent;
import com.azure.core.models.CloudEventDataFormat;

CloudEvent cloudEvent = new CloudEvent(
"/myapp/orders", // source
"order.created", // type
BinaryData.fromObject(orderData), // data
CloudEventDataFormat.JSON // dataFormat
);
cloudEvent.setSubject("orders/12345");
cloudEvent.setId(UUID.randomUUID().toString());

cloudClient.sendEvent(cloudEvent);

发布 CloudEvents 批处理

java
List<CloudEvent> cloudEvents = Arrays.asList(
    new CloudEvent("/app", "event.typ
e1", BinaryData.fromString("data1"), CloudEventDataFormat.JSON), new CloudEvent("/app", "event.type2", BinaryData.fromString("data2"), CloudEventDataFormat.JSON) );

cloudClient.sendEvents(cloudEvents);

code
### 异步发布
java
asyncClient.sendEvent(event)
.subscribe(
unused -> System.out.println("事件发送成功"),
error -> System.err.println("错误: " + error.getMessage())
);

// 发送多个事件
asyncClient.sendEvents(events)
.doOnSuccess(unused -> System.out.println("所有事件已发送"))
.doOnError(error -> System.err.println("失败: " + error))
.block(); // 如有需要则阻塞

code
### 自定义事件数据类
java
public class OrderData {
private String orderId;
private double amount;
private String customerId;

public OrderData(String orderId, double amount) {
this.orderId = orderId;
this.amount = amount;
}

// Getters 和 setters
}

// 使用示例
OrderData order = new OrderData("ORD-123", 150.00);
EventGridEvent event = new EventGridEvent(
"orders/" + order.getOrderId(),
"MyApp.Order.Created",
BinaryData.fromObject(order),
"1.0"
);

code
## 接收事件

解析 EventGridEvent

java import com.azure.messaging.eventgrid.EventGridEvent;

// 从 JSON 字符串解析(例如 webhook 负载)
String jsonPayload = "[{\"id\": \"...\", ...}]";
List<EventGridEvent> events = EventGridEvent.fromString(jsonPayload);

for (EventGridEvent event : events) {
System.out.println("事件类型: " + event.getEventType());
System.out.println("主题: " + event.getSubject());
System.out.println("事件时间: " + event.getEventTime());

// 获取数据
BinaryData data = event.getData();
OrderData orderData = data.toObject(OrderData.class);
}

code
### 解析 CloudEvent
java
import com.azure.core.models.CloudEvent;

String cloudEventJson = "[{\"specversion\": \"1.0\", ...}]";
List<CloudEvent> cloudEvents = CloudEvent.fromString(cloudEventJson);

for (CloudEvent event : cloudEvents) {
System.out.println("类型: " + event.getType());
System.out.println("来源: " + event.getSource());
System.out.println("ID: " + event.getId());

MyEventData data = event.getData().toObject(MyEventData.class);
}

code
### 处理系统事件
java
import com.azure.messaging.eventgrid.systemevents.*;

for (EventGridEvent event : events) {
if (event.getEventType().equals("Microsoft.Storage.BlobCreated")) {
StorageBlobCreatedEventData blobData =
event.getData().toObject(StorageBlobCreatedEventData.class);
System.out.println("Blob URL: " + blobData.getUrl());
}
}

code
## Event Grid 命名空间 (MQTT/Pull)

从命名空间主题接收

java import com.azure.messaging.eventgrid.namespaces.EventGridReceiverClient; import com.azure.messaging.eventgrid.namespaces.EventGridReceiverClientBuilder; import com.azure.messaging.eventgrid.namespaces.models.*;

EventGridReceiverClient receiverClient = new EventGridReceiverClientBuilder()
.endpoint("<namespace-endpoint>")
.credential(new AzureKeyCredential("<key>"))
.topicName("my-topic")
.subscriptionName("my-subscription")
.buildClient();

// 接收事件
ReceiveResult result = receiverClient.receive(10, Duration.ofSeconds(30));

for (ReceiveDetails detail : result.getValue()) {
CloudEvent event = detail.getEvent();
System.out.println("事件: " + event.getType());

// 确认事件
rece
iverClient.acknowledge(Arrays.asList(detail.getBrokerProperties().getLockToken()));
}

code
### 拒绝或释放事件
java
// 拒绝(不再重试)
receiverClient.reject(Arrays.asList(lockToken));

// 释放(稍后重试)
receiverClient.release(Arrays.asList(lockToken));

// 延迟释放
receiverClient.release(Arrays.asList(lockToken),
new ReleaseOptions().setDelay(ReleaseDelay.BY_60_SECONDS));

code
## 错误处理
java
import com.azure.core.exception.HttpResponseException;

try {
client.sendEvent(event);
} catch (HttpResponseException e) {
System.out.println("Status: " + e.getResponse().getStatusCode());
System.out.println("Error: " + e.getMessage());
}

code
## 环境变量
bash
EVENT_GRID_TOPIC_ENDPOINT=https://<topic-name>.<region>.eventgrid.azure.net/api/events
EVENT_GRID_ACCESS_KEY=<your-access-key>
```

最佳实践

1. 批量发送事件:尽可能在一次调用中发送多个事件。
2. 幂等性:包含唯一的事件 ID 以进行去重。
3. 架构验证:使用强类型事件数据类。
4. 重试逻辑:SDK 已内置重试,但对于持续失败的情况请考虑使用死信队列(dead-letter)。
5. 事件大小:将事件控制在 1MB 以内(基础层为 64KB)。

触发词

  • "Event Grid Java"
  • "publish events Azure"
  • "CloudEvent SDK"
  • "event-driven messaging"
  • "pub/sub Azure"
  • "webhook events"

适用场景

本技能适用于执行概览中所描述的工作流或操作。

局限性

  • 仅在任务明确符合上述范围时使用此技能。
  • 不要将输出结果视为针对特定环境的验证、测试或专家评审的替代方案。
  • 如果缺少必要的输入、权限、安全边界或成功标准,请停止操作并请求澄清。