Streaming 5 — State, Checkpoint & Exactly-once

22 thg 7, 2026 3 lượt xem
#data-engineering
#streaming
#checkpoint
#exactly-once
#state

Mô hình tinh thần: stream có trí nhớ

Một phép biến đổi stateless như filter hay map chỉ nhìn từng bản ghi rồi quên ngay — bản ghi thứ N không liên quan gì bản ghi thứ N-1. Nhưng phần lớn giá trị của stream processing lại nằm ở những phép cần nhớ quá khứ:

  • Tổng hợp (aggregation): "tổng số tiền chuyển của tài khoản A trong 5 phút qua" — phải cộng dồn, tức phải nhớ tổng đang chạy.
  • Join theo thời gian: ghép sự kiện authorize với capture của cùng một giao dịch thẻ đến sau vài giây — phải giữ bên đến trước để chờ bên kia.
  • Dedup: loại bản ghi trùng txn_id — phải nhớ những id đã thấy.
  • Phát hiện chuỗi (pattern/CEP): "3 lần rút tiền thất bại rồi 1 lần thành công trong 60 giây" — phải theo dõi trạng thái chuyển tiếp.

Tất cả đều là stateful stream processing: bộ toán tử duy trì một khối state sống giữa các sự kiện. Câu hỏi lớn của cả bài này là: khối trí nhớ đó sống ở đâu, và làm sao không mất nó khi máy chết giữa dòng dữ liệu vô tận? Một job batch chết thì chạy lại từ đầu; một stream job chạy 24/7 xử lý dòng vô hạn không thể "chạy lại từ đầu" — nó phải khôi phục đúng trạng thái đang có rồi đi tiếp. Đó là lý do state và checkpoint là trái tim của mọi engine streaming nghiêm túc (Flink, Spark Structured Streaming, Kafka Streams).

Bài này nối tiếp Streaming 4 — Windowing: mỗi window đang mở chính là một mẩu state đang tích luỹ. Ở đây ta hỏi: state đó lưu ở đâu và làm sao bền vững.

State sống ở đâu: local state + state backend

Trực giác sai lầm phổ biến: "cứ để state trong một database ngoài (Redis/Postgres) cho chắc". Với thông lượng hàng trăm nghìn sự kiện/giây, mỗi sự kiện phải gọi mạng ra DB ngoài để đọc-ghi state sẽ giết chết throughput và thêm điểm chết. Engine streaming hiện đại làm ngược lại: state nằm local, ngay cạnh phép tính, trên chính node đang xử lý.

Hai lớp cần phân biệt rạch ròi:

  • Local state (state đang làm việc): cấu trúc dữ liệu engine đọc/ghi với độ trễ cực thấp trên node. Có hai lựa chọn phổ biến trong Flink:
    • Heap/in-memory backend: state là object Java trên JVM heap — nhanh nhất, nhưng bị giới hạn bởi RAM và áp lực GC; hợp state nhỏ.
    • RocksDB backend: state lưu trong một embedded key-value store (LSM-tree) ghi ra đĩa local. Cho phép state lớn hơn RAM (hàng trăm GB đến TB mỗi node), đổi lại mỗi truy cập có chi phí serialize/deserialize và đọc đĩa. Đây là lựa chọn mặc định cho state lớn: cửa sổ dài, bảng dedup rộng, join giữ nhiều bản ghi.
  • State backend / durable store: nơi engine sao lưu local state một cách bền vững — thường là object store (S3, GCS) hoặc HDFS. Đây là bản sao để phục hồi, không phải nơi đọc-ghi nóng.

Mấu chốt: đọc-ghi nóng ở local (nhanh), sao lưu bền vững ra xa (an toàn). RocksDB không phải "DB ngoài" — nó là thư viện nhúng chạy trong process của task, dữ liệu trên đĩa local; còn checkpoint mới là thứ đẩy ra store bền vững.

Keyed state — vì sao state luôn gắn với key

Để song song hoá, stream được phân vùng theo key (keyBy trong Flink, key của Kafka Streams). Sau keyBy(account_id), mọi sự kiện của cùng một tài khoản luôn tới cùng một task. Nhờ vậy state được chia thành keyed state: mỗi key có "ngăn" state riêng, và một task chỉ giữ state cho tập key nó phụ trách. Điều này cho hai thứ:

  1. Song song không tranh chấp: task chỉ chạm state của key mình — không khoá, không đồng bộ chéo node.
  2. Rescale được: khi tăng/giảm parallelism, engine phân phối lại các "ngăn" keyed state cho số task mới (Flink dùng key groups làm đơn vị phân phối lại).

Các loại keyed state hay gặp: ValueState (một giá trị/khoá — ví dụ tổng đang chạy), ListState (danh sách — các sự kiện chờ join), MapState (map — bảng dedup id đã thấy), ReducingState/AggregatingState (gộp dần).

// Flink: đếm số giao dịch chống-gian-lận theo từng tài khoản (keyed ValueState)
public class TxnCounter extends KeyedProcessFunction<String, Txn, Alert> {
    private transient ValueState<Long> countState;

    @Override
    public void open(Configuration cfg) {
        // state được engine quản lý & checkpoint; KHÔNG phải biến thường
        countState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("txn-count", Long.class));
    }

    @Override
    public void processElement(Txn txn, Context ctx, Collector<Alert> out) throws Exception {
        Long c = countState.value();          // đọc state của đúng key hiện tại
        c = (c == null ? 0L : c) + 1;
        countState.update(c);                 // ghi lại — sẽ được checkpoint
        if (c > 10) out.collect(new Alert(txn.accountId, c)); // luật minh hoạ
    }
}

Điểm cần nhớ: state phải khai báo qua API state của engine (như getState(...)), không phải một biến instance thường. Chỉ state do engine quản lý mới được đưa vào checkpoint và khôi phục được.

Checkpoint: chụp ảnh nhất quán một hệ phân tán đang chạy

Vấn đề nan giải: job gồm nhiều toán tử, chạy song song trên nhiều máy, dữ liệu vẫn đang chảy. Làm sao chụp một "ảnh" của toàn bộ state cộng vị trí đọc nguồn sao cho nhất quán — nghĩa là ảnh đó tương ứng với đúng một điểm cắt hợp lệ trong dòng sự kiện, chứ không phải một mớ chắp vá (toán tử A đã xử lý sự kiện #100 nhưng toán tử B mới tới #60)? Dừng cả cluster để chụp thì mất throughput.

Flink giải bài này bằng thuật toán barrier bất đồng bộ (asynchronous barrier snapshotting) — một biến thể của thuật toán chụp ảnh phân tán Chandy–Lamport cổ điển. Ý tưởng: chèn các checkpoint barrier vào chính dòng dữ liệu, trôi cùng sự kiện; khi barrier đi qua một toán tử, toán tử đó chụp state của mình. State không dừng dòng chảy.

Luồng một checkpoint:

  1. Coordinator tiêm barrier: JobManager định kỳ (ví dụ mỗi 30 giây) bảo các source chèn barrier n vào dòng. Source đồng thời ghi lại vị trí đọc hiện tại (Kafka offset) vào snapshot — đây là "state" của source.
  2. Barrier trôi theo dòng: barrier đi cùng dữ liệu qua từng toán tử. Mọi sự kiện trước barrier thuộc checkpoint n; sự kiện sau thuộc checkpoint sau. Barrier chia dòng thành "trước" và "sau".
  3. Toán tử chụp khi barrier tới: khi barrier n tới một toán tử, nó bất đồng bộ snapshot state của mình (đẩy bản sao RocksDB/heap ra durable store) rồi cho barrier chảy tiếp. Việc snapshot chạy nền, không chặn xử lý sự kiện mới lâu.
  4. Barrier alignment (với toán tử nhiều đầu vào): một toán tử join có 2 luồng vào; nó phải chờ barrier n từ cả hai luồng rồi mới chụp — để ảnh nhất quán. Trong khi chờ, luồng nào tới barrier trước sẽ bị "giữ" (buffer) các sự kiện sau barrier. Alignment đảm bảo nhất quán nhưng thêm độ trễ khi có nghẽn; Flink có unaligned checkpoint để đổi nhất quán-alignment lấy tốc độ khi backpressure nặng.
  5. Hoàn tất: khi tất cả toán tử (đến tận sink) ack barrier n, coordinator đánh dấu checkpoint n hoàn tất và bền vững. Đây là một điểm phục hồi hợp lệ: {offset nguồn + state mọi toán tử} khớp đúng một điểm cắt trong dòng.

Vì sao offset nguồn phải nằm trong checkpoint: phục hồi = khôi phục state tua nguồn về đúng offset đã lưu, rồi phát lại từ đó. Nếu state và offset lệch nhau, phát lại sẽ đếm thiếu hoặc thừa. Chúng phải được chụp cùng một ảnh nhất quán.

Checkpoint vs Savepoint

Cùng cơ chế, khác mục đích:

CheckpointSavepoint
Ai kích hoạtEngine, tự động, định kỳNgười vận hành, thủ công
Mục đíchPhục hồi sau lỗiNâng cấp code, đổi parallelism, di trú cluster
Vòng đờiEngine tự dọn bản cũNgười dùng giữ, có chủ đích
Định dạng/ổn địnhTối ưu tốc độ, có thể gắn engineỔn định, portable qua phiên bản

Nói ngắn: checkpoint để máy tự chữa lành; savepoint để con người chủ động thay đổi job mà không mất trí nhớ. Trước khi deploy phiên bản luật chống gian lận mới, ta trigger savepoint, dừng job, khởi động job mới từ savepoint — mọi cửa sổ đang mở, mọi tổng đang chạy được mang nguyên sang.

Từ checkpoint tới exactly-once

Có checkpoint mới chỉ giải quyết state không mất. Nhưng khi phục hồi, engine tua nguồn về offset đã lưu và phát lại các sự kiện từ đó — nghĩa là một số sự kiện đã được xử lý và có thể đã ghi ra sink trước khi crash sẽ được xử lý lại. Nếu không xử lý khéo, ta có bản ghi trùng ở đích. Đây chính là ranh giới at-least-once vs exactly-once.

  • At-least-once: mỗi sự kiện được xử lý ít nhất một lần; phát lại sau phục hồi có thể tạo trùng ở sink. Rẻ, đơn giản, chấp nhận được cho metric best-effort.
  • Exactly-once (effect): hiệu ứng lên kết quả cuối xảy ra đúng một lần, dù sự kiện có thể được xử lý lại. Không có nghĩa "không bao giờ xử lý lại" — mà là làm cho lần xử lý lại không đổi kết quả cuối.

Bên trong engine, state của Flink là exactly-once nhờ chính cơ chế checkpoint (khôi phục state khớp offset). Vấn đề khó nằm ở biên với thế giới bên ngoài — cái sink ta ghi ra. Có hai con đường để sink đạt exactly-once:

Con đường 1 — Idempotent sink

Làm cho việc ghi lặp không đổi trạng thái cuối: ghi cùng dữ liệu hai lần = ghi một lần. Đây đúng là công thức từ Batch 4 — Idempotency, áp cho streaming:

at-least-once delivery + sink idempotent = exactly-once effect.

Cách làm: sink có khóa duy nhất tất định và dùng upsert theo key thay vì append. Ví dụ ghi vào một key-value store hay bảng có primary key: PUT(key=account_id, value=running_total) — phát lại ghi đè cùng key với cùng giá trị, vô hại. Hoặc bảng đích với UNIQUE(txn_id) + ON CONFLICT DO NOTHING. Không cần giao dịch phân tán, đơn giản và mạnh — nhưng chỉ dùng được khi kết quả biểu diễn được dưới dạng upsert theo khóa (không hợp với sink chỉ-append thuần như file log).

Con đường 2 — Two-phase commit (transactional sink)

Khi sink phải append và không thể upsert (ví dụ ghi ra một topic Kafka khác, hay một hệ transactional), ta cần buộc việc ghi ra sink hoàn tất nguyên tử cùng lúc với checkpoint. Cơ chế: two-phase commit (2PC), dùng chính barrier checkpoint làm tín hiệu điều phối. Flink hiện thực qua TwoPhaseCommitSinkFunction; Kafka cung cấp transactions (producer transaction + read_committed consumer) để hiện thực.

Logic 2PC gắn với checkpoint:

  1. Pre-commit (pha 1): khi barrier n tới, sink flush dữ liệu của giai đoạn này vào một giao dịch đang mở ở hệ đích, nhưng chưa commit. Đồng thời state toán tử được snapshot. Dữ liệu đã ra đích nhưng còn "ẩn" — consumer đọc ở mức read_committed chưa nhìn thấy.
  2. Commit (pha 2): chỉ khi coordinator xác nhận checkpoint n hoàn tất trên toàn job, nó gọi notifyCheckpointComplete, và sink mới commit giao dịch. Lúc này dữ liệu mới hiện ra cho downstream.

Vì sao đúng exactly-once:

  • Crash trước khi checkpoint hoàn tất: giao dịch chưa commit → khi phục hồi, engine abort giao dịch dang dở đó và phát lại từ checkpoint trước → dữ liệu "ẩn" bị vứt, không trùng.
  • Crash sau khi checkpoint hoàn tất nhưng trước khi commit tới đích: khi phục hồi, sink biết checkpoint n đã hoàn tất nên commit lại (resume) giao dịch đúng của n → dữ liệu hiện đúng một lần. Đây là lý do commit phải idempotent/khôi phục được theo transaction id ổn định.

Điểm cần lưu ý thực tế: Kafka có transaction.max.timeout.ms — nếu khoảng giữa hai checkpoint dài hơn timeout giao dịch, giao dịch pre-commit có thể bị broker huỷ → mất dữ liệu. Nên checkpoint interval phải nhỏ hơn transaction timeout. (Đây là cấu hình định tính cần cân chỉnh theo cụm; xem docs Flink/Kafka để lấy giá trị đúng phiên bản.)

So sánh nhanh: idempotent sink đơn giản, không cần giao dịch phân tán, nhưng đòi kết quả upsert-được. 2PC tổng quát hơn (append-only, cross-system), nhưng phức tạp, thêm độ trễ end-to-end (downstream chỉ thấy dữ liệu sau khi checkpoint hoàn tất — độ trễ bằng bội của checkpoint interval) và đòi hệ đích hỗ trợ transaction.

Recovery: khôi phục rồi đi tiếp

Khi một task chết, engine không chạy lại từ đầu dòng vô hạn — nó phục hồi từ checkpoint hoàn tất gần nhất:

Ba mảnh phải khớp nhau như một: (1) state khôi phục về đúng ảnh checkpoint n, (2) nguồn tua về đúng offset của n, (3) sink xử lý phần phát lại một cách idempotent hoặc qua 2PC. Thiếu bất kỳ mảnh nào là mất exactly-once. Cũng vì cần tua lại nguồn, nguồn phải replayable (đọc lại được từ offset cũ) — Kafka làm được nhờ log giữ lại; một socket TCP thuần thì không, nên không thể exactly-once với nguồn không phát lại được.

Bối cảnh rộng hơn: exactly-once thượng nguồn còn dựa vào CDC/log nguồn tin cậy — xem Streaming 7 — CDC & Event-driven về Debezium và log-based CDC làm nguồn phát-lại-được cho pipeline. Và ranh giới event time vs processing time ảnh hưởng cái gì được tính lại khi phát lại — xem Streaming 3 — Event & Processing time.

Use case thực tế

Bối cảnh (số liệu minh hoạ): NCB chạy một job Flink phát hiện gian lận thẻ real-time: đọc luồng giao dịch từ Kafka (~8.000 sự kiện/giây giờ cao điểm), giữ keyed state theo card_id — cửa sổ trượt 5 phút đếm số lần và tổng tiền, cộng bảng dedup txn_id. State tổng ~120 GB, dùng RocksDB backend, checkpoint mỗi 30 giây ra S3. Kết quả (cảnh báo + giao dịch làm giàu) ghi ra một topic Kafka fraud_alerts mà hệ chặn thẻ tiêu thụ.

  • Sự cố: 20:14, một node TaskManager bị OOM-kill giữa giờ cao điểm.
  • Nếu at-least-once + sink thường: phục hồi tua Kafka về offset checkpoint gần nhất và phát lại ~30 giây sự kiện. Các cảnh báo trong 30 giây đó bị ghi trùng vào fraud_alerts → hệ chặn thẻ nhận lệnh chặn hai lần cho cùng giao dịch, có nguy cơ khoá nhầm/đếm sai tỷ lệ fraud.
  • Sau khi bật exactly-once (2PC + Kafka transactions): khi node chết, Flink nạp lại 120 GB state từ checkpoint S3 vào node thay thế, tua Kafka về offset đã lưu, abort giao dịch fraud_alerts dang dở của checkpoint chưa hoàn tất, rồi phát lại. Downstream đọc read_committed nên không bao giờ thấy bản ghi của giao dịch bị abort. Tổng cục: mỗi cảnh báo hiện đúng một lần, hệ chặn thẻ không nhận lệnh trùng. Đánh đổi: độ trễ end-to-end tăng thêm cỡ một checkpoint interval (~vài chục giây tối đa) vì downstream chỉ thấy dữ liệu sau khi checkpoint commit — chấp nhận được cho luồng chặn thẻ.
  • Vận hành: khi cập nhật luật fraud, đội SRE trigger savepoint, dừng job, khởi động phiên bản mới từ savepoint — mọi cửa sổ 5 phút đang mở được mang nguyên, không "quên" lịch sử vừa tích luỹ, không bỏ sót gian lận ngay sau khi deploy.

Ghi nhớ

  • Stateful là bản chất của tổng hợp/join/dedup/pattern trên stream: toán tử phải nhớ quá khứ giữa các sự kiện; state là trái tim của engine streaming.
  • State sống local (heap hoặc RocksDB cho state lớn hơn RAM) để đọc-ghi nóng nhanh; state backend/durable store (S3/HDFS) chỉ giữ bản sao để phục hồi. RocksDB là store nhúng trên đĩa local, không phải DB ngoài.
  • Keyed state: keyBy phân vùng theo key → mỗi task giữ state cho tập key của mình → song song không tranh chấp và rescale được (key groups).
  • Checkpoint = ảnh nhất quán của {state mọi toán tử + offset nguồn}, chụp bằng barrier bất đồng bộ kiểu Chandy–Lamport: barrier trôi trong dòng, toán tử chụp khi barrier tới; toán tử nhiều đầu vào phải align barrier để nhất quán.
  • Savepoint cùng cơ chế nhưng do người vận hành trigger, để nâng cấp code/đổi parallelism mà không mất trí nhớ; checkpoint là để máy tự phục hồi.
  • at-least-once + sink idempotent = exactly-once effect — không phải "không xử lý lại", mà là làm lần xử lý lại không đổi kết quả cuối.
  • Hai con đường exactly-once ở sink: idempotent (upsert theo khóa tất định, đơn giản) hoặc two-phase commit (pre-commit khi barrier tới, commit khi checkpoint hoàn tất toàn cục — dùng Kafka transactions; tổng quát nhưng thêm độ trễ).
  • Recovery khớp ba mảnh: nạp state từ checkpoint + tua nguồn về offset đã lưu + sink idempotent/2PC. Nguồn phải replayable (Kafka) mới đạt exactly-once.

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ẻ!