Streaming 1 — Xử lý luồng & vì sao Real-time
Mô hình tinh thần: streaming là "xử ngay khi sự kiện tới, trên một dòng không có điểm cuối"
Trước khi đi vào định nghĩa, hãy giữ một hình ảnh: nếu batch là gom cả ngày giao dịch rồi tối đến tính sổ một lượt, thì streaming là một nhân viên đứng ngay tại cửa băng chuyền, cầm từng gói hàng lên xử lý ngay khi nó trôi qua — băng chuyền không bao giờ dừng, không có "gói cuối cùng". Bạn không chờ dữ liệu tích luỹ thành khối rồi khoá lại; bạn phản ứng với từng sự kiện ngay khi nó xuất hiện, và làm việc đó mãi mãi.
Đây là sự khác biệt bản chất nhất: batch xử bounded data (khối hữu hạn, đã đóng), streaming xử unbounded data (dòng vô tận, không bao giờ kết thúc). Mọi khó khăn và mọi giá trị của streaming đều mọc ra từ chữ unbounded này.
Trong ngân hàng, streaming là địa hạt của những thứ không được phép chờ đến đêm: chặn giao dịch thẻ nghi gian lận trong vài trăm mili-giây, cảnh báo thấu chi tức thời, cập nhật dashboard rủi ro theo phút, cá nhân hoá ưu đãi ngay khi khách vừa quẹt thẻ.
Bài này là bài mở màn cho series Streaming & Real-time Data chuyên sâu gồm 9 bài. Mục tiêu: cho bạn bản đồ tổng thể — streaming là gì, vì sao cần real-time, khác batch ở đâu, những khái niệm cốt lõi và thách thức trung tâm — để các bài sau đào sâu từng mảnh.
Series này nói về nguyên lý cross-cutting, không phụ thuộc một công cụ. Khi cần công cụ cụ thể, ta cross-link tới series chuyên sâu: Kafka (nền tảng log/message), Flink (stream processor), Spark Structured Streaming (micro-batch). Song song có series xử lý Batch để bạn thấy nửa còn lại của bức tranh.
Stream processing là gì
Stream processing là mô hình xử lý dữ liệu liên tục, từng sự kiện một, ngay khi sự kiện tới, trên một nguồn dữ liệu được coi là vô hạn (unbounded). Ba đặc điểm định nghĩa:
- Dữ liệu là một dòng vô tận: sự kiện đến liên tục, không có điểm kết thúc. Chương trình streaming về nguyên tắc chạy mãi (long-running), không "xong rồi thoát" như một job batch.
- Xử lý được kích hoạt bởi sự xuất hiện của dữ liệu: mỗi khi có sự kiện mới, logic được chạy — không đợi lịch, không đợi "đủ file".
- Độ trễ thấp là mục tiêu thiết kế: từ lúc sự kiện xảy ra đến lúc có kết quả thường tính bằng mili-giây → giây, chứ không phải giờ/ngày.
Điểm mấu chốt phân biệt với batch nằm ở biên dữ liệu. Vì dòng không bao giờ đóng, bạn không thể "chờ có đủ dữ liệu rồi tính" — bạn buộc phải trả kết quả tăng dần (incremental) trên dữ liệu đang chảy, và phải tự định nghĩa "một khoảng thời gian" bằng cửa sổ / windowing thay vì có sẵn một lô đóng.
Vì sao real-time: giá trị dữ liệu giảm theo thời gian
Câu hỏi quan trọng nhất không phải "streaming làm được gì" mà là "vì sao phải trả giá phức tạp để có real-time". Câu trả lời cốt lõi là time value of data — giá trị hành động được của một dữ liệu giảm theo thời gian kể từ khi sự kiện xảy ra.
Một giao dịch thẻ nghi gian lận: nếu bạn phát hiện trong 200 mili-giây, bạn chặn được nó — giá trị cực cao. Phát hiện sau 2 giờ, tiền đã ra khỏi hệ thống, bạn chỉ còn điều tra và bồi hoàn — giá trị thấp hơn nhiều. Phát hiện trong báo cáo EOD sáng hôm sau, gần như chỉ còn giá trị thống kê. Cùng một sự thật, giá trị hành động rơi theo thời gian.
Các use case mà real-time tạo ra giá trị không thể thay thế bằng batch:
- Phát hiện & chặn gian lận (fraud): chấm điểm rủi ro từng giao dịch thẻ/chuyển khoản trước khi phê duyệt. Đây là use case ngân hàng kinh điển của streaming — xem Streaming 9 — ngân hàng real-time.
- Cảnh báo & giám sát (alerting/monitoring): thấu chi, số dư chạm ngưỡng, đăng nhập bất thường, giao dịch từ địa điểm lạ — báo ngay cho khách/đội vận hành.
- Dashboard real-time: theo dõi dòng tiền, khối lượng giao dịch, tình trạng hệ thống theo giây/phút thay vì chờ báo cáo cuối ngày.
- Cá nhân hoá tức thời (real-time personalization): gợi ý ưu đãi/sản phẩm ngay trong ngữ cảnh khách vừa hành động (vừa quẹt thẻ tại cửa hàng → gợi ý hoàn tiền).
Nguyên tắc chọn: real-time có giá trị khi độ trễ đổi được thành hành động. Nếu không ai làm gì khác nhờ biết sớm hơn, thì trả tiền cho streaming là lãng phí — cứ dùng batch.
Streaming vs batch: cùng một trục đánh đổi
Streaming và batch không phải "cái nào tốt hơn" mà nằm trên một trục đánh đổi giữa biên dữ liệu, độ trễ, thông lượng/chi phí và độ phức tạp vận hành. Giữa hai cực còn có micro-batch (cắt dòng thành lô rất nhỏ, xử liên tục) — mô hình của Spark Structured Streaming, là cầu nối tư duy batch sang streaming.
| Tiêu chí | Batch | Micro-batch | Streaming thuần |
|---|---|---|---|
| Biên dữ liệu | bounded | "cắt lát" unbounded | unbounded |
| Độ trễ điển hình | phút → giờ → ngày | giây → phút | mili-giây → giây |
| Thông lượng/chi phí | cao nhất | cao | thấp hơn (state tốn kém) |
| Mô hình lập trình | đơn giản nhất | tư duy batch lặp nhanh | phức tạp (state, thời gian) |
| Tái chạy (re-run) | rất dễ (lô bất biến) | vừa | khó (cần replay + state) |
| Xử lý thời gian | có sẵn lô đóng | ranh giới lô | cần watermark, windowing |
| Ví dụ ngân hàng | EOD, sao kê, phân loại nợ | dashboard cập nhật mỗi phút | chặn gian lận thẻ tức thời |
Ba khác biệt bản chất (không chỉ là "nhanh hơn"):
- Bounded vs unbounded: batch có điểm đầu–cuối rõ nên kết quả tất định; streaming không bao giờ "thấy hết dữ liệu" nên phải trả kết quả tăng dần và chấp nhận dữ liệu đến trễ sau khi đã tính.
- Độ trễ đổi lấy độ phức tạp: đẩy độ trễ xuống mili-giây buộc bạn quản lý state, checkpoint, event time vs processing time, backpressure — những thứ batch không cần lo.
- Chạy một lần vs chạy mãi: job batch khởi động rồi kết thúc; ứng dụng streaming là dịch vụ long-running 24/7, nên fault-tolerance và khả năng khôi phục sau sự cố là yêu cầu bắt buộc, không phải tuỳ chọn.
Thực tế thường lai: streaming cho cảnh báo/gian lận, batch cho báo cáo quy định. Cách ghép hai nửa này (Lambda vs Kappa) là chủ đề của Streaming 2.
Khái niệm cốt lõi: event, stream, producer/consumer, log, offset
Toàn bộ thế giới streaming được dựng trên một bộ khái niệm nhỏ nhưng chặt chẽ. Nắm chắc bộ này là hiểu 80% mọi công cụ (Kafka, Flink, Pulsar...).
- Event (sự kiện): một sự thật đã xảy ra, bất biến, gắn mốc thời gian — ví dụ
{txn_id, account, amount, timestamp}. Event mô tả "đã có gì xảy ra", không phải "trạng thái hiện tại". Đây là đơn vị nguyên tử của streaming — Streaming 7 — CDC & event-driven đào sâu tư duy này. - Stream (luồng): một chuỗi event vô hạn, có thứ tự theo thời gian, thường được nhóm theo chủ đề (topic trong Kafka).
- Producer: bên sinh ra và ghi (append) event vào stream (ATM, POS, mobile app, CDC từ database).
- Consumer: bên đọc và xử lý event. Nhiều consumer đọc độc lập cùng một stream ở tốc độ khác nhau — fraud engine, dashboard, kho dữ liệu cùng ăn một nguồn.
- Log: cấu trúc dữ liệu nền tảng — một chuỗi append-only, có thứ tự, bất biến. Bạn chỉ thêm vào cuối, không sửa/xoá giữa chừng. Đây chính là "trái tim" của Kafka và là ý tưởng trong bài kinh điển "The Log" của Jay Kreps. Chi tiết ở Kafka — kiến trúc.
- Offset: vị trí đọc của một consumer trong log (số thứ tự bản ghi). Vì log bất biến, mỗi consumer chỉ cần nhớ offset của mình để biết đã xử đến đâu — và có thể "tua lại" (replay) bằng cách đặt offset về quá khứ. Đây là nền tảng cho khả năng khôi phục và tái xử lý.
Một đặc tính then chốt: log tách bạch việc ghi khỏi việc đọc. Producer không cần biết ai đọc; consumer đọc theo tốc độ của mình và tự quản offset. Đây là gốc rễ của kiến trúc event-driven phân rã (decoupled).
Bảo đảm giao hàng: at-least-once vs exactly-once
Trong một hệ phân tán có lỗi mạng, máy chết, retry — câu hỏi "mỗi event được xử bao nhiêu lần" là trung tâm của tính đúng đắn. Ba mức bảo đảm:
- At-most-once: mỗi event xử nhiều nhất một lần — có thể mất khi lỗi. Nhanh, đơn giản, nhưng hiếm khi chấp nhận được trong ngân hàng.
- At-least-once: mỗi event xử ít nhất một lần — không mất, nhưng có thể trùng (khi retry sau lỗi). Đây là mặc định phổ biến và rẻ nhất trong nhóm "không mất dữ liệu".
- Exactly-once: hiệu ứng đúng một lần — không mất, không trùng. Đắt và khó nhất về mặt kỹ thuật.
Điểm tinh tế thường bị hiểu sai: "exactly-once" thực tế là "exactly-once hiệu ứng", không phải "gửi đúng một gói mạng". Nó đạt được bằng một trong hai cách:
- Idempotent sink: thiết kế thao tác ghi sao cho ghi lại lần hai không đổi kết quả (ví dụ upsert theo khoá thay vì append) → cộng với at-least-once là đủ có hiệu ứng exactly-once.
- Transactional / two-phase commit: gắn việc dịch chuyển offset và việc ghi kết quả vào cùng một giao dịch nguyên tử, để không có trạng thái "đã ghi nhưng chưa commit offset".
Cơ chế đầy đủ (checkpoint, transactional producer, two-phase commit) là chủ đề riêng của Streaming 5 — State & exactly-once. Ở đây chỉ cần nhớ mô hình tinh thần: real-time nghiêng về at-least-once + sink idempotent, vì nó vừa an toàn vừa rẻ hơn exactly-once "cứng".
Thách thức trung tâm của streaming
Ba thách thức dưới đây là lý do streaming khó hơn batch — và là xương sống của cả series:
- Out-of-order & dữ liệu đến trễ (late data): event không đến theo đúng thứ tự xảy ra — mạng trễ, thiết bị offline rồi đồng bộ sau. Một giao dịch xảy ra lúc 10:00:00 có thể tới hệ thống lúc 10:00:07, sau giao dịch 10:00:05. Điều này buộc phải phân biệt event time vs processing time và dùng watermark để quyết định "chờ đến bao giờ thì chốt một cửa sổ" — xem Streaming 4 — Windowing.
- State (trạng thái): nhiều phép tính streaming cần nhớ thông tin qua nhiều event — đếm số giao dịch 5 phút gần nhất, số dư đang chạy, join hai luồng. State này phải sống trong bộ nhớ long-running, phải bền vững qua sự cố, và có thể rất lớn. Quản lý state là nội dung Streaming 5; quan hệ giữa "dòng event" và "bảng trạng thái" là Streaming 6 — stream-table duality.
- Fault-tolerance (chịu lỗi) cho dịch vụ 24/7: vì ứng dụng chạy mãi, máy sẽ chết trong lúc chạy. Hệ thống phải khôi phục state và tiếp tục từ đúng offset mà không mất, không trùng — thường bằng checkpoint định kỳ cộng khả năng replay từ log. Kèm theo là backpressure: khi consumer chậm hơn producer, hệ phải ghìm nhịp để không vỡ bộ nhớ.
Ba thách thức này đan vào nhau: xử lý out-of-order cần state (nhớ cửa sổ đang mở), và giữ state đúng qua sự cố cần fault-tolerance. Chúng là lý do vì sao không thể "batch nhưng nhanh hơn" là ra streaming.
Một pseudocode rất rút gọn cho khung một stateful streaming job (đếm giao dịch theo cửa sổ để chấm rủi ro), minh hoạ cả state lẫn event time:
# Khung stream job có state theo event-time window (minh hoạ, trung lập công cụ)
# Ý tưởng: đếm số giao dịch của mỗi tài khoản trong cửa sổ 5 phút (theo THỜI GIAN SỰ KIỆN)
def process(stream):
for event in stream: # dòng vô tận: vòng lặp không kết thúc
key = event.account_id
etime = event.event_time # dùng thời gian SỰ KIỆN, không phải lúc nhận
# WINDOWING theo event time (tumbling 5') — không dựa vào lúc bản ghi tới
window = floor_to_5min(etime)
# STATE: đếm luỹ kế theo (tài khoản, cửa sổ), bền vững qua checkpoint
state.increment(key=(key, window), by=1)
# WATERMARK: khi watermark vượt qua cuối cửa sổ -> chốt & phát kết quả
if watermark(stream) >= window_end(window):
count = state.get((key, window))
if count > FRAUD_THRESHOLD:
emit_alert(key, window, count) # sink idempotent theo (key, window)
# dữ liệu đến TRỄ sau đây được xử theo chính sách allowed_lateness
# checkpoint định kỳ (offset + state) để khôi phục exactly-once sau sự cố
Lưu ý kỹ thuật: mấu chốt của tính đúng đắn ở đây là dùng event_time (thời gian sự kiện xảy ra) chứ không phải thời gian nhận, và chốt cửa sổ theo watermark thay vì theo đồng hồ hệ thống — đúng hai chủ đề Streaming 3 và Streaming 4 sẽ mổ xẻ.
Bản đồ series: 9 bài
Lộ trình 9 bài của series:
- Streaming 1 — Xử lý luồng & vì sao Real-time (bài này) — bản đồ tổng thể, time value of data, khái niệm cốt lõi, thách thức.
- Streaming 2 — Lambda & Kappa — hai kiến trúc ghép batch + streaming, đánh đổi và khi nào dùng.
- Streaming 3 — Event time vs processing time — hai trục thời gian, vì sao out-of-order buộc phải phân biệt.
- Streaming 4 — Windowing & watermark — tumbling/sliding/session, watermark & lateness để chốt cửa sổ.
- Streaming 5 — State & exactly-once — quản lý state, checkpoint, transactional/idempotent sink.
- Streaming 6 — Stream–table duality — dòng event ↔ bảng trạng thái, changelog, materialized view.
- Streaming 7 — CDC & event-driven — biến thay đổi database thành luồng (Debezium), kiến trúc hướng sự kiện.
- Streaming 8 — Quality & observability — chất lượng dữ liệu luồng, lag/throughput, giám sát pipeline chạy 24/7.
- Streaming 9 — Ngân hàng real-time — ghép tất cả thành pipeline fraud/cảnh báo thời gian thực cho ngân hàng.
Use case thực tế
Bối cảnh (minh hoạ, số liệu giả định): NCB muốn chặn giao dịch thẻ nghi gian lận trước khi phê duyệt. Yêu cầu nghiệp vụ: mỗi giao dịch quẹt thẻ phải được chấm điểm rủi ro và trả quyết định cho phép/chặn trong < 300 mili-giây, vì cổng thanh toán không thể giữ khách chờ.
Vì sao đây là bài toán bắt buộc streaming, không batch nào thay được:
- Time value of data cực dốc: nếu chờ EOD mới phát hiện, tiền đã ra khỏi hệ thống — giá trị hành động gần như bằng 0. Chỉ real-time mới chặn được.
- Cần state theo thời gian thực: luật gian lận thường là "≥ 5 giao dịch trong 2 phút", "quẹt ở 2 quốc gia cách nhau < 10 phút" — cần nhớ lịch sử gần của từng thẻ trong bộ nhớ luồng, đúng phần state + windowing.
- Nhiều consumer trên một luồng: cùng dòng
card_transactions, fraud engine đọc để chấm điểm, dashboard đọc để hiển thị khối lượng, kho dữ liệu đọc để lưu lịch sử — nhờ log + offset độc lập. - Không được mất, hạn chế trùng: chọn at-least-once + sink idempotent (dedup theo
txn_id) — an toàn về tiền mà không phải trả giá exactly-once "cứng".
Con số minh hoạ: luồng ~3.000 giao dịch/giây giờ cao điểm, engine chấm điểm ở p99 ≈ 180ms, cảnh báo phát trong < 1 giây; song song, một job batch EOD vẫn chạy đêm để đối chiếu và huấn luyện lại mô hình — minh hoạ kiến trúc lai batch + streaming. Nếu ép bài toán này chạy batch, ngân hàng mất khả năng chặn — cái giá không đo bằng chi phí hạ tầng mà bằng tiền thất thoát và rủi ro tuân thủ.
Ghi nhớ
- Streaming = xử từng sự kiện ngay khi tới, trên dòng vô tận (unbounded); đối lập batch xử một lô hữu hạn (bounded) đã đóng. Mọi khó khăn và giá trị đều mọc từ chữ "unbounded".
- Vì sao real-time = time value of data: giá trị hành động của dữ liệu giảm theo thời gian. Chỉ trả tiền cho streaming khi biết sớm hơn đổi được thành hành động khác (chặn, cảnh báo, cá nhân hoá).
- Batch, micro-batch, streaming nằm trên trục biên dữ liệu ↔ độ trễ ↔ độ phức tạp — không có cái "tốt nhất", chỉ có cái hợp yêu cầu độ trễ.
- Bộ khái niệm cốt lõi: event (sự thật bất biến) → stream/topic (dòng có thứ tự) → producer/consumer (ghi/đọc tách bạch) → log (append-only bất biến) → offset (vị trí đọc, cho phép replay).
- At-least-once (không mất, có thể trùng) là mặc định rẻ và an toàn; exactly-once là hiệu ứng đúng một lần, đạt bằng idempotent sink hoặc transactional/two-phase commit.
- Ba thách thức trung tâm đan vào nhau: out-of-order/late data (cần event time + watermark), state (nhớ qua nhiều event), fault-tolerance cho dịch vụ 24/7 (checkpoint + replay + backpressure).
- Đừng ép streaming vào bài toán vốn là batch (EOD, sao kê) — sẽ trả giá bằng state/watermark mà không thêm giá trị; và đừng ép batch vào bài toán cần chặn tức thời (fraud) — sẽ mất giá trị hành động.
Nguồn tham khảo
- Streaming Systems — Tyler Akidau, Slava Chernyak, Reuben Lax (O'Reilly): mô hình bounded vs unbounded, event time/processing time, watermark — nền tảng lý thuyết của streaming hiện đại.
- Designing Data-Intensive Applications — Martin Kleppmann (O'Reilly): Chương 11 "Stream Processing" và ý tưởng log là cấu trúc nền tảng.
- Fundamentals of Data Engineering — Joe Reis & Matt Housley (O'Reilly): so sánh batch vs streaming và time value of data trong vòng đời dữ liệu.
- "The Log: What every software engineer should know about real-time data's unifying abstraction" — Jay Kreps (LinkedIn Engineering).
- "Questioning the Lambda Architecture" — Jay Kreps (O'Reilly Radar): góc nhìn về khi nào cần tách batch/streaming và ý tưởng Kappa.
- The Dataflow Model — Akidau et al., VLDB 2015: mô hình thống nhất xử lý bounded & unbounded, xử lý out-of-order với watermark.
- Apache Kafka — tài liệu chính thức (log, topic, offset, consumer group): https://kafka.apache.org/documentation/
- Apache Flink — tài liệu chính thức (event time, state, checkpoint, exactly-once): https://nightlies.apache.org/flink/flink-docs-stable/
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ẻ!