Streaming 2 — Kiến trúc Lambda vs Kappa
Mô hình tinh thần: hai cách trả lời câu hỏi "vừa nhanh vừa đúng"
Mọi hệ thống dữ liệu real-time đều va phải một mâu thuẫn cốt lõi: kết quả nhanh thường không đầy đủ, kết quả đầy đủ thường không nhanh. Bản ghi đến muộn (late data), sự cố khiến job dừng giữa chừng, logic có bug phải tính lại — nếu bạn chỉ có một đường stream chạy liên tục thì rất khó "tính lại toàn bộ lịch sử cho đúng". Nhưng nếu bạn chỉ có batch thì lại quá chậm cho những câu hỏi cần trả lời trong vài giây.
Có hai trường phái kiến trúc trả lời mâu thuẫn này theo hai hướng ngược nhau:
- Lambda (Nathan Marz): chạy song song hai đường — một đường batch cho tính đúng, một đường stream cho tính nhanh — rồi hợp nhất kết quả khi truy vấn.
- Kappa (Jay Kreps): chỉ giữ một đường stream duy nhất, và khi cần "tính lại cho đúng" thì phát lại (replay) toàn bộ lịch sử từ log.
Bài này là bài 2 của series Streaming & Real-time Data chuyên sâu. Ta sẽ mổ xẻ từng kiến trúc, chỉ ra nỗi đau thật của Lambda (trùng lặp logic hai nơi), cách Kappa né được nó nhờ replay, khi nào chọn cái nào, và tại sao ngành đang hội tụ về một mô hình thống nhất batch-stream.
Lambda Architecture — hai đường song song
Lambda Architecture do Nathan Marz đề xuất (khoảng 2011, sau đúc kết trong sách Big Data, Marz & Warren). Ý tưởng: dữ liệu thô là master dataset bất biến, chỉ ghi thêm (append-only), và mọi view phục vụ truy vấn đều là hàm thuần tính từ dữ liệu thô đó. Vì tính lại toàn bộ thì chậm, Marz chia thành ba tầng:
| Tầng | Nhiệm vụ | Đặc tính |
|---|---|---|
| Batch layer | Giữ master dataset bất biến, định kỳ tính lại batch view trên toàn bộ lịch sử | Chính xác, đầy đủ, độ trễ cao (giờ) |
| Speed layer | Bù đắp khoảng thời gian batch chưa kịp xử, tính real-time view trên dữ liệu mới nhất | Nhanh (giây), gần đúng, tạm thời |
| Serving layer | Nhận truy vấn, hợp nhất batch view + real-time view trả kết quả | Đọc thấp độ trễ |
Nguyên tắc vận hành: batch layer là nguồn chân lý. Speed layer chỉ lấp "khoảng trống" từ lần chạy batch gần nhất đến hiện tại. Khi batch chạy xong một chu kỳ mới và bao phủ luôn khoảng thời gian đó, real-time view tương ứng bị vứt đi — sai số của speed layer chỉ là tạm thời và tự động biến mất. Đây gọi là tính chất "eventual accuracy": sai lệch của tầng nhanh luôn được tầng chậm sửa lại.
Ưu điểm của Lambda
- Chống lỗi con người: master dataset bất biến, nên bug logic hay dữ liệu hỏng đều sửa được bằng cách tính lại batch từ đầu — không mất dữ liệu gốc.
- Cân bằng độ trễ/chính xác: có ngay câu trả lời gần đúng (speed) và câu trả lời chuẩn (batch) mà không phải hy sinh cái nào.
- Batch layer đơn giản để suy luận: xử lý dữ liệu bounded, dễ test, dễ backfill (xem Batch 1 — tổng quan).
Nhược điểm cốt lõi: trùng lặp logic ở hai nơi
Đây là điểm mà mọi phê phán Lambda đều xoáy vào. Cùng một logic nghiệp vụ (ví dụ "đếm số giao dịch nghi ngờ gian lận theo khách hàng") phải được cài đặt hai lần: một lần trong hệ batch (Spark, SQL), một lần trong hệ stream (Flink, Kafka Streams), bằng hai mô hình lập trình khác nhau.
# Cùng một bài toán "đếm giao dịch/khách/ngày" — Lambda viết 2 bản:
# Batch layer (SQL trên kho dữ liệu, chạy EOD)
SELECT customer_id, DATE(txn_ts) d, COUNT(*) c
FROM transactions
GROUP BY customer_id, DATE(txn_ts);
# Speed layer (Kafka Streams / Flink, chạy liên tục)
stream.groupBy(txn -> txn.customerId)
.windowedBy(TimeWindows.of(Duration.ofDays(1)))
.count(); # ← LOGIC SONG SONG, dễ trôi lệch khỏi bản batch
Hệ quả thực tế:
- Chi phí bảo trì gấp đôi: mỗi thay đổi nghiệp vụ phải sửa hai codebase, và phải giữ cho chúng cho kết quả khớp nhau — nếu không, batch view và real-time view mâu thuẫn tại điểm giao.
- Hai bộ kỹ năng, hai cụm hạ tầng: đội ngũ phải giỏi cả batch lẫn stream; vận hành hai hệ phân tán riêng.
- Khó debug: khi số liệu lệch, phải truy xem lệch do batch hay do speed hay do khâu hợp nhất ở serving layer.
Kappa Architecture — chỉ một đường stream + replay
Năm 2014, Jay Kreps (đồng tác giả Apache Kafka) viết bài kinh điển "Questioning the Lambda Architecture", lập luận: sự phức tạp của Lambda không phải bản chất bài toán mà là cái giá của việc duy trì hai hệ. Ông đề xuất bỏ hẳn batch layer — kiến trúc sau này cộng đồng gọi là Kappa.
Ý tưởng chìa khoá: coi log (Kafka) là nguồn chân lý bền vững, giữ lại đủ lâu (hoặc vô hạn). Khi đó ta chỉ cần một engine stream. Câu hỏi "làm sao tính lại cho đúng khi có bug?" — thứ mà Lambda dùng batch layer để giải — được Kappa trả lời bằng reprocessing qua replay:
Quy trình reprocessing của Kappa (không có batch layer riêng):
- Khi cần sửa logic hoặc backfill, khởi động một job stream mới (v2) với code mới, cho nó đọc lại log từ offset 0 (hoặc từ mốc cần thiết) và ghi ra một output table mới.
- Job v2 chạy hết lịch sử rồi bắt kịp dòng dữ liệu hiện tại — vì replay chạy nhanh hơn thời gian thực.
- Khi v2 đã ngang hàng, chuyển truy vấn sang output v2, tắt v1 và xoá output cũ.
Điểm tinh tế: cùng một mã stream vừa xử lý real-time vừa dùng để tính lại lịch sử — chỉ khác điểm bắt đầu đọc (offset). Không còn hai bản logic. Đây chính là thứ Kappa mua được bằng cách né bỏ batch layer. Ý tưởng "log là nền tảng" được Kreps khai triển kỹ trong bài "The Log" — nền cho tư duy stream-table duality.
Điều kiện để Kappa khả thi
Kappa không miễn phí. Nó đòi hỏi:
- Log giữ đủ dữ liệu: Kafka phải retention dài (hoặc dùng tiered storage / log compaction) để replay được toàn bộ lịch sử cần thiết. Nếu chỉ giữ 7 ngày thì không thể tính lại 2 năm.
- Engine stream đủ mạnh: hỗ trợ state lớn, checkpoint và exactly-once để replay cho kết quả tất định, khớp với lần chạy trước.
- Sink chịu được double-write: trong lúc v1 và v2 chạy song song, hạ tầng lưu trữ phải chứa được hai bản output.
Khi nào dùng Lambda, khi nào dùng Kappa
Không có bên "thắng tuyệt đối". Chọn theo bản chất bài toán:
| Tiêu chí | Nghiêng về Lambda | Nghiêng về Kappa |
|---|---|---|
| Nguồn chân lý | Đã có kho/lake batch khổng lồ, lịch sử nằm ở đó | Log (Kafka) giữ được lịch sử đủ dài |
| Logic batch vs stream | Khác nhau bản chất (ví dụ ML training nặng chạy batch) | Gần như đồng nhất, chỉ khác độ trễ |
| Chi phí lưu trữ log | Không muốn giữ log dài | Chấp nhận retention dài / tiered storage |
| Độ chín của đội stream | Đội mạnh batch, stream còn non | Đội vững streaming, exactly-once |
| Ưu tiên | Tận dụng hệ batch sẵn có | Giảm tối đa trùng lặp logic, một codebase |
Một cách nói gọn: Lambda phù hợp khi bạn đã có sẵn thế giới batch và chỉ thêm một lớp real-time; Kappa phù hợp khi bạn thiết kế mới, coi stream là trung tâm và muốn tránh nợ kỹ thuật của hai codebase. Trong thực tế nhiều hệ ngân hàng là hybrid: fraud/cảnh báo real-time theo tinh thần Kappa, còn báo cáo tài chính/phân loại nợ EOD vẫn thuần batch.
Xu hướng hội tụ batch-stream (unified)
Phê phán sâu nhất với Lambda dẫn tới một nhận thức quan trọng: batch không phải một phạm trù tách biệt, mà là một trường hợp đặc biệt của streaming — xử lý một dòng dữ liệu bounded (có điểm kết thúc) thay vì unbounded. Nếu engine coi mọi thứ là stream, thì "batch" chỉ là "stream có kết thúc".
Đây là luận điểm trung tâm của mô hình Dataflow (Akidau et al., Google, VLDB 2015) — nền tảng của Apache Beam — và của sách "Streaming Systems" (Akidau, Chernyak, Lax). Hệ quả kiến trúc:
- Apache Beam: một API duy nhất mô tả pipeline, chạy được cả ở chế độ batch lẫn stream trên nhiều runner (Flink, Dataflow, Spark). "Write once, run batch or streaming".
- Apache Flink: coi batch là stream bounded, cùng một runtime; DataStream API và bảng/SQL thống nhất hai chế độ.
- Apache Spark Structured Streaming: chính là micro-batch/continuous trên cùng engine và cùng API DataFrame với batch (xem Structured Streaming).
- Delta/Iceberg: bảng đọc-ghi được bởi cả job batch lẫn stream, giúp reprocessing kiểu Kappa trở nên thực dụng ngay trên lakehouse (xem streaming CDC trên Iceberg).
Khi engine đã thống nhất, ranh giới Lambda/Kappa mờ đi: bạn viết logic một lần, chạy chế độ streaming cho real-time, chạy chế độ batch (trên chính pipeline đó) để backfill/tính lại — được cái tính-lại-cho-đúng của Lambda mà không phải viết logic hai lần như nhược điểm của Lambda. Đây chính là lời hứa gốc của Kappa được hiện thực hoá bằng công cụ chín hơn.
Reprocessing qua replay — trái tim của cả hai
Dù Lambda (tính lại ở batch layer) hay Kappa (replay job mới từ log), năng lực nền tảng đều là reprocessing: khả năng tạo lại output từ dữ liệu gốc bất biến. Điều kiện để reprocessing đúng:
- Dữ liệu gốc bất biến, đánh version được: master dataset (Lambda) hoặc log giữ nguyên bản ghi theo offset (Kappa).
- Logic tất định (deterministic): cùng input + cùng code phải ra cùng output. Tránh phụ thuộc
now(), random, thứ tự không ổn định. - Output ghi ra nơi mới rồi swap: không sửa tại chỗ output đang phục vụ; tính bản mới xong mới chuyển — tránh người dùng thấy trạng thái nửa vời.
- Xử lý late data nhất quán: reprocessing phải áp cùng quy tắc watermark/lateness như luồng gốc để số liệu khớp (chủ đề bài sau trong series).
Use case thực tế
Số liệu dưới đây là minh hoạ, không phải số thật của NCB.
Giả sử NCB xây hệ cảnh báo gian lận thẻ real-time đếm số giao dịch bất thường theo khách trong cửa sổ 24h.
- Giai đoạn đầu (Lambda): đội đã có sẵn kho dữ liệu và job Spark EOD tính "hồ sơ rủi ro khách hàng". Họ thêm speed layer bằng Kafka Streams để cảnh báo trong ~2 giây, còn batch layer tính lại mỗi đêm cho chuẩn. Vấn đề phát sinh: một lần đổi định nghĩa "giao dịch nghi ngờ", đội phải sửa hai codebase; suốt 3 ngày real-time view và batch view lệch nhau ~4% tại điểm giao vì bản stream bị sửa sau.
- Chuyển sang Kappa/unified: nhóm chuyển sang Flink, giữ Kafka retention 30 ngày + tiered storage cho lịch sử. Khi cần sửa logic, họ chạy một job Flink phiên bản mới đọc lại từ offset đầu, tính ra bảng
fraud_score_v2, mất ~40 phút để replay 30 ngày rồi bắt kịp real-time, sau đó chuyển truy vấn sang v2. Chỉ một bản logic, không còn lệch batch/stream. - Kết quả minh hoạ: thời gian đưa một thay đổi logic ra production giảm từ ~2 ngày (đồng bộ hai codebase) xuống còn một lần deploy + replay; số vụ số liệu lệch giữa hai tầng về 0 vì không còn hai tầng.
Ghi nhớ
- Lambda (Nathan Marz): batch layer (đúng, chậm) + speed layer (nhanh, tạm) + serving layer (hợp nhất). Master dataset bất biến; sai số speed layer là tạm thời, batch sửa lại.
- Nỗi đau lớn nhất của Lambda: trùng lặp logic hai nơi — cùng nghiệp vụ viết hai lần bằng hai mô hình, tốn gấp đôi bảo trì và dễ lệch kết quả.
- Kappa (Jay Kreps, "Questioning the Lambda Architecture"): bỏ batch layer, chỉ một engine stream; tính lại bằng replay log (Kafka) với job phiên bản mới rồi swap output.
- Điều kiện Kappa: log giữ đủ lịch sử, engine hỗ trợ state + exactly-once tất định, sink chịu double-write.
- Chọn cái nào: Lambda khi đã có hệ batch lớn/logic hai bên khác bản chất; Kappa khi coi stream là trung tâm và muốn một codebase.
- Hội tụ batch-stream: streaming là tổng quát của batch (batch = stream bounded); Beam/Flink/Spark cho một logic chạy cả hai chế độ — hiện thực hoá lời hứa Kappa.
- Nền tảng chung của cả hai là reprocessing: dữ liệu gốc bất biến + logic tất định + ghi output mới rồi swap.
Nguồn tham khảo
- Jay Kreps — "Questioning the Lambda Architecture" (O'Reilly Radar, 2014).
- Jay Kreps — "The Log: What every software engineer should know about real-time data's unifying abstraction" (LinkedIn Engineering, 2013).
- Nathan Marz & James Warren — "Big Data: Principles and best practices of scalable realtime data systems" (Manning).
- Tyler Akidau, Slava Chernyak, Reuven Lax — "Streaming Systems" (O'Reilly).
- Akidau et al. — "The Dataflow Model" (Google, VLDB 2015).
- Martin Kleppmann — "Designing Data-Intensive Applications" (O'Reilly), chương xử lý batch & stream.
- Apache Flink — Unified batch/stream processing: https://flink.apache.org/
- Apache Beam — Programming model: https://beam.apache.org/documentation/basics/
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ẻ!