Chuyển đến nội dung chính

第 11 課:資料一致性模式 — 寄件匣、CDC 與最終一致性

CAP 定理實踐、最終一致性、發件箱模式、Debezium 的變更資料擷取 (CDC)、冪等消費者和端到端資料一致性策略。

🏗️ 建築 — 第 11 課 第 11 課:資料一致性模式 — 寄件匣、CDC 和最終一致性

雲端原生微服務架構

第 3 部分:微服務中的資料管理

亞洲開發網

第 11 課:資料一致性模式 — 寄件匣、CDC 與最終一致性

簡介

當每個微服務都有自己的資料庫時,最大的挑戰是確保服務之間的資料一致性。本課程深入探討解決雙重寫入、可靠事件發布和最終一致性的模式。


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 流以實現近乎即時的事件傳遞
  • 冪等消費者確保精確一次處理,儘管至少一次交付
  • 最終一致性是不可避免的——針對它而不是反對它進行設計
  • 生產需要監控一致性滯後