Iceberg 7 — Streaming & CDC: nạp dữ liệu thời gian thực
Vì sao đưa streaming vào lakehouse lại khó
Apache Iceberg sinh ra cho phân tích trên object storage với các batch job lớn. Nhưng nghiệp vụ ngày càng đòi dữ liệu near-real-time: số dư tài khoản, giao dịch thẻ, cảnh báo gian lận không thể chờ ETL chạy đêm. Khi ta cố nạp dòng dữ liệu liên tục vào một table format thiết kế cho batch, ba mâu thuẫn cốt lõi lộ ra.
Thứ nhất — nhiều commit nhỏ sinh small files. Streaming nghĩa là commit thường xuyên: mỗi vài giây tới vài phút một lần. Mỗi commit Iceberg ghi ít nhất một data file cho mỗi partition nó chạm, cộng một snapshot mới. Một job commit mỗi 30 giây, chạy 10 task song song, phân vùng theo ngày sẽ để lại hàng vạn file tí hon mỗi ngày. Small files bào mòn tốc độ query (planning chậm, mở nhiều file trên S3) và đội chi phí API — căn bệnh mà bài Bảo trì & Hiệu năng dành trọn để chữa, và streaming làm nặng hơn batch nhiều lần.
Thứ hai — cần exactly-once. Job streaming có thể chết và restart bất cứ lúc nào. Nếu sau restart nó ghi lại những bản ghi đã commit, bảng sẽ đếm trùng — với dữ liệu tài chính, một giao dịch ghi hai lần là lỗi nghiêm trọng. Ta cần mỗi bản ghi nguồn xuất hiện đúng một lần trong bảng đích, kể cả khi có failure ở giữa.
Thứ ba — upsert theo khóa. Batch analytics thường chỉ append. Nhưng dữ liệu nghiệp vụ thật thì thay đổi: tài khoản đổi số dư, khách cập nhật địa chỉ, giao dịch bị hủy. Muốn bảng lakehouse phản chiếu trạng thái hiện tại của CSDL nguồn, ta phải upsert (update nếu khóa đã tồn tại, insert nếu chưa) và delete theo khóa chính — chứ không chỉ nối thêm dòng. Đây là bài toán khó nhất, và là lý do CDC tồn tại.
Streaming ingest: Flink và Spark ghi thẳng vào Iceberg
Hướng đơn giản nhất là streaming append: đọc một luồng sự kiện (thường từ Kafka) và ghi liên tục vào bảng Iceberg. Hai engine phổ biến:
- Apache Flink — engine streaming đúng nghĩa, xử lý theo từng bản ghi (record-at-a-time), độ trễ thấp. Iceberg có connector Flink cả cho DataStream API lẫn Flink SQL.
- Spark Structured Streaming — mô hình micro-batch: gom sự kiện thành các lô nhỏ theo chu kỳ trigger rồi ghi mỗi lô như một commit. Độ trễ cao hơn Flink chút nhưng tiện nếu đội đã quen Spark.
Điểm mấu chốt vận hành là kiểm soát tần suất commit qua checkpoint interval. Cả Flink lẫn Spark chỉ commit vào Iceberg khi tới checkpoint (Flink) hoặc hết một trigger interval (Spark). Interval này là nút xoay đánh đổi trực tiếp:
| Checkpoint interval | Độ trễ (dữ liệu tới bảng) | Số file sinh ra |
|---|---|---|
| Ngắn (vd 10–30 giây) | Thấp — dữ liệu tươi | Nhiều — small files trầm trọng |
| Dài (vd 5–15 phút) | Cao hơn | Ít file hơn, mỗi file to hơn |
Không có lựa chọn miễn phí: muốn dữ liệu tươi thì phải commit dày, mà commit dày thì đẻ nhiều file. Nguyên tắc thực tế là chọn interval dài nhất mà nghiệp vụ còn chấp nhận được (ví dụ 1–5 phút thay vì 10 giây nếu dashboard chỉ cần cập nhật mỗi vài phút), rồi giao phần small files còn lại cho compaction định kỳ dọn. Đừng ép độ trễ xuống thấp hơn mức nghiệp vụ thật sự cần.
Exactly-once đến từ đâu? Iceberg commit là atomic (nhờ swap con trỏ metadata nguyên tử — xem ACID). Flink/Spark phối hợp checkpoint với commit Iceberg: chỉ commit khi checkpoint hoàn tất, và lưu offset Kafka đã đọc trong cùng checkpoint. Nếu job chết giữa chừng, nó rollback về checkpoint gần nhất và đọc lại từ offset đã lưu — dữ liệu chưa commit thì viết lại, đã commit thì không đụng. Kết quả: mỗi bản ghi vào bảng đúng một lần, dựa trên tính atomic của Iceberg commit chứ không phải phép màu.
Ví dụ Flink SQL tạo bảng Iceberg và stream append từ Kafka (MINH HOẠ — không chạy trong sandbox PostgreSQL):
-- MINH HOẠ (Flink SQL) — nguồn Kafka
CREATE TABLE src_txn (
txn_id STRING,
account_id BIGINT,
amount DECIMAL(18,2),
kind STRING,
event_ts TIMESTAMP(3),
WATERMARK FOR event_ts AS event_ts - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'core.transactions',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'debezium-json',
'scan.startup.mode' = 'group-offsets'
);
-- MINH HOẠ (Flink SQL) — sink Iceberg qua catalog
CREATE TABLE lake.db.txn (
txn_id STRING, account_id BIGINT, amount DECIMAL(18,2),
kind STRING, event_ts TIMESTAMP(3)
) PARTITIONED BY (days(event_ts))
WITH ('write.format.default'='parquet');
-- checkpoint interval đặt ở cấu hình job: execution.checkpointing.interval = 60s
INSERT INTO lake.db.txn SELECT txn_id, account_id, amount, kind, event_ts FROM src_txn;
Câu INSERT INTO … SELECT này chạy mãi mãi: nó là một streaming job, mỗi 60 giây commit một snapshot mới vào lake.db.txn. Catalog lake phải là catalog Iceberg mà cả Flink lẫn engine đọc (Trino/Spark) cùng thấy — xem Catalogs & Engines.
CDC: bắt thay đổi từ CSDL nguồn vào Iceberg
Streaming append ở trên tốt cho luồng sự kiện chỉ-thêm (như log giao dịch). Nhưng khi nguồn là một bảng OLTP thay đổi tại chỗ (bảng số dư trong core banking, bảng khách hàng), ta cần Change Data Capture (CDC): bắt từng thao tác INSERT/UPDATE/DELETE ở nguồn rồi tái hiện lên bảng Iceberg.
Chuỗi CDC điển hình:
- Debezium cắm vào transaction log của CSDL nguồn (WAL của PostgreSQL, binlog của MySQL, redo log của Oracle). Nó không query bảng mà đọc log nên không tạo tải lên OLTP và bắt được mọi thay đổi theo đúng thứ tự, kèm cả giá trị trước và sau của mỗi dòng.
- Debezium publish mỗi thay đổi thành một message lên Kafka (mỗi bảng nguồn thường một topic). Message chứa
op(c=create,u=update,d=delete), ảnhbefore/after, và khóa chính. Kafka làm buffer bền vững giữa nguồn và đích, đảm bảo thứ tự trong mỗi partition và cho phép replay — xem Kafka storage & reliability. - Flink (hoặc Spark) đọc topic CDC và MERGE vào bảng Iceberg: khóa đã tồn tại thì update, chưa có thì insert,
op=dthì delete.
Bước 3 là nơi Iceberg thể hiện năng lực row-level. Với engine batch/micro-batch, ta viết MERGE INTO tường minh (MINH HOẠ — không chạy trong sandbox):
-- MINH HOẠ (Spark SQL) — upsert một micro-batch CDC vào bảng đích
MERGE INTO lake.db.accounts t
USING cdc_batch s -- một lô bản ghi CDC vừa đọc từ Kafka
ON t.account_id = s.account_id
WHEN MATCHED AND s.op = 'd' THEN DELETE
WHEN MATCHED AND s.op IN ('u','c') THEN UPDATE SET
balance = s.balance, currency = s.currency, updated_at = s.event_ts
WHEN NOT MATCHED AND s.op <> 'd' THEN INSERT
(account_id, customer_id, balance, currency, updated_at)
VALUES (s.account_id, s.customer_id, s.balance, s.currency, s.event_ts);
Với Flink, connector Iceberg hỗ trợ upsert mode: khai báo khóa chính (primary key) trên bảng và bật 'write.upsert.enabled'='true', Flink tự dịch dòng CDC thành cặp delete + insert theo khóa, không cần viết MERGE tay.
Một cạm bẫy quan trọng: thứ tự. Nếu hai thay đổi cùng một khóa bị xử lý sai thứ tự (update cũ ghi đè update mới), số dư sẽ sai. Vì thế phải partition Kafka theo khóa chính (mọi thay đổi của một account đi cùng một partition → giữ đúng thứ tự) và dùng cột thời gian/LSN để giải quyết khi cần.
Vì sao merge-on-read hợp với streaming CDC
CDC nghĩa là rất nhiều thao tác update/delete nhỏ liên tục. Iceberg có hai cách hiện thực hóa row-level change (xem chi tiết ở ACID & Time Travel):
- Copy-on-write (CoW): mỗi update viết lại nguyên cả data file chứa dòng bị đổi. Đọc nhanh (file đã sạch) nhưng ghi cực đắt — một update chạm một dòng có thể phải viết lại file 512 MB. Không hợp streaming: tần suất ghi cao sẽ làm write amplification bùng nổ.
- Merge-on-read (MoR): update/delete chỉ ghi thêm một delete file nhỏ đánh dấu dòng nào bị vô hiệu, dữ liệu mới ghi vào data file mới. Ghi nhanh, độ trễ thấp — đúng cái streaming cần. Cái giá là lúc đọc engine phải hợp nhất data file với delete file, và delete file càng tích tụ đọc càng chậm.
Vậy MoR đổi chi phí ghi lấy chi phí đọc trả sau. Đó là đánh đổi đúng cho CDC — nhưng chỉ đúng nếu ta trả nợ đọc đều đặn bằng nén định kỳ. Compaction đọc data ⊕ delete, viết lại data file đã áp delete, rồi loại các delete file đã tiêu hóa. Không nén, bảng MoR CDC sẽ chậm dần không phanh.
So sánh ngắn: Iceberg vs Hudi vs Delta
Ba open table format lớn đều làm ACID trên object storage, nhưng thế mạnh khác nhau về streaming/CDC:
| Tiêu chí | Iceberg | Hudi | Delta Lake |
|---|---|---|---|
| Xuất phát điểm | Batch analytics, hidden partitioning | Streaming upsert / incremental (sinh ra tại Uber cho đúng bài toán này) | Batch + streaming trên Spark/Databricks |
| Upsert/CDC | Làm được qua MERGE + MoR; ngày càng tốt | Nguyên bản — index khóa chính sẵn, upsert là công dân hạng nhất | Làm được qua MERGE; mạnh trên Spark |
| Điểm mạnh riêng | Trung lập engine, evolution mạnh, chuẩn mở rộng | Record-level index, incremental query, auto-compaction ngầm | Tích hợp Databricks sâu, Delta chi tiết |
Cách đọc bảng này: Hudi được thiết kế từ đầu cho upsert streaming — có index ánh xạ khóa→file nên tìm dòng cần update nhanh, và tự nén ngầm. Nếu workload là CDC upsert khối lượng cực lớn, độ trễ thấp, ít quan tâm hệ sinh thái engine, Hudi có thể tự nhiên hơn. Delta mạnh nếu bạn ở trong hệ Databricks/Spark.
Khi nào Iceberg là đủ? Trong đa số bối cảnh ngân hàng: khi độ trễ mục tiêu là phút chứ không phải giây, khi cần nhiều engine cùng đọc (Trino cho BI, Spark cho ML, Flink cho ingest) và coi trọng tính trung lập engine + schema/partition evolution, thì Iceberg + MoR + compaction định kỳ đáp ứng tốt. Bạn cần một chuẩn mở, nhiều engine, dễ vận hành lâu dài hơn là độ tươi mili-giây — đó là lý do Iceberg thắng phần lớn use case lakehouse dù Hudi "streaming hơn".
Kiến trúc medallion trên Iceberg streaming
Streaming CDC thường tổ chức theo kiến trúc medallion ba tầng, mỗi tầng là (một hoặc nhiều) bảng Iceberg:
- Bronze (thô): hứng nguyên message CDC từ Kafka, gần như append-only, giữ cả
op/before/after— bản sao trung thực của luồng thay đổi, dùng để replay/audit. Độ trễ thấp nhất, small files nhiều nhất. - Silver (sạch): MERGE bronze thành trạng thái hiện tại theo khóa — bảng
accounts/transactionsphản chiếu core banking, đã khử trùng lặp, ép kiểu, áp quy tắc chất lượng. Đây là bảng dùng MoR + compaction nặng nhất. - Gold (nghiệp vụ): bảng tổng hợp phục vụ dashboard/report (số dư theo chi nhánh, doanh số theo ngày), thường batch/micro-batch đọc silver rồi aggregate, độ tươi nới lỏng hơn.
Bảo trì là bắt buộc, không phải tùy chọn
Với bảng batch, bỏ quên compaction vài tuần thì bảng chậm dần. Với bảng streaming CDC dùng MoR, bỏ quên bảo trì là tự sát về hiệu năng — vì cả small files lẫn delete file đều tích tụ liên tục 24/7. Hai việc tối thiểu phải chạy tự động:
- Compaction thường xuyên (
rewrite_data_files): gộp small files và gộp delete file vào data file. Bảng streaming nóng nên nén mỗi 1–6 giờ trên các partition mới, dùngwheređể chỉ đụng phần vừa ghi. Với CDC, mục tiêu chính không phải trị small files thông thường mà là không để delete file tồn đọng làm đọc chậm. - expire_snapshots thường xuyên: streaming commit rất nhiều snapshot (mỗi phút một cái → hàng nghìn snapshot/ngày). Nếu không expire, metadata phình khủng khiếp và storage không bao giờ được giải phóng (file cũ vẫn bị snapshot cũ giữ). Chạy hằng ngày, giữ 3–7 ngày với bảng thô.
Chi tiết procedure, nhịp chạy và cách tự động hóa bằng Airflow nằm trọn ở Bảo trì & Hiệu năng. Điều cần khắc cốt: streaming ingest và maintenance job là một cặp không thể tách rời — thiết kế pipeline streaming mà chưa thiết kế lịch bảo trì đi kèm là thiết kế dở dang.
Use case thực tế
Bối cảnh (số liệu ước lượng minh họa). NCB muốn có dữ liệu số dư tài khoản và giao dịch gần thời gian thực cho analytics: dashboard giám sát dòng tiền, đối soát trong ngày, và feature cho mô hình cảnh báo gian lận — thay vì chờ batch ETL chạy đêm (dữ liệu trễ tới 24 giờ). Nguồn là core banking trên Oracle, không được phép query trực tiếp vì sợ tải lên hệ thống giao dịch.
Kiến trúc. Đội dựng pipeline CDC theo đúng chuỗi trên:
- Debezium đọc redo log của Oracle cho hai bảng
ACCOUNTSvàTRANSACTIONS, không tạo tải query lên core. Mỗi thay đổi thành message lên Kafka, topic partition theoaccount_idđể giữ đúng thứ tự thay đổi của mỗi tài khoản. - Flink consume hai topic, checkpoint mỗi 60 giây (exactly-once), ghi tầng bronze (thô, append) rồi MERGE vào tầng silver ở chế độ merge-on-read:
op=u/cupsert theoaccount_id/txn_id,op=ddelete. Bảng silver phân vùng theodays(event_ts). - Airflow chạy lịch bảo trì:
rewrite_data_files(sort theoaccount_id, gộp delete file) mỗi giờ trên partition ngày hiện tại;expire_snapshotsgiữ 7 ngày mỗi đêm;remove_orphan_fileshằng tuần.
Vấn đề gặp và cách chỉnh. Bản đầu đội đặt checkpoint 10 giây cho "tươi nhất có thể". Kết quả: mỗi partition ngày sinh ~50.000 file cộng hàng nghìn delete file, query đối soát chậm dần. Đội nới checkpoint lên 60 giây (nghiệp vụ chỉ cần độ tươi ~1 phút, không cần 10 giây) và tăng tần suất compaction lên mỗi giờ. Sau điều chỉnh:
- Độ trễ end-to-end (thay đổi ở core → thấy trên bảng silver): khoảng 1–3 phút, đủ cho dashboard và đối soát trong ngày.
- Số data file mỗi partition ngày sau nén ổn định ở ~200–400 file thay vì hàng vạn; delete file không còn tồn đọng nhờ compaction hằng giờ.
- Query đối soát lọc theo ngày + nhóm
account_idchạy trong ~10 giây thay vì 1–2 phút, chủ yếu nhờ planning ngắn và pruning theoaccount_id(đã sort). - Tải lên core banking gần như không đổi vì Debezium chỉ đọc log, không query bảng — điều kiện then chốt để bộ phận vận hành core chấp thuận.
Đánh đổi đã chấp nhận: không đuổi theo độ trễ mili-giây; chọn "phút" đủ cho analytics và giữ small files trong tầm kiểm soát. Các con số trên là ước lượng minh họa cho một pipeline CDC điển hình, không phải đo lường chính thức.
Ghi nhớ
- Streaming vào lakehouse khó vì ba lý do: nhiều commit nhỏ → small files (chữa bằng compaction); cần exactly-once (tránh đếm trùng dữ liệu tài chính); cần upsert/delete theo khóa chứ không chỉ append.
- Streaming ingest: Flink (record-at-a-time, độ trễ thấp) hoặc Spark Structured Streaming (micro-batch) ghi vào Iceberg. Checkpoint interval là nút xoay đánh đổi độ trễ vs số file — chọn interval dài nhất nghiệp vụ chấp nhận. Exactly-once đến từ việc phối hợp checkpoint với Iceberg commit atomic và offset Kafka đã lưu.
- CDC: Debezium đọc transaction log nguồn (không tải lên OLTP) → publish thay đổi lên Kafka (buffer bền, giữ thứ tự theo khóa) → Flink/Spark MERGE upsert/delete vào Iceberg. Partition Kafka theo khóa chính để không sai thứ tự.
- Merge-on-read hợp streaming: update/delete chỉ ghi delete file nhỏ (ghi nhanh, độ trễ thấp), engine gộp lúc đọc — đổi chi phí ghi lấy chi phí đọc trả sau. Bắt buộc nén định kỳ để gộp delete file, nếu không đọc chậm dần vô phanh (xem ACID & MoR).
- So với Hudi/Delta: Hudi mạnh upsert/CDC nguyên bản (index khóa, auto-compaction); Delta mạnh trong hệ Databricks. Iceberg đủ khi độ trễ mục tiêu là phút, cần nhiều engine cùng đọc và coi trọng tính trung lập engine + evolution.
- Medallion trên Iceberg: bronze (CDC thô, append) → silver (trạng thái hiện tại, MoR + nén nặng) → gold (bảng nghiệp vụ tổng hợp).
- Bảo trì là bắt buộc: bảng streaming CDC phải chạy compaction (gộp small files + delete file) mỗi 1–6 giờ và expire_snapshots hằng ngày — streaming ingest và maintenance job là một cặp không tách rời.
- Xem thêm: tổng quan Iceberg, ACID & time travel, catalogs & engines, bảo trì & hiệu năng.
Nguồn tham khảo
- Apache Iceberg Documentation — Flink Writes / Flink DDL & Upsert
- Apache Iceberg Documentation — Spark Structured Streaming
- Apache Iceberg Documentation — Spark Writes: mục
MERGE INTOvà merge-on-read (copy-on-write vs merge-on-read) - Debezium Documentation — kiến trúc CDC, connector cho Oracle/PostgreSQL/MySQL, cấu trúc message change event (
op,before/after) - Apache Kafka Documentation — Kafka Connect
- Apache Flink Documentation — Fault Tolerance / Checkpointing
- Apache Hudi Documentation và Delta Lake Documentation — tham chiếu so sánh table format
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ẻ!