簡介
事件驅動架構(EDA)解決了服務之間的緊密耦合問題。服務不是直接呼叫它,而是將事件發佈到 Apache Kafka - 其他服務訂閱並做出反應。 Quarkus 將 SmallRye Reactive Messaging 與 Kafka 連接器結合使用。
何時使用事件而不是 REST/gRPC?
| 圖案 | 使用案例 | 電子商務範例 |
|---|---|---|
| 休息 | 請求-回應,立即需要結果 | 查看庫存 |
| gRPC | 高效能,內部 | 批次庫存查詢 |
| 卡夫卡事件 | 即發即忘、解耦 | 訂單已建立 → 傳送電子郵件 |
設定 Kafka 開發服務
依賴關係
<dependency>
<groupId>io.quarkus</groupId>
<artifactId>quarkus-messaging-kafka</artifactId>
</dependency>
零配置
開發服務可自動啟動 Kafka 容器(Redpanda):
# Không cần cấu hình gì cho dev mode!
# Kafka broker tự start, tạo topics tự động
# Production config
%prod.kafka.bootstrap.servers=${KAFKA_BOOTSTRAP_SERVERS}
領域事件
定義事件
// Shared event schemas (có thể đặt trong common module)
public sealed interface OrderEvent {
String orderId();
String orderNumber();
Instant timestamp();
record OrderCreated(
String orderId,
String orderNumber,
String customerId,
String customerEmail,
List<OrderItemInfo> items,
BigDecimal totalAmount,
Instant timestamp
) implements OrderEvent {}
record OrderConfirmed(
String orderId,
String orderNumber,
Instant timestamp
) implements OrderEvent {}
record OrderPaid(
String orderId,
String orderNumber,
String paymentId,
BigDecimal amount,
Instant timestamp
) implements OrderEvent {}
record OrderCancelled(
String orderId,
String orderNumber,
String reason,
Instant timestamp
) implements OrderEvent {}
record OrderItemInfo(
Long productId,
String productName,
int quantity,
BigDecimal unitPrice
) {}
}
public sealed interface PaymentEvent {
record PaymentCompleted(
String paymentId,
String orderId,
BigDecimal amount,
String method,
Instant timestamp
) implements PaymentEvent {}
record PaymentFailed(
String paymentId,
String orderId,
String reason,
Instant timestamp
) implements PaymentEvent {}
}
事件製作者 — 訂單服務
import org.eclipse.microprofile.reactive.messaging.Channel;
import org.eclipse.microprofile.reactive.messaging.Emitter;
import io.smallrye.reactive.messaging.kafka.Record;
@ApplicationScoped
public class OrderEventPublisher {
@Inject
@Channel("order-events-out")
Emitter<OrderEvent> orderEventEmitter;
public void publishOrderCreated(Order order) {
OrderEvent.OrderCreated event =
new OrderEvent.OrderCreated(
order.id.toString(),
order.orderNumber,
order.customerId,
order.customerEmail,
order.items.stream()
.map(i -> new OrderEvent.OrderItemInfo(
i.productId, i.productName,
i.quantity,
i.unitPrice.amount()))
.toList(),
order.totalAmount.amount(),
Instant.now());
orderEventEmitter.send(event);
Log.infof("Published OrderCreated: %s",
order.orderNumber);
}
public void publishOrderPaid(Order order) {
orderEventEmitter.send(
new OrderEvent.OrderPaid(
order.id.toString(),
order.orderNumber,
order.paymentId,
order.totalAmount.amount(),
Instant.now()));
}
public void publishOrderCancelled(Order order,
String reason) {
orderEventEmitter.send(
new OrderEvent.OrderCancelled(
order.id.toString(),
order.orderNumber,
reason,
Instant.now()));
}
}
Kafka 設定-生產者
# application.properties (Order Service)
mp.messaging.outgoing.order-events-out.connector=smallrye-kafka
mp.messaging.outgoing.order-events-out.topic=order-events
mp.messaging.outgoing.order-events-out.value.serializer=io.quarkus.kafka.client.serialization.ObjectMapperSerializer
mp.messaging.outgoing.order-events-out.key.serializer=org.apache.kafka.common.serialization.StringSerializer
事件消費者-通知服務
import org.eclipse.microprofile.reactive.messaging.Incoming;
import io.smallrye.common.annotation.Blocking;
@ApplicationScoped
public class OrderEventConsumer {
@Inject
NotificationTemplateService templateService;
@Inject
EmailNotificationSender emailSender;
@Incoming("order-events-in")
@Blocking // Nếu xử lý blocking (DB, email)
public void onOrderEvent(OrderEvent event) {
Log.infof("Received order event: %s",
event.getClass().getSimpleName());
switch (event) {
case OrderEvent.OrderCreated created -> {
Notification notification =
templateService.createOrderConfirmation(
created.customerEmail(),
created.orderNumber(),
created.totalAmount());
notification.persist();
emailSender.send(notification);
}
case OrderEvent.OrderPaid paid -> {
// Gửi email xác nhận thanh toán
Log.infof("Order %s paid with %s",
paid.orderNumber(), paid.paymentId());
}
case OrderEvent.OrderCancelled cancelled -> {
// Gửi email thông báo hủy
Log.infof("Order %s cancelled: %s",
cancelled.orderNumber(),
cancelled.reason());
}
default -> Log.warnf(
"Unknown order event: %s", event);
}
}
}
Kafka 設定 — 消費者
# application.properties (Notification Service)
mp.messaging.incoming.order-events-in.connector=smallrye-kafka
mp.messaging.incoming.order-events-in.topic=order-events
mp.messaging.incoming.order-events-in.group.id=notification-service
mp.messaging.incoming.order-events-in.auto.offset.reset=earliest
mp.messaging.incoming.order-events-in.value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
mp.messaging.incoming.order-events-in.failure-strategy=dead-letter-queue
自訂解串器
import io.quarkus.kafka.client.serialization.ObjectMapperDeserializer;
public class OrderEventDeserializer
extends ObjectMapperDeserializer<OrderEvent> {
public OrderEventDeserializer() {
super(OrderEvent.class);
}
}
產品服務 — 庫存消費者
@ApplicationScoped
public class StockEventConsumer {
@Inject
ProductRepository productRepo;
@Incoming("order-events-for-stock")
@Blocking
@Transactional
public void onOrderEvent(OrderEvent event) {
switch (event) {
case OrderEvent.OrderConfirmed confirmed -> {
// Stock đã reserved qua gRPC khi confirm
Log.infof("Order %s confirmed, stock reserved",
confirmed.orderNumber());
}
case OrderEvent.OrderCancelled cancelled -> {
// Release stock
// (cần OrderItems info - hoặc query Order Service)
Log.infof("Order %s cancelled, releasing stock",
cancelled.orderNumber());
}
default -> {} // ignore
}
}
}
寄件匣模式-保證原子性
問題:將訂單儲存到資料庫 + 發布事件 - 需要兩者都成功(或兩者都失敗)。
// Outbox table
@Entity
@Table(name = "outbox_events")
public class OutboxEvent extends PanacheEntity {
@Column(name = "aggregate_type", length = 50)
public String aggregateType; // "Order"
@Column(name = "aggregate_id", length = 100)
public String aggregateId;
@Column(name = "event_type", length = 100)
public String eventType; // "OrderCreated"
@Column(columnDefinition = "TEXT")
public String payload; // JSON
@Column(name = "created_at")
public LocalDateTime createdAt;
@Column(name = "published")
public boolean published = false;
@PrePersist
void onCreate() { createdAt = LocalDateTime.now(); }
}
// Trong Order Service
@Transactional
public OrderDTO createOrder(CreateOrderRequest request) {
Order order = new Order();
// ... build order ...
order.persist();
// Cùng transaction: lưu event vào outbox
OutboxEvent outbox = new OutboxEvent();
outbox.aggregateType = "Order";
outbox.aggregateId = order.id.toString();
outbox.eventType = "OrderCreated";
outbox.payload = objectMapper.writeValueAsString(
OrderEvent.OrderCreated.from(order));
outbox.persist();
return OrderDTO.from(order);
// Transaction commit → cả order + outbox event đều saved
}
// Scheduled poller: đọc outbox → publish to Kafka
@ApplicationScoped
public class OutboxPoller {
@Inject @Channel("order-events-out")
Emitter<String> emitter;
@Scheduled(every = "5s")
@Transactional
void pollAndPublish() {
List<OutboxEvent> pending = OutboxEvent
.list("published = false ORDER BY createdAt ASC");
for (OutboxEvent event : pending) {
try {
emitter.send(event.payload);
event.published = true;
} catch (Exception e) {
Log.errorf("Failed to publish outbox: %s",
event.id);
}
}
}
}
死信隊列
# Khi consumer fail xử lý message
mp.messaging.incoming.order-events-in.failure-strategy=dead-letter-queue
mp.messaging.incoming.order-events-in.dead-letter-queue.topic=order-events-dlq
mp.messaging.incoming.order-events-in.dead-letter-queue.value.serializer=org.apache.kafka.common.serialization.StringSerializer
練習
- 定義
OrderEvent具有事件類型的密封接口 - 創建
OrderEventPublisher訂單服務 3.創建OrderEventConsumer在通知服務中 4.配置Kafka通道(生產者+消費者) - 實現發件箱模式確保原子性
- 測試:建立訂單 → 驗證通知服務接收事件 → 發送電子郵件
總結
- SmallRye 反應式訊息傳送 —
@Incoming(消耗),@Channel+Emitter(生產) - 開發服務自動啟動 Kafka (Redpanda) — 零配置
- 密封介面用於型別安全的網域事件
@Blocking適用於需要阻塞 I/O(資料庫、電子郵件)的消費者- 寄件匣模式確保資料庫寫入+事件發布原子
- 失敗訊息的死信佇列 - 沒有遺失事件
下一篇:容錯 - 斷路器、重試、隔板、超時。