Streaming 7 — CDC Streaming & Event-driven

22 thg 7, 2026 2 lượt xem
#event-driven
#cdc
#data-engineering
#streaming
#debezium

Mô hình tinh thần: database là một cái log đang giả vờ làm bảng

Mọi database giao dịch (OLTP) đều ghi một transaction log trước khi sửa dữ liệu: PostgreSQL gọi là WAL (Write-Ahead Log), MySQL gọi là binlog, Oracle gọi là redo log, SQL Server có transaction log. Đây không phải chi tiết phụ — đó là nguồn sự thật thật sự của database. Bảng bạn nhìn thấy chỉ là materialized view của việc phát lại (replay) toàn bộ log đó. Chính cơ chế này giúp database phục hồi sau crash và cho replica đồng bộ với primary.

Change Data Capture (CDC) dựa trên một nhận thức đơn giản: nếu log đã chứa mọi INSERT/UPDATE/DELETE theo đúng thứ tự commit, thì thay vì hỏi bảng "có gì mới không?", ta chỉ cần đọc cái log đó và phát mỗi thay đổi ra ngoài như một sự kiện. Đây là cầu nối giữa thế giới database at rest và thế giới stream in motion — và nó khớp trực tiếp với ý tưởng stream–table duality: changelog của một bảng chính là một stream.

CDC là cách phổ biến nhất hiện nay để đưa dữ liệu OLTP vào lakehouse/analytics gần thời gian thựckhông tải nặng lên database nguồn — điều kiện then chốt để bộ phận vận hành core banking chấp thuận.

Query-based CDC vs Log-based CDC

Có hai họ CDC, và sự khác biệt quyết định toàn bộ chất lượng pipeline.

Tiêu chíQuery-based (polling)Log-based
Cách hoạt độngSELECT ... WHERE updated_at > :last theo lịchĐọc WAL/binlog/redo trực tiếp
Tải lên nguồnCao (quét bảng lặp lại, tranh chấp lock)Rất thấp (đọc log tuần tự)
Bắt được DELETE?Không (hàng biến mất, query không thấy) (log ghi cả delete)
Bắt trạng thái trung gian?Không (chỉ thấy giá trị cuối mỗi lần poll) (mọi lần update đều có bản ghi)
Độ trễBằng chu kỳ poll (giây → phút)Gần tức thời (mili-giây → giây)
Yêu cầu cộtCần cột updated_at đáng tinKhông cần thay đổi schema

Query-based là kiểu CDC mô tả trong batch incremental: dùng high-water-mark trên updated_at. Nó đơn giản nhưng mù với DELETE và các update bị "nuốt" giữa hai lần poll. Log-based khắc phục cả hai và là chuẩn cho streaming CDC — đổi lại phải có quyền đọc log và cấu hình database (bật WAL logical, binlog ROW format, supplemental logging cho Oracle).

Debezium: log-based CDC thành sự kiện Kafka

Debezium là nền tảng CDC mã nguồn mở chạy như connector nguồn trong Kafka Connect. Với mỗi loại DB (PostgreSQL, MySQL, Oracle, SQL Server, MongoDB) có một connector riêng biết cách đọc log của DB đó.

Cách Debezium làm việc, hai giai đoạn:

  1. Snapshot ban đầu: lần đầu chạy, connector chụp trạng thái hiện tại của bảng (một loạt event op = "r" — read) để consumer có nền đầy đủ, không chỉ các thay đổi tương lai.
  2. Streaming: sau snapshot, nó chuyển sang đọc log liên tục, phát mỗi thay đổi ra topic. Vị trí đã đọc trong log (LSN cho Postgres, SCN cho Oracle, GTID/file+pos cho MySQL) được lưu trong offset topic — connector restart thì đọc tiếp từ đúng chỗ, không bỏ sót không lặp thừa.

Cấu trúc một change event (Debezium envelope, đơn giản hoá):

{
  "op": "u",                         // c=create, u=update, d=delete, r=read(snapshot)
  "ts_ms": 1721500000000,            // thời điểm connector xử lý
  "source": {
    "lsn": 34255112,                 // vị trí trong log — dùng để sắp thứ tự
    "table": "account",
    "txId": 90112
  },
  "before": { "id": 42, "balance": 1000000 },   // trạng thái cũ (null nếu op=c)
  "after":  { "id": 42, "balance":  850000 }    // trạng thái mới (null nếu op=d)
}

beforeafter là điều query-based không bao giờ cho được. Với DELETE, Debezium còn phát thêm một tombstone (message value = null, cùng key) để các topic log-compacted biết dọn key đó.

Khi CDC event đổ xuống lakehouse, sink thường là một job MERGE/upsert theo khóa chính (xem Iceberg streaming CDC) hoặc consumer Kafka thông thường (xem Kafka streaming). Với Oracle, hệ sinh thái đọc redo/LogMiner hoặc GoldenGate được bàn kỹ ở series Oracle CDC.

Thứ tự (ordering) và idempotency của CDC event

Đây là hai tính chất khiến CDC "đúng" hay "sai lệch âm thầm". Cả hai bắt buộc phải xử lý ở tầng consumer.

Thứ tự. Các thay đổi trên cùng một khóa (cùng một account_id) phải được xử lý đúng thứ tự commit, nếu không trạng thái cuối sẽ sai (áp update cũ đè lên update mới). Debezium bảo toàn thứ tự trong một partition Kafka, và Kafka chỉ đảm bảo thứ tự trong partition. Nguyên tắc vàng:

Partition topic CDC theo khóa chính của bảng. Mọi thay đổi của một hàng rơi vào cùng partition → cùng một consumer xử lý tuần tự → không bao giờ đảo thứ tự. Thứ tự giữa các khóa khác nhau thường không cần quan tâm.

Khi cần so sánh "cái nào mới hơn", đừng dùng thời gian tường (wall clock) — dùng cột đơn điệu trong log: lsn/scn/(binlog_file, pos). Kỹ thuật này gọi là last-writer-wins theo log position: consumer chỉ áp event nếu lsn của nó lớn hơn lsn đã ghi cho khóa đó.

Idempotency. CDC pipeline gần như luôn là at-least-once: sau một restart/rebalance, vài event có thể được phát lại. Vì thế consumer phải idempotent — áp cùng một event hai lần cho ra cùng kết quả (xem sâu ở state & exactly-once). Ba cách phổ biến:

  • Upsert theo khóa (MERGE/INSERT ... ON CONFLICT): áp lại cùng after không đổi kết quả — cách đơn giản và mạnh nhất.
  • Khử trùng theo log position: lưu lsn cuối cùng đã áp cho mỗi khóa, bỏ qua event có lsn ≤ giá trị đã lưu.
  • Idempotent sink / transactional sink: sink tự khử trùng theo event id.
-- Áp một CDC event idempotent vào bảng trạng thái hiện tại (silver)
-- Chỉ ghi đè nếu event MỚI HƠN (log position lớn hơn) — chống cả trùng lẫn đảo thứ tự
MERGE INTO silver.account AS t
USING cdc_batch AS s
  ON t.id = s.id
WHEN MATCHED AND s.op = 'd' AND s.lsn > t.src_lsn THEN
  DELETE
WHEN MATCHED AND s.op IN ('u','c','r') AND s.lsn > t.src_lsn THEN
  UPDATE SET balance = s.balance, src_lsn = s.lsn
WHEN NOT MATCHED AND s.op IN ('c','u','r') THEN
  INSERT (id, balance, src_lsn) VALUES (s.id, s.balance, s.lsn);

Dual-write problem và Outbox Pattern

CDC thuần chép trạng thái bảng. Nhưng nhiều hệ cần phát sự kiện nghiệp vụ ("khoản vay được giải ngân", "tài khoản bị đóng băng") chứ không chỉ "hàng X đổi cột Y". Cám dỗ tự nhiên là: trong một request, vừa UPDATE database vừa publish message lên Kafka. Đây chính là dual-write problem.

Không có transaction phân tán nào bao trọn "DB + Kafka" một cách rẻ và bền. Nếu commit DB xong rồi mới publish, một crash ở giữa để lại DB đã đổi nhưng sự kiện mất; nếu publish trước rồi commit DB fail, ta phát ra sự kiện cho một thay đổi chưa từng xảy ra.

Outbox Pattern giải quyết bằng cách quy mọi thứ về một transaction duy nhất trên một database. Thay vì ghi Kafka, service ghi sự kiện vào một bảng outbox trong cùng transaction nghiệp vụ. Vì bảng nghiệp vụ và bảng outbox nằm trong cùng một database, chúng commit atomic — hoặc cả hai, hoặc không gì cả. Sau đó Debezium CDC đọc chính bảng outbox và phát lên Kafka.

Điểm mấu chốt: không còn dual-write. Chỉ có một hành động ghi bền — commit database. Việc chuyển từ outbox lên Kafka do CDC lo, và CDC vốn at-least-once + bảo toàn thứ tự nên sự kiện chắc chắn được phát (có thể lặp → consumer idempotent). Debezium có sẵn Outbox Event Router (một SMT — Single Message Transform) để đọc bảng outbox và định tuyến message theo cột aggregate_type/aggregate_id, đặt key = aggregate_id để giữ đúng thứ tự sự kiện của mỗi aggregate.

-- Ghi nghiệp vụ + sự kiện trong CÙNG transaction → atomic, hết dual-write
BEGIN;
  UPDATE account SET status = 'frozen' WHERE id = 42;
  INSERT INTO outbox (id, aggregate_type, aggregate_id, event_type, payload)
  VALUES (gen_random_uuid(), 'account', '42', 'AccountFrozen',
          '{"account_id":42,"reason":"fraud_hold","by":"rule_engine"}');
COMMIT;   -- Debezium sẽ đọc bản ghi outbox này từ WAL và publish lên Kafka

Event-driven architecture: ba kiểu sự kiện

CDC và outbox đều sinh ra "sự kiện", nhưng sự kiện chứa gì quyết định kiểu kiến trúc. Theo phân loại kinh điển của Martin Fowler, có ba kiểu, và trộn lẫn chúng là nguồn gốc nhiều hệ khó hiểu.

KiểuSự kiện chứa gìƯuNhược
Event NotificationChỉ "có chuyện xảy ra" + idRất gọn, coupling thấpConsumer phải gọi ngược lại nguồn để lấy chi tiết → phụ thuộc
Event-Carried State TransferToàn bộ trạng thái liên quan (after)Consumer tự chủ, không gọi ngượcMessage lớn, có thể trùng lặp dữ liệu
Event SourcingChuỗi sự kiện nguồn sự thậtAudit hoàn hảo, replay, time-travelPhức tạp; cần snapshot; xử lý schema evolution khó
  • Event Notification: "Order 42 changed" — nhẹ nhưng tạo coupling ngầm (consumer phải query lại nguồn).
  • Event-Carried State Transfer (ECST): sự kiện mang cả after (như change event của Debezium). Consumer dựng được bản sao trạng thái cục bộ mà không cần hỏi lại nguồn — chính là mô hình CDC vào lakehouse. Đây là kiểu phổ biến nhất cho analytics.
  • Event Sourcing: không lưu trạng thái hiện tại làm nguồn sự thật; thay vào đó lưu toàn bộ chuỗi sự kiện bất biến (append-only event store), và trạng thái hiện tại được tính bằng cách replay các sự kiện. Rất khác CDC: CDC phái sinh sự kiện từ một database vốn lưu trạng thái; event sourcing coi sự kiện là bản gốc, còn bảng chỉ là projection. Đừng nhầm hai khái niệm này — đây là lỗi thuật ngữ phổ biến nhất trong lĩnh vực.

Phân biệt cốt lõi cần nhớ:

  • CDC = phái sinh changelog từ một DB lưu trạng thái. State-first, sự kiện là hệ quả.
  • Event Sourcing = sự kiện nguồn sự thật, trạng thái là hệ quả (projection). Event-first.

Cả hai đều tạo ra stream, đều cần idempotency và ordering, nhưng động cơ thiết kế ngược nhau.

Use case thực tế

Bối cảnh (minh hoạ NCB). Core banking trên Oracle. Cần: (1) dashboard số dư và giao dịch gần thời gian thực cho vận hành, (2) fraud engine phản ứng khi tài khoản bị đóng băng, (3) không được thêm tải query lên core.

Thiết kế. Debezium (đọc redo qua LogMiner) đẩy CDC của bảng account, transaction lên Kafka, partition theo account_id. Nhánh analytics: Flink MERGE upsert vào bảng silver Iceberg (ECST — mang after), độ trễ end-to-end ~1–3 phút. Nhánh nghiệp vụ: các thao tác đóng băng tài khoản ghi bảng outbox trong cùng transaction; Debezium Outbox Router phát AccountFrozen lên topic account-events; fraud consumer nhận trong ~vài giây.

Kết quả (số minh hoạ). Tải đọc thêm lên core ≈ 0 (chỉ đọc log). Không còn sự cố "DB đổi nhưng thông báo mất" vì đã bỏ dual-write. Sau một lần connector restart, ~vài nghìn event bị phát lại nhưng không gây lệch dữ liệu nhờ MERGE idempotent theo scn. Thứ tự mỗi tài khoản luôn đúng nhờ partition theo account_id. Các con số là ước lượng minh hoạ, không phải đo lường chính thức.

Ghi nhớ

  • Database là một cái log giả vờ làm bảng. Log-based CDC đọc WAL/binlog/redo — nguồn sự thật thật sự — nên bắt được cả DELETE và mọi update trung gian, tải rất thấp lên nguồn; query-based (polling updated_at) mù với delete và các thay đổi bị nuốt giữa hai lần poll.
  • Debezium = log-based CDC chạy trong Kafka Connect: snapshot ban đầu (op=r) rồi streaming; lưu offset theo log position (LSN/SCN/GTID) để restart không sót không thừa; change event có before/after/op/source.
  • Thứ tự: partition topic theo khóa chính; so sánh "mới hơn" bằng log position (LSN/SCN), không bằng wall clock.
  • Idempotency: CDC là at-least-once → consumer phải idempotent. Upsert theo khóa hoặc khử trùng theo log position (áp event chỉ khi lsn lớn hơn giá trị đã lưu).
  • Outbox pattern diệt dual-write: ghi bảng nghiệp vụ + bảng outbox trong một transaction atomic, rồi để CDC phát outbox lên Kafka. Chỉ còn một hành động ghi bền duy nhất.
  • Ba kiểu event-driven: notification (nhẹ, coupling ngầm) / event-carried state transfer (mang after, tự chủ — mô hình CDC→lakehouse) / event sourcing (sự kiện là nguồn sự thật, trạng thái là projection).
  • CDC ≠ Event Sourcing: CDC phái sinh sự kiện từ DB state-first; event sourcing coi sự kiện là bản gốc, event-first. Đừng nhầm.
  • CDC là cách chuẩn đưa OLTP vào lakehouse/analytics real-time, khớp trực tiếp với stream–table duality và là nền cho ngân hàng real-time.

Nguồn tham khảo

Bài viết liên quan

So sánh các định dạng dữ liệu (CSV, JSON, XML, Avro, Parquet, ORC) và lý do lưu theo cột nhanh hơn cho phân tích. Bài đi sâu vào row vs columnar storage, nén (Snappy/gzip/zstd), schema evolution, OLTP vs OLAP, object storage và partitioning để tối ưu chi phí lẫn tốc độ truy vấn.

13 thg 7, 2026 10

Data Engineering là ngành xây dựng và vận hành hệ thống biến dữ liệu thô thành dữ liệu sạch, tin cậy, sẵn sàng cho phân tích và AI. Bài giới thiệu vai trò Data Engineer trong vòng đời dữ liệu (nguồn → ingestion → storage → transformation → serving), phân biệt với Analyst/Scientist/ML Engineer, bức tranh hệ sinh thái công cụ và bài toán đưa dữ liệu core banking sang kho phân tích.

13 thg 7, 2026 9

Stream processing là gì, khác biệt batch vs stream (bounded/unbounded), micro-batch (Spark) vs true streaming, Apache Flink là gì và định vị so với Spark Structured Streaming và Kafka Streams. Kiến trúc runtime JobManager/TaskManager, các tầng API, triết lý streaming-first và bối cảnh phát hiện gian lận ngân hàng.

13 thg 7, 2026 8

Vì sao một máy không đủ và cần xử lý phân tán: từ MapReduce, Hadoop/HDFS đến Apache Spark in-memory. Kiến trúc driver–executor–cluster manager, các mức trừu tượng RDD/DataFrame/Dataset, cơ chế lazy evaluation với DAG, và vì sao shuffle (wide dependency) là phần tốn kém nhất. Kèm PySpark, Spark SQL và các kỹ thuật tối ưu (partition, broadcast join, cache, chống skew) cùng khi nào KHÔNG nên dùng Spark.

13 thg 7, 2026 7

Cảm nhận của bạn

Bình luận

Bạn cần để viết bình luận.

Chưa có bình luận. Hãy là người đầu tiên chia sẻ!