Streaming 3 — Event Time, Processing Time & Watermark

22 thg 7, 2026 2 lượt xem
#data-engineering
#streaming
#flink
#watermark
#event-time
#beam

Mô hình tinh thần: "khi nào sự kiện xảy ra" khác "khi nào ta xử lý nó"

Trong xử lý batch cổ điển, câu hỏi thời gian thường bị che giấu: bạn đọc một partition ngày 2026-07-20, mặc nhiên coi mọi dòng trong đó thuộc về ngày đó. Nhưng streaming buộc bạn đối mặt với sự thật khó chịu: thời điểm một sự kiện xảy ra ngoài đời thực, và thời điểm hệ thống của bạn nhìn thấy nó, là hai con số khác nhau — và thường lệch nhau rất xa.

Toàn bộ độ khó của xử lý luồng "đúng đắn" nằm ở khoảng lệch này. Nếu bạn tính tổng giao dịch theo cửa sổ 5 phút mà lại dùng nhầm loại thời gian, kết quả sẽ sai một cách âm thầm — không lỗi, không cảnh báo, chỉ là con số lệch. Bài này (thuộc series Streaming) làm rõ ba loại thời gian, giải thích cơ chế watermark để "biết mình đã thấy đủ dữ liệu chưa", và cách xử lý dữ liệu đến trễ. Đây là nền tảng bắt buộc trước khi bàn về windowingstate / exactly-once.

Ba loại thời gian (three domains of time)

Mô hình chuẩn — được hệ thống hoá trong Google Dataflow model và áp dụng trong Apache Beam / Apache Flink — phân biệt:

Loại thời gianĐịnh nghĩaAi gán?Đặc tính
Event timeThời điểm sự kiện thực sự xảy ra tại nguồnThiết bị/ứng dụng sinh sự kiện nhúng vào payloadBất biến; là "sự thật" nhưng đến sai thứ tự
Ingestion timeThời điểm sự kiện đi vào hệ thống streamingBroker/source (vd Kafka gán khi ghi vào topic)Đơn điệu tăng dần; xấp xỉ event time nếu đường truyền nhanh
Processing timeThời điểm một toán tử đang xử lý sự kiệnĐồng hồ tường (wall clock) của máy đang chạyĐơn giản, latency thấp; không tất định, phụ thuộc tải & tốc độ

Điểm mấu chốt: processing time luôn ≥ event time, và khoảng chênh (gọi là skew, hay event-time lag) biến động liên tục. Ingestion time nằm giữa — nó là một dạng "event time do hệ thống gán", ổn định hơn nhưng không phản ánh đúng lúc sự kiện thực sự phát sinh (vd một giao dịch offline chỉ vào Kafka khi điện thoại có mạng trở lại).

Ví dụ ngân hàng: một khách quẹt thẻ lúc 10:00:02 (event time). Điện thoại POS mất sóng, gói tin chỉ về tới Kafka lúc 10:00:20 (ingestion time), và job Flink chống gian lận xử lý nó lúc 10:00:22 (processing time). Nếu bạn gom "giao dịch trong phút 10:00" theo processing time, giao dịch này rơi nhầm sang cửa sổ khác.

Vì sao out-of-order & skew là điều tất yếu

Đây không phải lỗi cấu hình mà là bản chất của hệ phân tán:

  • Độ trễ mạng & retry: gói tin đi qua nhiều chặng, có thể mất và gửi lại → đến muộn, đến lệch thứ tự.
  • Client offline (đặc thù mobile): app ghi sự kiện vào bộ đệm cục bộ khi mất mạng, đẩy hàng loạt (batch) khi có mạng lại — có thể trễ vài giây tới vài giờ.
  • Xử lý song song & phân mảnh: nhiều partition Kafka, nhiều task đọc song song với tốc độ khác nhau; thứ tự toàn cục không được bảo toàn giữa các partition.
  • Backpressure & rebalance: khi hệ thống nghẽn hoặc consumer group cân bằng lại, một luồng có thể "khựng" rồi ồ ạt bù.

Đọc sơ đồ: e3 xảy ra trước e2 (theo event time), nhưng vì client offline nên nó được xử lý sau e2. Nếu chỉ nhìn thứ tự đến (processing order), ta thấy e1 → e2 → e3; nhưng theo event time đúng phải là e1 → e3 → e2. Đó chính là out-of-order.

Watermark: "tôi tin rằng đã thấy hết dữ liệu tới thời điểm T"

Nếu dữ liệu đến trễ và sai thứ tự, làm sao biết khi nào một cửa sổ event-time đã đủ dữ liệu để chốt kết quả? Ta không thể chờ mãi (chờ vô hạn = không bao giờ có kết quả). Câu trả lời là watermark.

Watermark là một mốc trong miền event time, kèm ngữ nghĩa: "hệ thống ước lượng rằng sẽ không còn (hoặc còn rất ít) sự kiện có event time nhỏ hơn W nữa". Khi watermark vượt qua thời điểm kết thúc một cửa sổ, cửa sổ đó được coi là "hoàn thành" và có thể phát kết quả (fire).

Có hai loại theo mô hình Dataflow:

  • Perfect watermark: biết chắc chắn mọi dữ liệu ≤ W đã tới (chỉ khả thi khi nguồn có thông tin đầy đủ, hiếm gặp).
  • Heuristic watermark: ước lượng dựa trên hiểu biết về nguồn (độ trễ tối đa quan sát được). Đây là loại phổ biến trong thực tế — và vì là ước lượng nên có thể sai: sự kiện đến sau watermark gọi là dữ liệu trễ (late data).

Watermark phải đơn điệu không giảm (monotonically non-decreasing): một khi đã tuyên bố "thấy đủ tới T", ta không được rút lại.

Chiến lược phổ biến nhất là bounded out-of-orderness: giả định độ trễ tối đa là một hằng số B. Watermark tại thời điểm bất kỳ = (event time lớn nhất đã thấy) − B. Trong Flink, forBoundedOutOfOrderness cài đặt chính xác công thức này (kỹ thuật là maxTs − B − 1ms để biên cửa sổ đóng đúng chuẩn nửa-mở).

// Flink: gán timestamp từ payload và sinh watermark bounded out-of-orderness
WatermarkStrategy<Txn> strategy = WatermarkStrategy
    .<Txn>forBoundedOutOfOrderness(Duration.ofSeconds(5))   // B = 5 giây
    .withTimestampAssigner((txn, ts) -> txn.getEventTimeMillis())
    .withIdleness(Duration.ofMinutes(1));   // đánh dấu partition "im lặng" để không kẹt watermark

DataStream<Txn> stream = source.assignTimestampsAndWatermarks(strategy);

Vài điểm cần đúng:

  • Watermark theo từng phân vùng, rồi lấy min: mỗi partition/subtask sinh watermark riêng; watermark của một toán tử downstream = giá trị nhỏ nhất trong các luồng đầu vào của nó. Vì thế một partition chậm/kẹt sẽ kéo lùi watermark toàn hệ — lý do cần withIdleness để bỏ qua partition không có dữ liệu.
  • Periodic vs punctuated: Flink mặc định phát watermark theo chu kỳ (mặc định 200ms, cấu hình pipeline.auto-watermark-interval); ngoài ra có thể phát theo từng sự kiện đặc biệt (punctuated) nếu payload có "dấu mốc".
  • Watermark lan truyền (propagation): trong Beam/Dataflow, mỗi PTransform có input watermark và output watermark; watermark chảy qua toàn pipeline, đảm bảo mọi tầng đều biết "đã thấy đủ tới đâu".

Allowed lateness & xử lý dữ liệu trễ

Vì heuristic watermark có thể "đi nhanh hơn thực tế", luôn có sự kiện đến khi watermark đã vượt qua nó — đó là late data. Có ba cách xử lý, đây là quyết định thiết kế, không có lựa chọn đúng tuyệt đối:

Chiến lượcCơ chếKhi nào dùng
Drop (mặc định)Bỏ luôn sự kiện trễ hơn watermark (Flink: quá allowedLateness; Beam: quá allowed_lateness)Chấp nhận mất mát nhỏ, ưu tiên đơn giản & latency
Side outputĐịnh tuyến late data sang một luồng riêng để xử lý sau / auditKhông được mất dữ liệu (đối soát, tuân thủ)
Update / re-fireGiữ state cửa sổ thêm một khoảng, mỗi khi có late data thì tính lại & phát bản cập nhậtCần kết quả đúng dần (eventual correctness); sink hỗ trợ upsert/retraction

Allowed lateness là khoảng thời gian (theo event time) mà cửa sổ vẫn giữ state sau khi watermark đã đi qua, để còn kịp thu nạp sự kiện trễ. Quá mốc windowEnd + allowedLateness, state bị giải phóng (garbage-collect) và mọi sự kiện trễ hơn nữa bị drop hoặc chuyển side output.

Trong Flink, cả ba hành vi được diễn đạt rõ ràng qua API cửa sổ:

OutputTag<Txn> lateTag = new OutputTag<Txn>("late-txn") {};

SingleOutputStreamOperator<FraudScore> result = stream
    .keyBy(Txn::getCardId)
    .window(TumblingEventTimeWindows.of(Time.minutes(5)))
    .allowedLateness(Time.seconds(30))       // giữ state thêm 30s sau khi watermark qua
    .sideOutputLateData(lateTag)             // quá 30s → đẩy sang side output thay vì drop
    .aggregate(new FraudAggregator());

DataStream<Txn> lateStream = result.getSideOutput(lateTag);   // đối soát/audit sau

Trong Beam/Dataflow, cơ chế trigger tách rời "khi nào phát kết quả" khỏi "cửa sổ chứa gì", cho phép phát early pane (sớm, chưa đủ dữ liệu), on-time pane (khi watermark qua) và late pane (khi có late data), kết hợp với AccumulationMode để chọn cộng dồn hay thay thế:

import apache_beam as beam
from apache_beam.transforms.window import FixedWindows
from apache_beam.transforms.trigger import (
    AfterWatermark, AfterProcessingTime, AfterCount,
    Repeatedly, AccumulationMode)

(events
 | beam.WindowInto(
       FixedWindows(5 * 60),                              # cửa sổ tumbling 5 phút theo event time
       trigger=AfterWatermark(
           early=Repeatedly(AfterProcessingTime(30)),     # kết quả sớm mỗi 30s (latency thấp)
           late=Repeatedly(AfterCount(1))),               # cập nhật ngay khi có 1 late data
       allowed_lateness=3600,                             # chấp nhận trễ tối đa 1 giờ (event time)
       accumulation_mode=AccumulationMode.ACCUMULATING)   # mỗi lần fire là kết quả tích luỹ mới nhất
 | beam.CombinePerKey(sum))

Trade-off cốt lõi: latency vs completeness

Mọi quyết định về watermark và lateness đều xoay quanh một đánh đổi duy nhất:

  • Watermark "hào phóng" (B lớn) / allowedLateness dài → chờ lâu hơn, thu được nhiều dữ liệu hơn (completeness cao), nhưng latency cao và tốn state (giữ cửa sổ lâu).
  • Watermark "gắt" (B nhỏ) / allowedLateness ngắn → phát kết quả sớm (latency thấp), nhưng completeness thấp — nhiều sự kiện bị coi là trễ và drop.

Không có cấu hình "đúng" phổ quát: bạn chọn điểm trên đường cong này theo yêu cầu nghiệp vụ. Mô hình trigger của Dataflow là cách "ăn cả hai": phát early pane để có phản hồi nhanh, rồi phát on-time và late pane để dần hội tụ về kết quả đúng — với điều kiện sink chịu được cập nhật/ghi đè (idempotent upsert), chủ đề nối sang state & exactly-once.

Use case thực tế

Chống gian lận thẻ realtime tại NCB (số liệu minh hoạ). Yêu cầu: cảnh báo khi một thẻ có ≥5 giao dịch trong 5 phút ở các quốc gia khác nhau. Đầu vào là luồng giao dịch từ nhiều kênh (POS, ATM, e-commerce), trong đó giao dịch từ app mobile có thể trễ do offline.

  • Nếu dùng processing time: một loạt giao dịch cùng phát sinh lúc 10:00 nhưng về hệ thống rải rác 10:00–10:03 sẽ bị gom sai cửa sổ → bỏ sót cụm gian lận hoặc báo động giả.
  • Chọn event time + watermark bounded out-of-orderness B = 3 phút (đo từ phân phối độ trễ thực tế: ~99% giao dịch về trong 3 phút). Cửa sổ tumbling 5 phút, allowedLateness = 2 phút, late data đẩy sang side output cho đội đối soát cuối ngày (EOD).
  • Kết quả minh hoạ: latency cảnh báo ~3–4 phút (chấp nhận được cho fraud), completeness ~99.7%; phần ~0.3% giao dịch quá trễ được xử lý bù trong batch EOD, không mất mát cho mục đích đối soát. So sánh: nếu hạ B xuống 30 giây, latency còn ~1 phút nhưng ~4% giao dịch bị coi là trễ → nguy cơ sót cảnh báo tăng mạnh.

Ghi nhớ

  • Ba miền thời gian: event time (lúc sự kiện xảy ra, bất biến, đến sai thứ tự), ingestion time (lúc vào hệ thống, đơn điệu), processing time (lúc xử lý, wall clock, không tất định). Processing time ≥ event time; khoảng chênh là skew.
  • Out-of-order & skew là tất yếu của hệ phân tán: mạng/retry, client mobile offline, xử lý song song nhiều partition, backpressure.
  • Watermark là mốc event time với ngữ nghĩa "đã thấy hết dữ liệu tới T"; heuristic nên có thể sai → sinh ra late data. Watermark phải đơn điệu không giảm.
  • Chiến lược phổ biến: bounded out-of-orderness, watermark = max_event_time − B. Watermark tính theo từng partition rồi lấy min; partition kẹt sẽ kéo lùi cả hệ → cần idleness.
  • Allowed lateness giữ state cửa sổ sau khi watermark qua để thu nạp late data; ba cách xử lý: drop / side output / update (re-fire).
  • Trigger (Beam) tách "khi nào phát" khỏi "cửa sổ chứa gì": early / on-time / late pane + accumulation mode để hội tụ dần về kết quả đúng.
  • Đánh đổi trung tâm: latency vs completeness — watermark/lateness càng rộng thì càng đầy đủ nhưng càng chậm và tốn state.
  • Muốn "update kết quả" cần sink idempotent/upsert; đọc tiếp windowingstate & exactly-once.

Nguồn tham khảo

  • Tyler Akidau et al., "The Dataflow Model: A Practical Approach to Balancing Correctness, Latency, and Cost in Massive-Scale, Unbounded, Out-of-Order Data Processing", VLDB 2015.
  • Tyler Akidau, Slava Chernyak, Reuven Lax, "Streaming Systems" (O'Reilly) — chương "The What, Where, When, and How".
  • Tyler Akidau, "Streaming 101" & "Streaming 102" (O'Reilly Radar) — nền tảng event time / watermark / trigger.
  • Apache Flink documentation — "Generating Watermarks", "Event Time and Watermarks", "Windows / Allowed Lateness & Side Output": https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/event-time/generating_watermarks/
  • Apache Beam documentation — "Streaming pipelines: Windowing, Triggers, Watermarks": https://beam.apache.org/documentation/programming-guide/#triggers
  • Martin Kleppmann, "Designing Data-Intensive Applications" (O'Reilly) — chương 11, Stream Processing (reasoning about time).
  • Joe Reis & Matt Housley, "Fundamentals of Data Engineering" (O'Reilly).

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