Streaming 4 — Windowing (cửa sổ thời gian)
Vì sao luồng vô hạn buộc phải có cửa sổ
Trong xử lý batch, câu hỏi "tổng số tiền giao dịch hôm nay là bao nhiêu?" có câu trả lời rõ ràng: đọc hết bảng của ngày, SUM, xong. Tập dữ liệu hữu hạn, có điểm đầu và điểm cuối, nên phép tổng hợp (aggregate) luôn hội tụ về một con số.
Trên một luồng (stream), dữ liệu không bao giờ kết thúc. Giao dịch thẻ, log đăng nhập, sự kiện từ core banking chảy về liên tục 24/7. Nếu bạn viết SELECT SUM(amount) FROM transactions trên một luồng, engine sẽ không bao giờ trả về kết quả — vì luôn còn dòng tiếp theo có thể tới. Aggregate trên tập vô hạn là một phép toán không hội tụ.
Windowing là cách giải quyết: thay vì tổng hợp trên "toàn bộ luồng", ta cắt luồng thành các cửa sổ hữu hạn theo thời gian (hoặc theo số lượng) và tổng hợp trong từng cửa sổ. "Tổng giao dịch" trở thành "tổng giao dịch mỗi 5 phút". Mỗi cửa sổ có ranh giới rõ ràng nên aggregate lại hội tụ — và ta có được một chuỗi kết quả liên tục theo thời gian.
Đây là ý tưởng nền tảng của toàn bộ stream analytics. Nhưng để làm đúng, ta cần trả lời hai câu hỏi tách biệt:
- Cắt cửa sổ theo tiêu chí nào? → các loại window.
- Khi nào phát ra kết quả của một cửa sổ? → trigger và mối liên hệ với watermark.
Bài này đi vào cả hai. Nó dựa trên Streaming 3 — event time, processing time & watermark (nếu bạn chưa nắm watermark, nên đọc trước) và dẫn sang Streaming 5 — state & exactly-once (cửa sổ chính là một dạng state có vòng đời).
Lưu ý sandbox: môi trường minh hoạ của tài liệu này dùng PostgreSQL, không có cluster Flink/Kafka. Toàn bộ mã Flink/Kafka Streams bên dưới là minh hoạ để đọc hiểu, không chạy trực tiếp trong sandbox.
Bốn loại window
Có bốn kiểu chia cửa sổ kinh điển, được hầu hết engine (Apache Flink, Apache Beam/Dataflow, Kafka Streams) hỗ trợ. Điểm khác biệt nằm ở kích thước cửa sổ và cách chúng chồng lấn nhau.
| Loại | Kích thước | Chồng lấn | Một sự kiện thuộc mấy cửa sổ | Ví dụ |
|---|---|---|---|---|
| Tumbling | Cố định | Không | Đúng 1 | Đếm giao dịch mỗi 5 phút |
| Sliding / Hopping | Cố định | Có | Nhiều (size / slide) | Trung bình trượt 10 phút, cập nhật mỗi 1 phút |
| Session | Thay đổi | Không | Đúng 1 | Gom hoạt động của một phiên đăng nhập |
| Global | Vô hạn | — | 1 (tất cả) | Đếm theo số lượng với trigger tuỳ biến |
Tumbling window (cố định, không chồng)
Cửa sổ tumbling (còn gọi fixed window trong Beam) chia trục thời gian thành các khoảng bằng nhau, liền kề, không chồng lấn: [00:00, 00:05), [00:05, 00:10), ... Mỗi sự kiện rơi vào đúng một cửa sổ dựa trên timestamp của nó.
Đây là loại phổ biến nhất và là mặc định tinh thần khi ai đó nói "mỗi 5 phút". Ranh giới cửa sổ được căn theo epoch (aligned to epoch), không phải theo lúc job khởi động — nên hai job khác nhau cùng size 5 phút sẽ có cùng ranh giới [..:00, ..:05), giúp kết quả có thể so khớp và tái lập. Dùng cho: đếm giao dịch/5 phút, throughput/giờ, doanh số EOD theo khung giờ.
Sliding / Hopping window (chồng lấn)
Cửa sổ sliding có hai tham số: kích thước (window size) và bước trượt (slide/advance). Nếu size = 10 phút và slide = 5 phút, ta có [00:00, 00:10), [00:05, 00:15), [00:10, 00:20)... — các cửa sổ chồng lên nhau. Một sự kiện tại phút 07 thuộc cả hai cửa sổ [00, 10) và [05, 15).
Số cửa sổ mà một sự kiện thuộc về = size / slide (làm tròn). Vì thế sliding tốn tài nguyên hơn tumbling đúng theo tỷ lệ đó — cùng một sự kiện được cộng vào nhiều state.
Dùng cho thống kê trượt (moving/rolling): "trung bình số giao dịch trong 10 phút gần nhất, cập nhật mỗi phút", giám sát tải, hay tính tỷ lệ lỗi mượt hơn. Khi slide == size, sliding suy biến thành tumbling — đó là lý do nhiều API coi tumbling là trường hợp đặc biệt.
Lưu ý thuật ngữ (dễ nhầm): Kafka Streams tách hai khái niệm.
TimeWindowsvớiadvanceBy < sizechính là hopping window (chồng lấn, ranh giới cố định) — thứ mà Flink/Beam gọi là sliding. Trong khi đóSlidingWindowscủa Kafka Streams (từ 2.7) là loại khác: cửa sổ được định nghĩa quanh từng bản ghi trong một khoảng thời gian, ranh giới do dữ liệu quyết định chứ không cố định. Khi đọc tài liệu, luôn kiểm tra engine đang dùng nghĩa nào của chữ "sliding".
Session window (theo khoảng nghỉ)
Cửa sổ session không có kích thước cố định. Nó gom các sự kiện lại thành một phiên và đóng phiên khi có một khoảng lặng (gap) đủ dài không có sự kiện nào. Tham số duy nhất là inactivity gap.
Ví dụ gap = 5 phút: nếu khách đăng nhập và thao tác lúc 00, 02, 06 rồi im lặng, sau đó thao tác lại lúc 14, 16 — engine tạo hai phiên: [00, 06] (đóng vì khoảng lặng 06→14 > 5') và [14, 16]. Độ dài mỗi phiên do dữ liệu quyết định, khác nhau tuỳ hành vi.
Session window đặc biệt hợp với hành vi người dùng: một phiên duyệt Internet Banking, một chuỗi thao tác chuyển tiền, một chuỗi giao dịch thẻ liền mạch. Về mặt cài đặt, engine tạo một cửa sổ nhỏ quanh mỗi sự kiện rồi hợp nhất (merge) các cửa sổ chồng/gần nhau lại — nên session là merging window, phức tạp về state hơn tumbling.
Global window
Cửa sổ global đặt tất cả sự kiện của một key vào một cửa sổ duy nhất, không bao giờ tự đóng theo thời gian. Bản thân nó vô dụng nếu không gắn trigger tuỳ biến để quyết định lúc phát kết quả — ví dụ "cứ mỗi 100 sự kiện thì phát một lần" (count window). Đây là cơ chế nền để xây các loại window theo số lượng hoặc theo điều kiện nghiệp vụ mà không dựa vào thời gian.
Trigger: khi nào phát ra kết quả?
Biết một sự kiện thuộc cửa sổ nào mới là nửa vấn đề. Nửa còn lại: khi nào engine "chốt sổ" và phát ra kết quả của cửa sổ đó? Đây là vai trò của trigger, và nó gắn chặt với watermark.
Nhắc lại từ Streaming 3: watermark là một mốc thời gian trôi theo luồng, mang nghĩa "engine tin rằng sẽ không còn sự kiện nào có event time nhỏ hơn mốc này nữa". Watermark cho phép engine biết khi nào một cửa sổ đã 'hoàn tất'.
Cơ chế mặc định (event-time trigger): cửa sổ phát kết quả khi watermark vượt qua điểm cuối của cửa sổ. Cửa sổ [00:00, 00:05) sẽ được tính và phát ra khi watermark ≥ 00:05. Đến lúc đó, engine tin rằng mọi giao dịch trước phút 05 đã tới đủ, nên tổng của cửa sổ là đáng tin.
Early, on-time và late firing
Trigger mặc định phát một lần khi watermark qua ranh giới — gọi là on-time firing. Nhưng thực tế cần linh hoạt hơn (mô hình của Beam/Dataflow trình bày rõ nhất):
- Early firing — phát kết quả sớm, mang tính tạm thời (speculative) trước khi watermark chốt cửa sổ, ví dụ "mỗi 1 phút phát một kết quả trung gian". Dùng khi cần dashboard cập nhật liên tục thay vì chờ đủ 5 phút mới thấy con số đầu tiên.
- On-time firing — phát khi watermark vượt điểm cuối cửa sổ. Đây là kết quả "chính thức".
- Late firing — khi có sự kiện đến trễ (event time thuộc cửa sổ đã chốt nhưng tới sau watermark). Nếu event còn nằm trong allowed lateness (thời gian gia hạn mà engine giữ state cửa sổ), engine phát lại kết quả đã cập nhật. Quá allowed lateness, cửa sổ bị huỷ state và sự kiện trễ bị bỏ (hoặc chuyển sang một luồng "side output" để xử lý riêng).
Accumulation mode
Khi một cửa sổ fire nhiều lần (early + on-time + late), kết quả các lần liên hệ nhau thế nào? Có hai chế độ tích luỹ:
- Accumulating — mỗi lần fire phát ra tổng tích luỹ đầy đủ đến thời điểm đó. Lần sau ghi đè lần trước. Hợp với sink có thể upsert theo key cửa sổ (ví dụ ghi vào bảng có khoá
window_start). - Discarding — mỗi lần fire chỉ phát ra phần thay đổi (delta) kể từ lần fire trước. Downstream phải tự cộng dồn. Hợp với sink kiểu append hoặc khi muốn giảm dữ liệu ghi ra.
Chọn sai accumulation mode là nguồn lỗi kinh điển: đếm bị nhân đôi (dùng accumulating nhưng downstream lại cộng dồn), hoặc thiếu (dùng discarding nhưng downstream lại ghi đè).
Allowed lateness vs. giữ state — đánh đổi
allowedLateness càng lớn thì càng "cứu" được nhiều sự kiện trễ, nhưng engine phải giữ state của cửa sổ lâu hơn → tốn bộ nhớ/đĩa và tăng áp lực cho checkpoint (xem Streaming 5). Đây là đánh đổi giữa độ đầy đủ (completeness) và chi phí tài nguyên, và nó là quyết định nghiệp vụ: báo cáo đối soát cuối ngày có thể chấp nhận độ trễ lớn để đầy đủ; cảnh báo gian lận realtime thì ưu tiên tốc độ, chấp nhận bỏ một ít sự kiện quá trễ.
Ví dụ 1 — Đếm giao dịch mỗi 5 phút (tumbling + event time + watermark)
Bài toán: đếm số giao dịch và tổng tiền theo từng cửa sổ 5 phút event time, cho mỗi chi nhánh. Đây là tumbling window kết hợp watermark cho phép trễ tối đa 1 phút. Ví dụ dùng Apache Flink (DataStream API):
// Flink — tumbling event-time window, cho phép trễ 1 phút.
// Minh hoạ: import và kiểu dữ liệu rút gọn cho dễ đọc.
DataStream<Txn> txns = env
.fromSource(kafkaSource, WatermarkStrategy
// watermark = maxEventTime - 1 phút (bounded out-of-orderness)
.<Txn>forBoundedOutOfOrderness(Duration.ofMinutes(1))
.withTimestampAssigner((txn, ts) -> txn.eventTimeMillis),
"txn-source");
DataStream<WindowResult> perBranch = txns
.keyBy(txn -> txn.branchId)
// tumbling 5 phút theo EVENT TIME (không phải processing time)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
// giữ state cửa sổ thêm 1 phút để nhận sự kiện đến trễ
.allowedLateness(Time.minutes(1))
// sự kiện quá trễ được đẩy ra side output thay vì bỏ im lặng
.sideOutputLateData(lateTag)
.aggregate(new CountAndSumAgg()); // đếm + cộng amount tăng dần
Vài điểm mấu chốt:
TumblingEventTimeWindows— dùng event time (thời điểm giao dịch thật) chứ không phải processing time (lúc engine nhận), nên khi hệ thống nghẽn và dữ liệu về trễ, con số vẫn gán đúng khung 5 phút.allowedLateness(1 phút)— cửa sổ[00:00,00:05)vẫn nhận sự kiện trễ tới khi watermark qua00:06, và fire lại (late firing) để cập nhật.aggregatedùngAggregateFunction(cộng tăng dần) thay vì gom hết rồi tính một lần — tiết kiệm state, chỉ giữ(count, sum)cho mỗi cửa sổ thay vì toàn bộ bản ghi.
Cùng ý tưởng bằng Kafka Streams trông như sau (hopping/tumbling qua TimeWindows):
// Kafka Streams — tumbling 5 phút; grace period = phần "allowed lateness".
KStream<String, Txn> txns = builder.stream("transactions");
txns.groupBy((k, txn) -> txn.branchId)
.windowedBy(
TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofMinutes(1)))
.aggregate(
CountSum::zero, // khởi tạo
(branch, txn, acc) -> acc.add(txn), // cộng dồn trong cửa sổ
Materialized.as("txn-5min-store")) // state store cho cửa sổ
.toStream()
.to("txn-count-5min");
ofSizeAndGrace(size, grace): grace chính là khoảng cho sự kiện trễ — tương đương allowed lateness. Nếu muốn hopping (chồng lấn), thêm .advanceBy(Duration.ofMinutes(1)).
Ví dụ 2 — Phát hiện chuỗi giao dịch trong một phiên (session window)
Bài toán chống gian lận: gom chuỗi giao dịch liền mạch của một khách/thẻ thành một phiên (session), rồi soi phiên nào có dấu hiệu bất thường — ví dụ nhiều giao dịch dồn dập trong một phiên ngắn. "Liền mạch" = không có khoảng nghỉ quá 2 phút.
// Flink — session window theo event time, gap 2 phút.
DataStream<SessionAlert> alerts = txns
.keyBy(txn -> txn.cardId)
// đóng phiên khi im lặng > 2 phút; độ dài phiên do dữ liệu quyết định
.window(EventTimeSessionWindows.withGap(Time.minutes(2)))
.process(new ProcessWindowFunction<Txn, SessionAlert, String, TimeWindow>() {
@Override
public void process(String cardId, Context ctx,
Iterable<Txn> txnsInSession,
Collector<SessionAlert> out) {
int n = 0; double sum = 0;
long start = ctx.window().getStart(), end = ctx.window().getEnd();
for (Txn t : txnsInSession) { n++; sum += t.amount; }
// heuristic minh hoạ: >= 5 giao dịch trong phiên < 3 phút => cảnh báo
if (n >= 5 && (end - start) < Duration.ofMinutes(3).toMillis()) {
out.collect(new SessionAlert(cardId, n, sum, start, end));
}
}
});
Ở đây EventTimeSessionWindows.withGap(2') tự cắt phiên theo hành vi: nếu khách quẹt thẻ 6 lần trong 90 giây rồi dừng, tất cả rơi vào một phiên, process thấy n = 6 và độ dài phiên < 3 phút → phát cảnh báo. Nếu hai lần quẹt cách nhau 3 phút, chúng thuộc hai phiên riêng biệt và không kích hoạt luật này. Session window diễn đạt "một tràng hoạt động" tự nhiên hơn nhiều so với việc ép vào cửa sổ 5 phút cố định — vì gian lận không đợi đúng ranh giới tumbling.
Chọn loại window nào?
Use case thực tế
Tại một ngân hàng như NCB (số liệu dưới đây là minh hoạ), giám sát giao dịch thẻ realtime kết hợp cả ba loại window trên cùng một luồng Kafka card-transactions:
- Tumbling 5 phút theo chi nhánh/kênh: dựng dashboard throughput và tổng tiền — ví dụ khung
10:00–10:05ghi nhận ~12.000 giao dịch, ~4,8 tỷ đồng. Watermark cho phép trễ 1 phút để các POS ở vùng mạng yếu vẫn được tính đúng khung. - Sliding 10 phút / trượt 1 phút: tính tỷ lệ giao dịch bị từ chối (decline rate) mượt theo thời gian; nếu vượt ngưỡng ~8% trong cửa sổ trượt → cảnh báo sự cố kênh thanh toán, phản ứng nhanh hơn nhiều so với chờ báo cáo giờ.
- Session window gap 2 phút theo số thẻ: phát hiện card testing — kẻ gian thử hàng loạt số thẻ với giao dịch nhỏ dồn dập; một phiên có ≥5 lần thử trong <3 phút được đẩy sang luồng cảnh báo để hệ thống risk chặn thẻ.
Sự kiện đến trễ quá allowedLateness không bị bỏ im lặng mà đưa vào side output để job đối soát cuối ngày (EOD) gom lại — nhờ vậy con số realtime nhanh mà báo cáo cuối ngày vẫn đầy đủ. Đây chính là ranh giới completeness ↔ latency được điều chỉnh theo từng mục đích.
Ghi nhớ
- Aggregate trên luồng vô hạn không hội tụ; window cắt luồng thành các đoạn hữu hạn để tổng hợp có nghĩa.
- Tumbling: cố định, không chồng, mỗi sự kiện thuộc 1 cửa sổ, căn theo epoch. Sliding/Hopping: cố định, chồng lấn, một sự kiện thuộc
size/slidecửa sổ — tốn tài nguyên hơn. - Session: kích thước thay đổi, cắt theo khoảng lặng (gap), là merging window — hợp với hành vi người dùng. Global: một cửa sổ, vô dụng nếu thiếu trigger tuỳ biến.
- Cẩn thận thuật ngữ: "sliding" của Kafka Streams khác hopping window và khác "sliding" của Flink/Beam.
- Trigger quyết định khi nào phát kết quả; mặc định fire khi watermark vượt cuối cửa sổ (on-time). Có thêm early firing (tạm thời) và late firing (khi có sự kiện trễ trong allowed lateness).
- Accumulation mode: accumulating (phát tổng đầy đủ, upsert) vs discarding (phát delta, cộng dồn) — chọn sai gây đếm nhân đôi hoặc thiếu.
- Window luôn dùng chung với event time + watermark;
allowedLatenesslớn = đầy đủ hơn nhưng giữ state lâu hơn, tốn checkpoint. - Cửa sổ là một dạng state có vòng đời — đọc Streaming 5 để hiểu cách nó được lưu và khôi phục; xem thêm ứng dụng ở Streaming 9 — realtime ngân hàng.
Nguồn tham khảo
- Streaming Systems — Tyler Akidau, Slava Chernyak, Reuven Lax (O'Reilly): trình bày mô hình what/where/when/how, windowing, trigger và accumulation chuẩn mực nhất.
- The Dataflow Model — Akidau et al., VLDB 2015: bài báo nền tảng về event time, window và trigger (cơ sở của Apache Beam).
- Apache Flink — Windows (docs chính thức): https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/operators/windows/
- Apache Flink — Generating Watermarks / Event Time: https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/event-time/generating_watermarks/
- Kafka Streams — Windowing (Confluent/Apache docs): https://kafka.apache.org/documentation/streams/developer-guide/dsl-api.html#windowing
- Apache Beam — Windowing & Triggers (programming guide): https://beam.apache.org/documentation/programming-guide/#windowing
- Designing Data-Intensive Applications — Martin Kleppmann (O'Reilly), chương "Stream Processing".
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ẻ!