Event Streams, CDC & Event Sourcing — Cốt lõi stream processing (DDIA)

Phong

Mở đầu

Khi xây dựng hệ thống phân tán, một trong những câu hỏi khó nhất là: làm sao để đồng bộ dữ liệu giữa các service, database, và cache một cách nhất quán? Batch processing (xử lý theo lô) giải quyết được một phần, nhưng với những hệ thống cần độ trễ thấp, batch không đáp ứng kịp.

Ảnh: Tibe De Kort — Pexels

Chương 11 của DDIA giới thiệu một cách tiếp cận khác: stream processing — xử lý dữ liệu khi nó vừa được sinh ra, thay vì chờ đến cuối ngày mới chạy batch. Trong phần 11.1, tác giả tập trung vào ba khái niệm nền tảng: Event Streams, Change Data Capture (CDC), và Event Sourcing.

Event Streams — Dữ liệu chảy như dòng sông

Trong thế giới batch, bạn có một lượng dữ liệu cố định (file, table snapshot) và bạn chạy job để xử lý nó. Kết quả là một file/output khác. Nhưng trong thế giới stream, dữ liệu không bao giờ ngừng chảy — mỗi sự kiện (event) là một bản ghi nhỏ, bất biến, được append vào một log.

Một event stream có thể được hiểu như một message queue nhưng có khả năng replay — tức là bạn có thể tua lại và đọc lại các sự kiện cũ. Đây là điểm khác biệt lớn so với message queue truyền thống (như RabbitMQ) nơi message bị xoá sau khi được consume.

Các hệ thống như Apache Kafka, AWS Kinesis, hay Pulsar đều dùng mô hình này: event được append vào log phân tán, và consumer có thể đọc ở bất kỳ vị trí nào.

Ảnh: Marek Prášil — Pexels

Change Data Capture (CDC) — Khi database tự kể chuyện

Change Data Capture là kỹ thuật "bắt" những thay đổi trong database và biến chúng thành event stream. Thay vì phải poll database định kỳ (SELECT * FROM table WHERE updated_at > last_check), CDC đọc trực tiếp từ write-ahead log (WAL) của database — cùng cái log mà database dùng để recovery.

Mỗi lần có một row được INSERT, UPDATE, hoặc DELETE, database ghi nó vào WAL trước. CDC consumer đọc WAL và biến dòng log đó thành một event có cấu trúc:

  • Insert → event chứa toàn bộ row mới
  • Update → event chứa row trước và row sau (tuỳ cấu hình)
  • Delete → event chứa row bị xoá

Điều này cho phép bạn đồng bộ dữ liệu giữa các hệ thống gần như real-time. Công cụ phổ biến: Debezium (connector Kafka), AWS DMS, PostgreSQL logical replication.

Ví dụ thực tế: Bạn có một database chính (PostgreSQL) và một search index (Elasticsearch). Thay vì dual-write (ghi cả 2 nơi — dễ bị inconsistent), bạn bật CDC trên PostgreSQL, để Debezium stream các change event vào Kafka, rồi Kafka consumer cập nhật Elasticsearch. Kết quả: Elasticsearch luôn sync với database chính, và bạn không cần sửa code application.

Event Sourcing — Lưu trữ hành trình, không chỉ đích đến

Event Sourcing là một ý tưởng cực kỳ thú vị: thay vì lưu trạng thái hiện tại của một entity, bạn lưu tất cả các sự kiện đã xảy ra với entity đó. Trạng thái hiện tại là kết quả của việc replay tất cả các sự kiện từ đầu.

Ảnh: Google DeepMind — Pexels

Ví dụ: Một tài khoản ngân hàng. Cách truyền thống: table accounts có cột balance = 1000. Cách Event Sourcing: bạn lưu các event như AccountOpened(amount=1000), Withdrawn(amount=200), Deposited(amount=500). Để biết balance hiện tại, bạn replay tất cả event: 1000 - 200 + 500 = 1300.

Lợi ích của Event Sourcing:

  • Audit trail đầy đủ: Bạn biết chính xác ai đã làm gì và khi nào
  • Debug dễ dàng: Có thể replay events trên môi trường dev để tái hiện bug
  • Temporal query: Trả lời câu hỏi "balance hồi tháng trước là bao nhiêu?" bằng cách replay events đến thời điểm đó
  • CQRS kết hợp tự nhiên: Command side ghi event, Query side snapshot state

Tuy nhiên, Event Sourcing cũng có chi phí:

  • Hệ thống event store cần xử lý lượng event khổng lồ
  • Schema evolution phức tạp — event đã ghi là bất biến, không thể sửa
  • Snapshot strategy cần thiết để tránh replay từ đầu mỗi lần

Mối quan hệ giữa CDC và Event Sourcing: CDC biến đổi database thành event stream (dựa trên log), còn Event Sourcing thiết kế application để lưu event ngay từ đầu. CDC là cách "retrofit" event stream cho database cũ; Event Sourcing là cách xây dựng hệ thống với event làm trung tâm.

Key Takeaways

  • Event streams là mô hình dữ liệu bất biến, có thể replay, phù hợp cho real-time processing
  • Change Data Capture (CDC) cho phép biến database transaction log thành event stream — đồng bộ dữ liệu giữa các hệ thống mà không cần dual-write
  • Event Sourcing lưu trữ sự kiện thay vì trạng thái, mở ra khả năng audit, debug, và temporal query
  • CDC và Event Sourcing có thể kết hợp: CDC để capture changes từ legacy DB, Event Sourcing cho hệ thống mới

Glossary

Thuật ngữ Ý nghĩa
Event stream Luồng sự kiện bất biến, có thứ tự, có thể replay từ đầu
Change Data Capture (CDC) Kỹ thuật bắt thay đổi từ database transaction log và biến thành event
Write-ahead log (WAL) Log ghi trước khi dữ liệu được apply vào database — dùng cho recovery và replication
Event Sourcing Pattern lưu tất cả sự kiện thay vì trạng thái hiện tại
Replay Đọc lại toàn bộ event stream từ đầu
Snapshot Ảnh chụp trạng thái tại một thời điểm, dùng để tránh replay quá dài
CQRS Command Query Responsibility Segregation — tách command (ghi) và query (đọc)

Kết

Chương 11.1 của DDIA đặt nền móng cho việc hiểu stream processing bằng cách giới thiệu ba khái niệm: event stream như một log bất biến, CDC như cầu nối giữa database và stream, và Event Sourcing như một cách tư duy lại về lưu trữ dữ liệu.

Trong các phần tiếp theo, Kleppmann sẽ đi sâu hơn về cách xử lý stream (join, windowing, exactly-once semantics) và các ứng dụng thực tế. Nếu bạn đang xây dựng hệ thống real-time, những khái niệm này là bắt buộc phải nắm.