Streaming 6 — Stream-Table Duality & Materialized View
Mô hình tinh thần: stream và table là hai mặt của một tờ giấy
Nếu chỉ nhớ một câu từ bài này, hãy nhớ câu của Jay Kreps (đồng sáng lập Kafka) trong bài kinh điển The Log:
Một stream là changelog của một bảng; một bảng là snapshot (ảnh chụp) của một stream tại một thời điểm.
Đây không phải ẩn dụ văn vẻ mà là một đẳng thức vận hành. Hai cấu trúc dữ liệu tưởng chừng đối lập — một cái chảy (vô hạn, chỉ nối thêm), một cái đứng yên (hữu hạn, tra cứu theo khóa) — thực ra là hai cách nhìn cùng một thông tin, và bạn chuyển đổi qua lại giữa chúng bất cứ lúc nào:
- Table → Stream: quan sát mọi thay đổi của một bảng theo thời gian, bạn thu được một dòng các bản cập nhật. Đây chính là changelog — và cũng chính là ý tưởng của CDC (Change Data Capture): đọc redo/WAL log của CSDL để biến bảng thành stream.
- Stream → Table: phát lại (replay) một dòng cập nhật từ đầu và áp từng thay đổi lên trạng thái, bạn dựng lại được bảng ở trạng thái hiện tại. Đây là cái mà stream processor gọi là materialize (vật chất hoá).
Chính tính hai chiều này là nền tảng lý thuyết của kiến trúc Kappa: nếu log là nguồn sự thật và mọi bảng chỉ là view được vật chất hoá lại từ log, thì ta không cần lớp batch riêng — chỉ cần replay log là tái tạo được mọi trạng thái.
Bảng bên trái chỉ giữ giá trị mới nhất mỗi key (acc:001 = 5.000.000). Stream bên phải giữ toàn bộ lịch sử cập nhật — và nếu bạn replay nó theo thứ tự rồi giữ lại bản ghi cuối cho mỗi key, bạn dựng lại đúng cái bảng đó. Không mất mát thông tin theo chiều table→stream; theo chiều stream→table thì "gập" lịch sử lại thành trạng thái hiện tại.
Event stream vs changelog stream — khác biệt quyết định ngữ nghĩa
Không phải stream nào cũng giống nhau. Có hai loại stream với ngữ nghĩa khác hẳn, và nhầm lẫn chúng là nguồn gốc của phần lớn bug trong stream processing:
| Event stream (dòng sự kiện) | Changelog stream (dòng cập nhật) | |
|---|---|---|
| Mỗi bản ghi là | một sự thật độc lập đã xảy ra | một cập nhật trạng thái theo key |
| Ngữ nghĩa | INSERT (append thuần) | UPSERT (insert-or-update theo key) |
| Key trùng nghĩa là | hai sự kiện khác nhau | phiên bản mới đè phiên bản cũ |
value = null nghĩa là | một sự kiện bình thường | tombstone — xoá key khỏi bảng |
| Vật thể hoá thành | log các fact (không "gập") | một bảng (giá trị mới nhất/key) |
| Ví dụ ngân hàng | "thẻ 123 quẹt 200k lúc 10:05" | "số dư acc:001 giờ là 5.000.000" |
Event stream là dòng các giao dịch, cú click, lần quẹt thẻ — mỗi cái xảy ra một lần và không bao giờ bị "sửa". Cộng dồn chúng lại không có nghĩa là đè lên nhau; hai giao dịch 200k là hai giao dịch.
Changelog stream là dòng các ảnh chụp trạng thái mới theo key — "số dư của tài khoản X bây giờ là Y". Bản ghi sau thay thế bản ghi trước cho cùng key. Đây chính là output tự nhiên của CDC và của các phép tổng hợp có trạng thái (aggregation).
Trong Kafka/Flink, sự khác biệt này được mã hoá trực tiếp: một topic/stream đọc theo ngữ nghĩa append cho ra event stream; đọc theo ngữ nghĩa upsert theo key (kèm tombstone) cho ra changelog → bảng.
Compacted topic — cái bảng nằm ngay trong Kafka
Làm sao lưu một bảng trong một hệ thống vốn chỉ biết log? Kafka trả lời bằng log compaction.
Topic Kafka bình thường xoá dữ liệu theo thời gian/dung lượng (retention: giữ 7 ngày rồi bỏ). Topic đặt cleanup.policy=compact thì khác: nó giữ lại ít nhất bản ghi mới nhất cho mỗi key, và dọn dần các phiên bản cũ hơn của cùng key ở nền. Kết quả là log đó, khi đọc từ đầu tới cuối, cho bạn đúng một giá trị hiện hành mỗi key — tức là một bảng vật chất hoá, bền vững, có thể replay.
- Key = khóa chính của bảng (mã tài khoản, mã khách...).
- Value mới nhất = trạng thái hiện tại của key đó.
- Tombstone (
value = null) báo compaction xoá hẳn key — tương đươngDELETEtrong bảng.
Đây là lý do Kafka Streams lưu trạng thái của KTable và store nội bộ (RocksDB) bằng cách backup lên một compacted topic: nếu instance chết, nó chỉ cần replay compacted topic là dựng lại y nguyên bảng — không mất trạng thái. Chi tiết cơ chế lưu trữ và compaction ở Kafka — Storage & Reliability. Cũng chính nhờ tính bền + replay được này mà state & exactly-once khả thi: trạng thái không nằm trong RAM dễ bay, mà tựa trên một log bền có thể tái tạo.
Một compacted topic vì thế vừa là stream (bạn subscribe để nhận mọi cập nhật mới), vừa là table (bạn đọc snapshot để biết trạng thái hiện tại) — hiện thân vật lý rõ nhất của duality.
KStream vs KTable — duality thành API
Kafka Streams (và ksqlDB xây trên nó) biến lý thuyết trên thành hai kiểu trừu tượng đối ngẫu:
- KStream — diễn giải topic như event stream. Mỗi bản ghi là một fact độc lập, append-only. Hai bản ghi cùng key = hai sự kiện. Dùng cho giao dịch, log, click.
- KTable — diễn giải topic như changelog stream. Mỗi bản ghi là một upsert theo key; bản ghi mới đè bản cũ;
null= xoá. KTable là một bảng được nuôi liên tục bởi một stream. Dùng cho trạng thái: số dư, hồ sơ khách, tỉ giá.
Và bạn chuyển đổi qua lại đúng như duality mô tả:
| Phép | Ý nghĩa |
|---|---|
stream.groupByKey().aggregate() → KTable | stream → table: gập event stream thành trạng thái |
table.toStream() → KStream | table → stream: phát ra changelog của mọi thay đổi |
KStream join KTable | enrichment: mỗi event tra cứu trạng thái hiện tại theo key |
Flink có khái niệm tương đương gọi là dynamic table: một bảng "động" được cập nhật liên tục bởi một stream, và mỗi continuous query trên nó lại phát ra một changelog stream (các dòng INSERT / UPDATE_BEFORE / UPDATE_AFTER / DELETE). Append-only query cho stream chỉ-thêm; query có aggregation/join cho updating (retract) stream. Về bản chất Flink SQL và ksqlDB đang nói cùng một ngôn ngữ duality, chỉ khác từ vựng — xem Flink SQL.
Materialized view real-time từ stream
Trong CSDL truyền thống, materialized view là kết quả một query được lưu sẵn để đọc nhanh, và phải refresh định kỳ (chạy lại query) nên luôn cũ. Stream processing lật ngược mô hình: view được cập nhật tăng dần (incremental) mỗi khi có sự kiện mới tới — không refresh cả bảng, chỉ áp đúng delta. Kết quả là một view luôn tươi, phản ánh trạng thái mới nhất trong mili-giây tới giây.
Điểm mấu chốt: output của một aggregation là một changelog stream, và khi vật chất hoá theo key nó thành một bảng (KTable / dynamic table). Đúng lại là duality — query chạy trên stream, kết quả là một bảng cập nhật liên tục.
Ba công cụ tiêu biểu, cùng nguyên lý:
- ksqlDB —
CREATE TABLE ... AS SELECTdựng một materialized view được nuôi bởi stream; hỗ trợ pull query (tra cứu điểm một key, như đọc bảng) và push query (EMIT CHANGES, đẩy mỗi khi giá trị đổi). Xem Kafka Streams & ksqlDB. - Flink SQL — dynamic table + continuous query; upsert-kafka connector để đọc/ghi changelog theo key.
- Materialize — một CSDL chuyên incremental view maintenance: bạn viết
CREATE MATERIALIZED VIEWbằng SQL chuẩn (kể cả join nhiều bảng), engine duy trì kết quả tăng dần với độ trễ thấp mà không refresh toàn phần.
So với materialized view của kho dữ liệu (data warehouse / dimensional modeling) vốn refresh theo lịch (phút/giờ), materialized view streaming cập nhật theo sự kiện — hợp cho dashboard vận hành, cảnh báo, đối soát trong ngày. Ngược lại, lịch sử đầy đủ của một chiều thay đổi chậm (SCD) chính là changelog stream được lưu lại — thêm một góc nhìn duality: bảng dim với lịch sử = table + toàn bộ changelog của nó.
SQL minh hoạ: materialized view + join stream-table
Ví dụ ksqlDB. Đầu tiên dựng một materialized view cộng dồn chi tiêu theo tài khoản từ event stream giao dịch (MINH HOẠ — cú pháp ksqlDB, không chạy trong sandbox PostgreSQL):
-- MINH HOẠ (ksqlDB) — event stream: mỗi bản ghi là 1 giao dịch (INSERT)
CREATE STREAM payments (
account_id BIGINT,
amount DECIMAL(18,2),
ts BIGINT
) WITH (kafka_topic='payments', value_format='json', timestamp='ts');
-- Materialized view: aggregation => KTable (changelog theo key, luôn tươi)
CREATE TABLE spend_by_account
WITH (kafka_topic='spend_by_account') -- topic này là compacted: bảng bền
AS
SELECT account_id,
SUM(amount) AS total_spend,
COUNT(*) AS txn_count
FROM payments
GROUP BY account_id
EMIT CHANGES;
-- Pull query: tra cứu ĐIỂM một key như đọc bảng (không stream)
SELECT total_spend, txn_count
FROM spend_by_account
WHERE account_id = 1001;
CREATE TABLE ... AS SELECT ... EMIT CHANGES biến một aggregation trên stream thành một bảng được cập nhật liên tục — output là changelog (mỗi lần tổng đổi, phát một upsert account_id → total mới). Backing store là compacted topic nên bảng bền và replay được.
Tiếp theo là join stream-table (enrichment) — mẫu dùng nhiều nhất của duality. Mỗi event giao dịch được làm giàu bằng trạng thái hiện tại của bảng hồ sơ tài khoản (một KTable nuôi bởi CDC):
-- MINH HOẠ (ksqlDB) — KTable trạng thái: hồ sơ tài khoản (upsert theo key)
CREATE TABLE accounts (
account_id BIGINT PRIMARY KEY,
segment STRING, -- phân khúc KH
branch STRING, -- chi nhánh
risk_flag STRING -- cờ rủi ro hiện hành
) WITH (kafka_topic='accounts_cdc', value_format='json');
-- STREAM (event) JOIN TABLE (state): tra cứu theo key tại thời điểm event tới
CREATE STREAM payments_enriched AS
SELECT p.account_id,
p.amount,
a.segment,
a.branch,
a.risk_flag
FROM payments p
JOIN accounts a ON p.account_id = a.account_id
EMIT CHANGES;
Ngữ nghĩa cần nắm chắc: trong stream-table join, chỉ phía stream (event) kích hoạt một dòng kết quả; phía table đóng vai bảng tra cứu, engine đọc giá trị hiện hành tại đúng thời điểm event tới. Đây là temporal join — khác join hai bảng tĩnh: nếu risk_flag của tài khoản đổi lúc 10:00, thì giao dịch lúc 09:59 vẫn thấy cờ cũ, giao dịch lúc 10:01 thấy cờ mới. Chính vì table được nuôi bởi một changelog (CDC từ core), nó luôn là snapshot mới nhất mà không cần query lại CSDL nguồn cho từng event.
Bản Flink SQL tương đương dùng upsert-kafka để khai báo phía bảng đọc theo changelog:
-- MINH HOẠ (Flink SQL) — bảng đọc changelog theo key (upsert + tombstone)
CREATE TABLE accounts (
account_id BIGINT,
segment STRING,
risk_flag STRING,
PRIMARY KEY (account_id) NOT ENFORCED
) WITH (
'connector' = 'upsert-kafka',
'topic' = 'accounts_cdc',
'key.format' = 'json',
'value.format' = 'json'
);
upsert-kafka chính là cách Flink nói "topic này là changelog của một bảng": mỗi bản ghi là upsert theo primary key, null value là tombstone/xoá. Đối lập với kafka connector thường (append-only, event stream).
Ứng dụng: khi nào nghĩ theo duality thì bài toán tự gỡ
- Enrichment real-time: làm giàu dòng giao dịch bằng hồ sơ/tỉ giá/hạn mức hiện hành — stream-table join, phía bảng nuôi bằng CDC. Không phải gọi DB cho từng event.
- Cache/lookup luôn tươi: thay vì cache tự viết dễ lệch, subscribe một compacted topic → có sẵn một bảng khóa-giá trị đồng bộ liên tục, chết-dựng-lại được.
- Materialized view vận hành: tổng/đếm/số dư theo key phục vụ dashboard, cảnh báo, đối soát trong ngày — cập nhật tăng dần, không refresh cả bảng.
- Tái tạo trạng thái sau sự cố: vì bảng = replay của log, mất instance chỉ cần đọc lại changelog — nền tảng của exactly-once & state.
- Bỏ lớp batch trùng lặp: coi log là nguồn sự thật, mọi bảng là view vật chất hoá lại từ log — tinh thần Kappa.
Use case thực tế
Bối cảnh (NCB — minh hoạ). Đội risk cần một màn hình đối soát số dư & cờ rủi ro theo tài khoản cập nhật gần tức thời để chặn giao dịch bất thường, thay cho báo cáo EOD chạy đêm. Yêu cầu: mỗi giao dịch phải thấy hồ sơ + cờ rủi ro hiện hành của tài khoản, và có một materialized view tổng chi tiêu trong ngày tra cứu theo account_id dưới 100ms.
Thiết kế theo duality.
- Table → Stream (CDC): Debezium đọc redo log của core banking, đẩy thay đổi bảng
ACCOUNTSlên topicaccounts_cdcdạng compacted (giữ trạng thái mới nhất mỗi tài khoản, tombstone khi đóng tài khoản). Topic này chính là bảng hồ sơ, ở dạng stream bền. - Stream event: giao dịch từ switch thẻ vào topic
payments(event stream, append-only). - Stream-table join: Flink/ksqlDB join
payments(event) vớiaccounts(KTable từ CDC) →payments_enriched, mỗi giao dịch mang theosegment,branch,risk_flagtại đúng thời điểm phát sinh. - Materialized view:
CREATE TABLE spend_by_account AS SELECT ... GROUP BY account_id EMIT CHANGES— pull query cho màn hình tra cứu điểm.
Kết quả minh hoạ (ước lượng cho một pipeline điển hình, không phải đo chính thức).
- Độ trễ giao dịch → thấy trên màn hình enrich: khoảng 1–3 giây, so với tới 24 giờ của báo cáo EOD cũ.
- Pull query tra cứu tổng chi tiêu theo
account_id: ~10–50ms (đọc điểm trên materialized view, không quét bảng lớn). - Tải lên core banking gần như không đổi: enrichment đọc từ KTable trong bộ nhớ/RocksDB, không query core cho từng giao dịch — điều kiện then chốt để vận hành core chấp thuận.
- Sau sự cố instance: trạng thái KTable dựng lại bằng cách replay compacted topic, không mất số liệu tổng hợp.
Đánh đổi đã chấp nhận: join là temporal nên nếu CDC trễ, một số giao dịch sát nút đổi cờ có thể thấy cờ cũ vài giây; đội chấp nhận vì thà trễ đúng còn hơn nhanh sai, và đặt ngưỡng cảnh báo độ trễ CDC.
Ghi nhớ
- Duality (Jay Kreps): stream = changelog của bảng; bảng = snapshot của stream tại một thời điểm. Table→stream = quan sát thay đổi (CDC); stream→table = replay + gập theo key (materialize).
- Event stream vs changelog stream là khác biệt ngữ nghĩa gốc: event = INSERT (fact độc lập, không gập); changelog = UPSERT theo key (bản mới đè bản cũ,
null= tombstone/xoá). Nhầm hai loại này là nguồn của phần lớn bug. - Compacted topic (
cleanup.policy=compact) giữ giá trị mới nhất mỗi key → là một bảng vật chất hoá, bền, replay được nằm ngay trong Kafka; vừa là stream (subscribe) vừa là table (snapshot). - KStream = event stream, KTable = changelog/bảng. Chuyển đổi:
aggregate()(stream→KTable),toStream()(KTable→changelog stream). Flink gọi là dynamic table, phát changelog+I/-U/+U/-D. - Materialized view streaming cập nhật tăng dần theo sự kiện (không refresh cả bảng như warehouse) → luôn tươi. ksqlDB
CREATE TABLE AS SELECT ... EMIT CHANGES(pull + push query); Flink SQL dynamic table; Materialize (incremental view maintenance). - Stream-table join = enrichment: chỉ phía event kích hoạt kết quả, phía table là lookup đọc trạng thái hiện hành tại thời điểm event (temporal join). Bảng nuôi bằng CDC nên không phải query DB nguồn cho từng event.
- Vì bảng = replay của log, khôi phục sau sự cố chỉ là đọc lại changelog — nền tảng của exactly-once & state và của Kappa.
Nguồn tham khảo
- Jay Kreps — The Log: What every software engineer should know about real-time data's unifying abstraction (Confluent / LinkedIn Engineering blog)
- Jay Kreps — Introducing Kafka Streams: Stream Processing Made Simple (Confluent blog) — trình bày stream-table duality
- Apache Kafka Documentation — Kafka Streams: Duality of Streams and Tables và Log Compaction
- ksqlDB Documentation — Materialized views, Streams & Tables, Pull vs Push queries
- Apache Flink Documentation — Dynamic Tables & Continuous Queries và upsert-kafka connector
- Materialize Documentation — Materialized views / incremental view maintenance
- "Designing Data-Intensive Applications" — Martin Kleppmann (O'Reilly), chương về logs, replication và derived data
- "Streaming Systems" — Akidau, Chernyak, Lax (O'Reilly), phần stream-and-table theory
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.
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.
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.
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.
Cảm nhận của bạn
Bình luận
Chưa có bình luận. Hãy là người đầu tiên chia sẻ!