
簡介
當每個微服務都有自己的資料庫時,最大的挑戰是確保服務之間的資料一致性。本課程深入探討解決雙重寫入、可靠事件發布和最終一致性的模式。
1. 雙寫問題
1.1 問題
當服務需要更新資料庫和發布事件時,如果其中任何一個失敗→不一致:
Scenario 1: DB thành công, Event thất bại
┌──────────────┐
│ Order Service│
│ │
│ 1. Save to DB ✅ (Order created)
│ 2. Publish to Kafka ❌ (Network error)
│ │
│ → Order tồn tại trong DB nhưng không ai biết
└──────────────┘
Scenario 2: Event thành công, DB thất bại
┌──────────────┐
│ Order Service│
│ │
│ 1. Publish to Kafka ✅ (OrderCreated sent)
│ 2. Save to DB ❌ (DB connection failed)
│ │
│ → Downstream services xử lý order không tồn tại
└──────────────┘
1.2 為什麼不使用分散式事務?
// ❌ Anti-pattern: Distributed transaction giữa DB và Kafka
@Transactional
public void createOrder(Order order) {
database.save(order); // Transaction participant 1
kafka.publish(orderCreated); // Transaction participant 2
// Kafka không support XA/2PC → KHÔNG THỂ atomic
}
2.寄件匣模式
2.1 概念
將事件記錄到與業務資料相同的資料庫事務中的寄件箱表。一個單獨的進程(Relay/Poller)讀取發件箱並將其發佈到訊息代理程式:
┌─────────────────────────────────────────────────┐
│ Order Service │
│ │
│ ┌───────────────────────────────────────────┐ │
│ │ Database Transaction (ACID) │ │
│ │ │ │
│ │ INSERT INTO orders (...) VALUES (...); │ │
│ │ INSERT INTO outbox (...) VALUES (...); │ │
│ │ │ │
│ │ → Both succeed or both fail │ │
│ └───────────────────────────────────────────┘ │
│ │
│ ┌───────────────────────────────────────────┐ │
│ │ Outbox Relay (separate process) │ │
│ │ │ │
│ │ 1. Poll outbox table │ │
│ │ 2. Publish event to Kafka │ │
│ │ 3. Mark outbox entry as published │ │
│ └────────────────────┬──────────────────────┘ │
└───────────────────────┼──────────────────────────┘
│
▼
┌─────────┐
│ Kafka │
└─────────┘
2.2 實施
寄件匣表:
CREATE TABLE outbox (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_type VARCHAR(255) NOT NULL, -- "Order"
aggregate_id VARCHAR(255) NOT NULL, -- "O-001"
event_type VARCHAR(255) NOT NULL, -- "OrderCreated"
payload JSONB NOT NULL, -- event data
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
published_at TIMESTAMPTZ, -- NULL = chưa publish
retries INT DEFAULT 0
);
CREATE INDEX idx_outbox_unpublished
ON outbox (created_at) WHERE published_at IS NULL;
業務邏輯:
@Service
public class OrderService {
@Transactional
public Order createOrder(CreateOrderCommand cmd) {
// 1. Save business data
Order order = new Order(cmd);
order.setStatus(OrderStatus.PENDING);
orderRepository.save(order);
// 2. Save outbox event (same transaction!)
OutboxEvent event = OutboxEvent.builder()
.aggregateType("Order")
.aggregateId(order.getId())
.eventType("OrderCreated")
.payload(toJson(new OrderCreatedPayload(order)))
.build();
outboxRepository.save(event);
return order;
}
}
寄件匣中繼(輪詢):
@Scheduled(fixedDelay = 500) // Poll mỗi 500ms
public void publishOutboxEvents() {
List<OutboxEvent> events = outboxRepository
.findUnpublishedOrderByCreatedAt(BATCH_SIZE);
for (OutboxEvent event : events) {
try {
kafkaTemplate.send(
event.getAggregateType().toLowerCase() + ".events",
event.getAggregateId(),
event.getPayload()
).get(); // Wait for Kafka ACK
event.setPublishedAt(Instant.now());
outboxRepository.save(event);
} catch (Exception e) {
event.incrementRetries();
if (event.getRetries() >= MAX_RETRIES) {
event.moveToDeadLetter();
}
outboxRepository.save(event);
}
}
}
2.3 優點和缺點
| 優勢 | 缺點 |
|---|---|
| 保證交付(至少一次) | 輪詢延遲(延遲) |
| 原子業務資料 | 寄件箱表需要清理 |
| 簡單的實作 | 至少一次 → 消費者必須是冪等的 |
| 無需改變資料庫技術 | 新增儲存開銷 |
3. 變更資料擷取 (CDC)
3.1 概念
CDC 從事務日誌(PostgreSQL 中的 WAL)捕獲資料庫變更(INSERT、UPDATE、DELETE)並直接串流到訊息代理程式:
┌──────────────────────────────────────────────┐
│ Order Service │
│ │
│ ┌────────────┐ ┌──────────────────────┐ │
│ │ Application│───▶│ PostgreSQL │ │
│ │ │ │ ┌──────────────────┐ │ │
│ └────────────┘ │ │ orders table │ │ │
│ │ └──────────────────┘ │ │
│ │ ┌──────────────────┐ │ │
│ │ │ outbox table │ │ │
│ │ └──────────────────┘ │ │
│ │ ┌──────────────────┐ │ │
│ │ │ WAL (Write-Ahead │ │ │
│ │ │ Log) │ │ │
│ │ └────────┬─────────┘ │ │
│ └──────────┼───────────┘ │
└───────────────────────────────┼───────────────┘
│
┌───────────▼───────────┐
│ Debezium Connector │
│ (reads WAL stream) │
└───────────┬───────────┘
│
┌───────────▼───────────┐
│ Apache Kafka │
│ topic: outbox.events │
└───────────────────────┘
3.2 Debezium 配置
{
"name": "order-outbox-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "order-db",
"database.port": "5432",
"database.user": "debezium",
"database.password": "${file:/secrets/db-password}",
"database.dbname": "orders",
"database.server.name": "order-service",
"plugin.name": "pgoutput",
"table.include.list": "public.outbox",
"transforms": "outbox",
"transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
"transforms.outbox.table.field.event.id": "id",
"transforms.outbox.table.field.event.key": "aggregate_id",
"transforms.outbox.table.field.event.type": "event_type",
"transforms.outbox.table.field.event.payload": "payload",
"transforms.outbox.route.topic.replacement": "${routedByValue}.events",
"transforms.outbox.table.expand.json.payload": true
}
}
3.3 CDC 與輪詢
| 標準 | 輪詢(寄件匣中繼) | CDC(Debezium) |
|---|---|---|
| 延遲 | 100ms-1s(輪詢間隔) | ~10ms(接近即時) |
| 資料庫載入 | 常見查詢 | 讀取 WAL(影響最小) |
| 基礎設施 | 簡單(無額外組件) | 需要 Debezium + Kafka Connect |
| 複雜性 | 低 | 中 |
| 可靠性 | 好 | 優(WAL = 事實來源) |
| 建議 | 小型/中型 | 大規模、低延遲 |
4.冪等消費者
4.1 為什麼我們需要冪等性?
至少一次傳遞可以發送重複訊息:
Producer ──msg──▶ Broker ──msg──▶ Consumer
│
Process ✅
│
ACK ──▶ Broker (network timeout!)
│
Broker re-delivers msg
│
Consumer receives msg AGAIN
Process AGAIN? → Duplicate!
4.2 實作策略
策略1:冪等性金鑰表
CREATE TABLE processed_events (
event_id UUID PRIMARY KEY,
processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
-- Consumer logic
BEGIN;
-- Check if already processed
INSERT INTO processed_events (event_id) VALUES ($event_id)
ON CONFLICT (event_id) DO NOTHING;
-- If inserted (not duplicate), process the event
IF FOUND THEN
-- Business logic here
UPDATE inventory SET quantity = quantity - $qty WHERE product_id = $pid;
END IF;
COMMIT;
策略2:自然冪等性
// SET operations are naturally idempotent
// "Set order status to CONFIRMED" — same result regardless of how many times
@EventHandler
public void on(PaymentCompletedEvent event) {
Order order = orderRepository.findById(event.getOrderId());
// Idempotent check: if already confirmed, skip
if (order.getStatus() == OrderStatus.CONFIRMED) {
return;
}
order.setStatus(OrderStatus.CONFIRMED);
orderRepository.save(order);
}
策略3:基於版本(樂觀鎖)
UPDATE orders
SET status = 'CONFIRMED', version = version + 1
WHERE id = 'O-001' AND version = 3;
-- If version doesn't match → someone else already updated → skip
-- Rows affected = 0 → duplicate processing detected
5.最終一致性策略
5.1 閱讀你自己寫的內容
確保用戶始終看到他們剛剛更改的數據:
Approach 1: Sticky Session
→ Route user requests đến cùng replica đã nhận write
Approach 2: Client-side Optimistic Update
→ UI cập nhật ngay trước khi server confirm, rollback nếu fail
Approach 3: Write-through Read Model
→ Sau khi write, đồng bộ invalidate read cache
Client ──POST /orders──▶ Order Service ──▶ Write DB
│
Client ──GET /orders/123──▶ Order Service │
↑ │ │
│ Check: event synced? │
│ Yes ──▶ Read DB │
│ No ──▶ Write DB│ (fallback)
└────────────────────────────────────────────┘
5.2 一致性監控
-- Monitor replication lag
SELECT
topic,
partition,
MAX(offset) - MIN(committed_offset) AS lag
FROM consumer_offsets
GROUP BY topic, partition
HAVING lag > 1000; -- Alert if lag > 1000 messages
# Prometheus alert for consistency lag
- alert: HighConsumerLag
expr: kafka_consumer_group_lag > 5000
for: 5m
labels:
severity: warning
annotations:
summary: "Consumer lag is high"
description: "Consumer group {{ $labels.group }} has lag {{ $value }}"
6. 端對端模式
6.1 推薦架構
┌─────────────────────────────────────────────────────────────┐
│ End-to-End Data Flow │
│ │
│ ┌──────────┐ ┌──────────────────────┐ ┌───────────┐ │
│ │ Service │ │ PostgreSQL │ │ Debezium │ │
│ │ Logic │───▶│ ┌────────┐ ┌───────┐ │───▶│ CDC │ │
│ │ │ │ │Business│ │Outbox │ │ │ Connector │ │
│ └──────────┘ │ │ Table │ │ Table │ │ └─────┬─────┘ │
│ │ └────────┘ └───────┘ │ │ │
│ └──────────────────────┘ │ │
│ ▼ │
│ ┌──────────┐ │
│ │ Kafka │ │
│ │ Topic │ │
│ └────┬─────┘ │
│ │ │
│ ┌────────────────────┼────┐ │
│ │ │ │ │
│ ┌─────▼────┐ ┌─────▼────┐ │
│ │ Consumer │ │ Consumer │ │
│ │ Idempot. │ │ Idempot. │ │
│ │ Check ✓ │ │ Check ✓ │ │
│ └──────────┘ └──────────┘ │
└─────────────────────────────────────────────────────────────┘
Guarantees:
✅ Atomic write (business data + outbox in same TX)
✅ At-least-once delivery (CDC from WAL)
✅ Exactly-once processing (idempotent consumers)
✅ Near real-time (<50ms latency)
總結
- 雙寫問題是微服務中資料不一致的根本原因
- 寄件匣模式透過將事件寫入相同資料庫交易來解決
- CDC (Debezium) 讀取 WAL 流以實現近乎即時的事件傳遞
- 冪等消費者確保精確一次處理,儘管至少一次交付
- 最終一致性是不可避免的——針對它而不是反對它進行設計
- 生產需要監控一致性滯後