← Về danh sách
AI & Dữ liệu
#스트림처리#윈도잉#워터마크#이벤트시간#ExactlyOnce
Cập nhật lần cuối · 2026-10-07

Cửa sổ hóa và ngữ nghĩa thời gian trong xử lý luồng (Event Time, Watermark, Exactly-Once)

1. Tổng quan

A. Định nghĩa

Xử lý luồng (Stream Processing) là một mô hình xử lý dữ liệu, biến đổi, tổng hợp và phân tích một cách liên tục ngay khi dữ liệu vô hạn (unbounded) — dữ liệu không có điểm bắt đầu và kết thúc cố định — vừa đến, tạo ra kết quả với độ trễ thấp.

Cửa sổ hóa (Windowing) là kỹ thuật thiết lập ranh giới nhằm cắt luồng vô hạn thành các đơn vị xử lý hữu hạn bằng cách nhóm các sự kiện theo tiêu chí như thời gian, số lượng hay phiên, còn ngữ nghĩa thời gian (Time Semantics) là mô hình bảo đảm độ chính xác của phép tổng hợp bằng cách phân biệt "thời điểm sự kiện thực sự xảy ra (event time)" với "thời điểm engine xử lý nó (processing time)".

Khó khăn bản chất của xử lý luồng nằm ở câu hỏi "đối với dữ liệu không bao giờ kết thúc, khi nào, cái gì và chính xác đến mức nào thì chốt một phép tính?". Xử lý theo lô (batch) có ranh giới hiển nhiên vì nó tính toán sau khi toàn bộ đầu vào đã tụ hội, nhưng luồng thì chảy không ngừng và, do độ trễ mạng, đến không theo thứ tự (out-of-order), nên engine phải tự định nghĩa thời điểm chốt phép tổng hợp. Cửa sổ hóa quyết định "coi cái gì là một nhóm (where in event time)", còn watermark quyết định "khi nào đóng nhóm đó và phát ra kết quả (when in processing time)". Dưới góc nhìn của Kỹ sư chuyên nghiệp, chủ đề này không phải là việc dùng API đơn thuần mà là bài toán quản trị đường ống dữ liệu thời gian thực nhằm thiết kế sự đánh đổi giữa độ trễ, độ chính xác và tính đầy đủ (completeness).

Về bản chất, xử lý luồng đối mặt với nghịch lý rằng "chờ đến khi dữ liệu tụ hội đầy đủ thì quá muộn, mà không chờ thì không đầy đủ". Trong khi batch chấp nhận độ trễ vì tính đầy đủ và thời gian thực thuần túy từ bỏ tính đầy đủ để triệt tiêu độ trễ, xử lý luồng hiện đại tham số hóa điểm trung gian đó qua ba núm vặn — cửa sổ, watermark và trigger — để cùng một mã nguồn có thể được điều chỉnh theo yêu cầu nghiệp vụ ở bất kỳ đâu giữa "kết quả nhanh nhưng xấp xỉ" và "kết quả chậm nhưng chính xác". Đó chính là đặc trưng định danh của nó.

B. Bối cảnh ra đời và sự cần thiết

Thứ nhất, chân trời thời gian của quyết định kinh doanh đã dịch từ "ngày hôm sau" sang "chính khoảnh khắc này". Hệ thống phát hiện gian lận (FDS), gợi ý thời gian thực, phát hiện bất thường thiết bị trong Internet vạn vật (IoT), và việc phản ánh tồn kho, giá thị trường biến động nhanh đều mất giá trị nếu xử lý bằng batch ban đêm. Vì quyết định phải được đưa ra trong độ trễ vài giây đến vài phút, một cấu trúc xử lý dữ liệu ngay trong lúc nó chảy — thay vì chờ nó tích tụ — đã trở nên cần thiết.

Thứ hai, trong môi trường thu thập phân tán, dữ liệu vốn dĩ đến muộn và xáo trộn. Một thiết bị di động ngoại tuyến rồi phục hồi sẽ gửi muộn các sự kiện từ vài phút trước, hoặc chênh lệch tốc độ xử lý theo từng phân vùng làm đảo thứ tự. Nếu tổng hợp theo thời gian xử lý, "đơn hàng phát sinh lúc 10h nhưng đến lúc 10h05" rơi vào cửa sổ sai và làm méo kết quả. Để khắc phục điều này cần đến xử lý theo event time dựa trên thời điểm phát sinh mà chính sự kiện mang theo, cùng với khái niệm watermark cho phép độ trễ nhưng vạch ra một giới hạn cho nó.

Thứ ba, là thời gian thực không có nghĩa được phép từ bỏ độ chính xác. Nếu sự kiện bị mất hoặc bị đếm trùng khi một nút chết hoặc khởi động lại do sự cố, hậu quả là chí mạng trong các lĩnh vực như tính cước, quyết toán và báo cáo tuân thủ. Do đó, bảo đảm xử lý chính xác một lần (exactly-once) — theo đó mỗi sự kiện được phản ánh vào kết quả đúng một lần ngay cả khi có sự cố — đã trở thành yêu cầu cốt lõi của các engine luồng.

Thứ tư, yêu cầu về mặt vận hành và chi phí cũng tăng lên. Một đường ống chạy không ngừng suốt ngày đêm phải giữ cho trạng thái của nó không phình to vô hạn, và khi logic thay đổi, nó phải có thể phát lại dữ liệu quá khứ để tái tạo kết quả (reprocessing). Hơn nữa, nếu không có cơ chế điều khiển áp lực ngược (backpressure) truyền tải lên thượng nguồn để tự bảo vệ khi có đột biến tức thời (traffic spike), toàn hệ thống sẽ sụp đổ vì cạn kiệt bộ nhớ. Như vậy, xử lý luồng đã tiến hóa thành một công nghệ đòi hỏi không chỉ độ chính xác mà cả khả năng vận hành bền vững.

C. Phổ các mức bảo đảm xử lý

Độ tin cậy của xử lý luồng được quy định bởi "mỗi sự kiện được phản ánh vào kết quả bao nhiêu lần khi có sự cố". Mức bảo đảm tối thiểu at-most-once không truyền lại nên cho phép mất mát nhưng không có trùng lặp. At-least-once ngăn mất mát bằng tái xử lý khi sự cố, nhưng cùng một sự kiện có thể được phản ánh hai lần, làm tổng và số đếm bị tính vượt. Exactly-once bảo đảm trạng thái (state) và đầu ra được phản ánh đúng một lần, nhất quán xuyên qua một sự cố. Tuy nhiên điều này không có nghĩa bản thân việc truyền qua mạng chỉ xảy ra một lần về mặt vật lý; đúng hơn, checkpoint và giao dịch làm cho "hiệu ứng (effect)" xảy ra một lần, nên nó cũng được gọi là effectively-once. Sự phân biệt này được hiện thực hóa qua các cơ chế checkpoint và giao dịch được đề cập ở mục chuyên sâu bên dưới.

Việc chọn mức bảo đảm không miễn phí. Mức bảo đảm càng cao thì chi phí phụ trội của snapshot trạng thái, điều phối giao dịch và commit trì hoãn càng lớn, và độ trễ càng cao. Do đó, thay vì thiết kế mọi đường ống đồng loạt là exactly-once, thông lệ chuẩn mực là áp dụng có chọn lọc: chỉ đặt exactly-once trên các tuyến mà độ chính xác gắn trực tiếp với trách nhiệm tài chính hay pháp lý — tính cước, quyết toán, báo cáo tuân thủ — còn các tuyến dung thứ một lượng nhỏ trùng lặp hoặc mất mát, như chỉ số bảng điều khiển hay tín hiệu gợi ý, thì đặt ở at-least-once để tối ưu chi phí.

2. Kiến trúc tổng thể và ngữ nghĩa thời gian

A. Kiến trúc đường ống xử lý luồng

Đường ống thời gian thực thường được xây dựng quanh một message broker dựa trên log bất biến, cùng với các toán tử xử lý có trạng thái, một kho checkpoint, và các sink lũy đẳng (idempotent) hoặc giao dịch.

flowchart LR
  SRC["Nguồn sự kiện(ứng dụng, IoT, log)"] --> BROKER["Message broker(log bất biến, phân vùng)"]
  BROKER --> OP["Toán tử luồng(cửa sổ, tổng hợp)"]
  OP --> ST["Kho trạng thái(Keyed State)"]
  OP --> CP["Checkpoint(snapshot)"]
  CP --> DFS["Kho bền vững(hệ thống tệp phân tán)"]
  OP --> SINK["Sink(lũy đẳng, giao dịch)"]
  SINK --> SERVE["Phục vụ(DB, bảng điều khiển, cảnh báo)"]

Message broker (ví dụ Apache Kafka) cung cấp một log bất biến mà thứ tự được bảo đảm theo từng phân vùng, trở thành nền tảng cho tái xử lý bằng cách tua lại offset để phát lại (replay) khi sự cố. Đặt thời gian lưu giữ (retention) là, chẳng hạn, 7 ngày cho phép tua vị trí tiêu thụ về bất kỳ điểm nào trong khoảng đó và áp dụng lại logic, còn số phân vùng quyết định đơn vị song song và do đó là cơ sở cho việc mở rộng thông lượng. Toán tử luồng duy trì trạng thái theo từng khóa và thực hiện tổng hợp theo cửa sổ, định kỳ chụp snapshot trạng thái làm checkpoint và lưu vào hệ thống tệp phân tán. Khi sự cố, nó khôi phục trạng thái từ checkpoint gần nhất và tiếp tục tiêu thụ từ offset của điểm đó, nên tái xử lý at-least-once là mặc định. Kết hợp điều này với một sink lũy đẳng hoặc giao dịch loại bỏ luôn cả phản ánh trùng lặp, hoàn thiện exactly-once. Trong cấu trúc này, chu kỳ checkpoint chi phối mức độ dày của các điểm khôi phục, còn hiệu năng của state backend chi phối thông lượng tổng hợp, nên tinh chỉnh cả hai cùng nhau là điểm khởi đầu của thiết kế vận hành.

B. Event time, processing time và watermark

Điểm khởi đầu của ngữ nghĩa thời gian là phân biệt ba thời điểm. Event time là thời điểm sự kiện thực sự xảy ra (ví dụ thời điểm chấp thuận thanh toán), được nhúng trong dữ liệu dưới dạng dấu thời gian. Thời gian thu nạp (ingestion time) là thời điểm nó vào broker hoặc engine, còn processing time là thời điểm toán tử thực sự xử lý nó. Tính tái lập và độ chính xác của kết quả đến từ xử lý theo event time, bởi khi phát lại cùng một đầu vào (tái xử lý), processing time khác nhau mỗi lần, trong khi tổng hợp theo event time luôn cho cùng một kết quả.

Vấn đề là đóng một cửa sổ theo event time đòi hỏi phán đoán rằng "sẽ không còn sự kiện thuộc cửa sổ đó đến nữa". Thiết bị ước lượng điều này là watermark. Watermark W(t) là chỉ báo tiến độ của engine khẳng định rằng "dữ liệu trước event time t (về đại thể) đã đến đủ", và khoảnh khắc watermark này vượt qua thời điểm kết thúc cửa sổ, phép tổng hợp của cửa sổ đó được chốt và kích hoạt (trigger). Trong thực tế, một watermark vô trật tự có chặn (bounded out-of-orderness) — event time lớn nhất quan sát được trừ đi độ trễ cho phép — thường được dùng. Ví dụ với độ trễ cho phép là 5 giây, W = max_event_time - 5s.

sequenceDiagram
  participant E as Luồng sự kiện
  participant W as Bộ sinh watermark
  participant WIN as Toán tử cửa sổ
  participant OUT as Sink kết quả
  E->>W: Sự kiện đến(gồm dấu thời gian event time)
  W->>WIN: Lan truyền watermark("hoàn tất đến trước t")
  WIN->>WIN: So sánh ranh giới cửa sổ với watermark
  alt watermark > kết thúc cửa sổ
    WIN->>OUT: Chốt và kích hoạt tổng hợp cửa sổ
  else dữ liệu muộn đến
    WIN->>WIN: Cập nhật nếu trong allowed lateness
    WIN->>OUT: Tách ra qua side output
  end

Để tham khảo, thời gian thu nạp là sự thỏa hiệp giữa event time và processing time; khi nguồn không có dấu thời gian đáng tin, nó là giải pháp thay thế thực dụng dùng thời điểm vào broker làm cơ sở để có được một mức ổn định thứ tự nào đó. Tuy nhiên nó không bảo đảm tính tái lập đầy đủ, vì kết quả có thể khác khi tái xử lý.

Watermark về bản chất là núm vặn hòa giải sự đánh đổi giữa tính đầy đủ và độ trễ. Đặt độ trễ cho phép lớn khiến kết quả chính xác do bao gồm cả dữ liệu đến muộn, nhưng việc chốt cửa sổ bị trì hoãn, làm tăng độ trễ. Đặt nó nhỏ thì nhanh nhưng bỏ sót dữ liệu muộn, khiến kết quả không đầy đủ. Vì vậy nhiều engine tiếp tục cập nhật kết quả bằng dữ liệu muộn trong một khoảng allowed lateness ngay cả sau khi cửa sổ đã chốt, và tách dữ liệu muộn hơn thế qua một side output vào một tuyến hiệu chỉnh riêng.

C. Các loại cửa sổ

Các cửa sổ được phân chia theo tiêu chí dùng để nhóm luồng vô hạn. Việc chọn loại phụ thuộc vào bản chất câu hỏi (là tổng hợp định kỳ hay phân tích các khoảng hoạt động); bảng dưới đây chỉ là tóm tắt bổ trợ, còn lý do chọn mỗi loại được giải thích bằng văn xuôi.

Loại cửa sổ Định nghĩa Đặc điểm Dùng tiêu biểu
Tumbling Kích thước cố định, không chồng lấn Ranh giới không chồng nên mỗi sự kiện rơi đúng 1 cửa sổ Tổng hợp doanh thu 1 phút, báo cáo theo giờ
Sliding Kích thước cố định, dịch theo khoảng cố định Cửa sổ chồng nên một sự kiện thuộc nhiều cửa sổ Trung bình trượt 5 phút (mỗi phút)
Session Dựa trên khoảng trống (gap) giữa các hoạt động Kích thước biến thiên, tự tách các khoảng hoạt động Phiên người dùng, khoảng vận hành thiết bị
Global Không ranh giới, trigger tùy chỉnh Người dùng định nghĩa trigger Tổng hợp dựa trên số lượng

Cửa sổ tumbling cắt thành các độ dài cố định như 1 phút hay 1 giờ mà không chồng lấn, và mỗi sự kiện thuộc đúng một cửa sổ. Nó phù hợp với các báo cáo định kỳ rõ ràng như "số giao dịch mỗi phút". Ranh giới của nó đơn giản nên chi phí quản lý trạng thái là thấp nhất, và khi một cửa sổ đóng, trạng thái của nó có thể bị loại bỏ ngay, khiến bộ nhớ dễ dự đoán. Mặt khác, vì ranh giới cố định, ta phải tính đến trong thiết kế việc một hiện tượng ngưỡng bị chia cắt qua hai cửa sổ, mỗi cửa sổ rơi dưới ngưỡng và có thể bỏ sót phát hiện.

Cửa sổ sliding tách kích thước (ví dụ 5 phút) khỏi khoảng dịch (ví dụ 1 phút) để các cửa sổ chồng lấn, và được dùng cho giám sát xu hướng liên tục như "trung bình trượt 5 phút được làm mới mỗi phút". Vì một sự kiện được bao gồm trùng lặp trong (kích thước ÷ khoảng) cửa sổ, kích thước trạng thái và lượng tính toán tăng tương ứng. Trong ví dụ trên, một sự kiện thuộc 5 cửa sổ, nên có khoảng gấp 5 lần gánh nặng trạng thái và tính toán so với tumbling đơn giản. Sự thỏa hiệp giữa độ nhạy xu hướng (khoảng ngắn hơn) và chi phí tài nguyên là điểm thiết kế cốt lõi.

Cửa sổ session chia các khoảng không theo ranh giới cố định mà theo khoảng trống (gap) giữa các hoạt động. Ví dụ với ngưỡng khoảng trống 30 phút, nếu các cú nhấp liên tục của người dùng gián đoạn 30 phút trở lên, phiên kết thúc và một phiên mới bắt đầu. Nó tự nhiên cho phân tích hành vi người dùng hay trích xuất các phiên vận hành thiết bị. Vì độ dài cửa sổ được xác định động theo dữ liệu, khi một sự kiện muộn lấp khoảng trống giữa hai phiên, nảy sinh phức tạp phải hợp nhất (merge) hai phiên đã kích hoạt thành một. Vì lý do này, cửa sổ session là loại mà thiết kế trạng thái và trigger khắt khe nhất.

Cửa sổ global kích hoạt qua một trigger do người dùng định nghĩa như "mỗi 100 bản ghi", không có ranh giới thời gian. Nó đáp ứng linh hoạt với xử lý dựa trên điều kiện số lượng hay ngưỡng khó biểu đạt bằng tổng hợp theo thời gian, nhưng vì người dùng hoàn toàn chịu trách nhiệm về trigger và dọn dẹp trạng thái, lạm dụng gây rò rỉ trạng thái và tăng bộ nhớ.

D. Trigger và chế độ tích lũy

Ngay cả với cùng một cửa sổ, "khi nào và phát kết quả bao nhiêu lần" được thiết kế riêng như trigger (Trigger), và "kết hợp với kết quả trước ra sao khi kích hoạt lại" như chế độ tích lũy (accumulation mode). Trong môi trường mà dữ liệu muộn quan trọng, phát một kết quả đầu tiên nhanh (speculative) khi watermark đến, rồi kích hoạt lại cửa sổ để hiệu chỉnh kết quả khi dữ liệu muộn đến sau, là hữu ích. Ở đây chế độ tích lũy (accumulating) phát ra giá trị tính lại toàn phần mỗi lần, còn chế độ tích lũy và rút lại (accumulating & retracting) phát ra một bản ghi hiệu chỉnh hủy giá trị trước cùng với giá trị mới, để hạ nguồn không tổng hợp trùng. Mẫu "kết quả sớm + hiệu chỉnh hậu kỳ" này là đóng góp cốt lõi của mô hình Dataflow và là chìa khóa để đạt độ trễ thấp mà không hy sinh tính đầy đủ.

3. So sánh và các trường hợp áp dụng

A. So sánh các engine tiêu biểu và lý do khác biệt

Phân loại Apache Flink Kafka Streams Spark Structured Streaming
Mô hình xử lý Streaming thực sự theo từng bản ghi Theo từng bản ghi (thư viện) Micro-batch (mặc định), liên tục (thử nghiệm)
Trạng thái và checkpoint Asynchronous Barrier Snapshotting (ABS) Topic changelog + RocksDB Checkpoint + WAL
exactly-once Bảo đảm qua trạng thái + sink giao dịch exactly_once_v2 (trong Kafka) Dựa trên sink lũy đẳng hoặc giao dịch
Đặc tính độ trễ Độ trễ thấp cỡ mili-giây Cỡ mili-giây (nhúng trong ứng dụng) Bằng chu kỳ batch (hàng trăm ms đến giây)

Khác biệt giữa ba engine bắt nguồn từ quan điểm nền tảng "streaming là gì". Flink coi luồng là công dân hạng nhất và xử lý liên tục theo từng bản ghi, nên mạnh về độ trễ thấp và điều khiển event time cùng watermark tinh vi. Spark Structured Streaming theo truyền thống coi luồng là một chuỗi các lô nhỏ — mô hình micro-batch — nên có lợi cho thông lượng và tích hợp với hệ sinh thái batch, nhưng phải trả giá bằng độ trễ cỡ chu kỳ batch (chế độ xử lý liên tục hạ độ trễ nhưng có ràng buộc về mức bảo đảm). Kafka Streams là một thư viện nhúng trong ứng dụng mà không cần cụm riêng, nên vận hành đơn giản; nó sao chép trạng thái sang một topic changelog để khôi phục và cung cấp exactly-once cho I/O nội bộ Kafka. Nói ngắn gọn, Flink là lựa chọn hợp lý khi độ trễ thấp và điều khiển thời gian tinh vi là trọng tâm, Kafka Streams cho xử lý nhúng nhẹ xoay quanh Kafka, và Spark khi tích hợp với hệ sinh thái batch và ML là quan trọng.

B. Các trường hợp áp dụng cụ thể

Hãy xét phát hiện gian lận thời gian thực (FDS). Các sự kiện chấp thuận thẻ được nhóm theo khóa số thẻ, số lần và số tiền chấp thuận được tổng hợp trong một cửa sổ tumbling 1 phút, và một cảnh báo chặn được kích hoạt nếu vượt ngưỡng. Vì theo bản chất của thanh toán di động, một số sự kiện đến muộn vài giây, độ trễ watermark được đặt 5 giây để bao gồm các lượt đến muộn, còn các sự kiện muộn hơn thế được gửi tới một side output và phản ánh vào quyết toán hậu kỳ. Phán đoán "cùng một khoảng 1 phút" chỉ nhất quán khi dựa trên thời điểm chấp thuận (event time), không phải processing time.

Như một trường hợp IoT công nghiệp, tách các phiên vận hành bằng cửa sổ session (ngưỡng khoảng trống 10 phút) từ các luồng cảm biến của hàng nghìn máy cho phép phát hiện sớm bất thường từ xu hướng rung động và nhiệt độ trong một phiên. Với thiết bị mà vận hành và dừng là không đều, cửa sổ session phản ánh các khoảng vận hành thực tế tự nhiên hơn cửa sổ cố định. Một dịch vụ gọi xe toàn cầu cũng được biết là vận hành xử lý luồng dựa trên Flink ở quy mô lớn cho khớp cung-cầu thời gian thực và dự đoán ETA, theo cùng mạch đó.

Trong gợi ý thương mại điện tử, một cửa sổ sliding được dùng. Tổng hợp các cú nhấp của mỗi người dùng trong 10 phút gần nhất với khoảng dịch 1 phút để phản ánh xu hướng sản phẩm quan tâm theo thời gian thực khiến gợi ý phản ứng nhạy với hành vi mới nhất khi phiên kéo dài. Ở đây, vì một sự kiện được bao gồm trùng lặp trong nhiều cửa sổ và trạng thái phình to, ta thiết kế TTL trạng thái cùng chu kỳ checkpoint (ví dụ 10 giây) để quản lý bộ nhớ và thời gian khôi phục. Chẳng hạn với 1 triệu người dùng hoạt động × trạng thái trung bình 1 KB, khoảng 1 GB trạng thái theo khóa được duy trì mọi lúc, nên quyết định chuyển state backend từ bộ nhớ sang RocksDB và làm hết hạn trạng thái người dùng nhàn rỗi bằng TTL ảnh hưởng trực tiếp đến thông lượng và chi phí.

4. Chuyên sâu: Cơ chế hiện thực exactly-once và xu hướng mới nhất

Exactly-once là bài toán bảo đảm "hiệu ứng chỉ một lần" chứ không phải "truyền chỉ một lần". Flink hiện thực điều này qua Asynchronous Barrier Snapshotting. Khi nguồn định kỳ chèn một barrier (barrier) đặc biệt vào dòng dữ liệu, mỗi toán tử chụp snapshot trạng thái của chính nó khi barrier đi qua và lưu lại. Khi snapshot của mọi toán tử được tụ hội, một checkpoint nhất quán toàn cục (globally consistent) hoàn thành, và khi sự cố thì trạng thái cùng offset nguồn được tua về điểm này cùng nhau. Điểm cốt lõi là suy giảm hiệu năng nhỏ vì snapshot được chụp bất đồng bộ mà không dừng xử lý. Đây là sự biến thể của lý thuyết snapshot phân tán (thuật toán Chandy-Lamport) cho phù hợp với dòng dữ liệu luồng.

Vì chỉ khôi phục trạng thái thì không thể ngăn trùng lặp của đầu ra đã rời đi tới sink, phía đầu ra kết hợp ghi lũy đẳng hoặc ghi giao dịch (two-phase commit). Ví dụ TwoPhaseCommitSinkFunction của Flink commit giao dịch ngoại vi khớp với thời điểm hoàn tất checkpoint, ràng buộc checkpoint và khả năng hiển thị đầu ra một cách nguyên tử. Kafka cung cấp exactly-once cho I/O nội bộ Kafka qua mẫu read-process-write, vốn ràng buộc một giao dịch (giao dịch producer) và commit offset consumer vào một giao dịch duy nhất, còn Kafka Streams đơn giản hóa điều này bằng thiết lập processing.guarantee=exactly_once_v2 (các cải tiến trong dòng KIP-447 được cho là đã tăng hiệu quả khi xử lý nhiều phân vùng).

Trong khi đó, Spark Structured Streaming bảo đảm khôi phục sự cố dựa trên một WAL (Write-Ahead Log) ghi offset và trạng thái vào checkpoint, bảo đảm độ chính xác end-to-end khi kết hợp với một sink lũy đẳng hoặc giao dịch. Kafka Streams sao chép trạng thái sang một topic changelog và khôi phục nó trên một instance khác, nên một ứng dụng có trạng thái có thể được mở rộng ngang và tái bố trí như một dịch vụ không trạng thái. Dù cách hiện thực cơ chế khác nhau theo từng engine, điểm cốt lõi là hiểu rằng nguyên lý chung — "một snapshot trạng thái nhất quán + bảo đảm lũy đẳng hoặc giao dịch cho đầu ra" — là như nhau.

Trong các xu hướng mới nhất, một số dòng chảy nổi bật. Thứ nhất, bốn trục được xác lập trong mô hình Dataflow của Google — "Cái gì (What), Ở đâu (Where, cửa sổ), Khi nào (When, watermark và trigger), và Ra sao (How, hiệu chỉnh dữ liệu muộn)" — đã lan tỏa qua Apache Beam như một trừu tượng trung lập với engine, củng cố hướng tái sử dụng một định nghĩa đường ống duy nhất trên nhiều engine thực thi. Thứ hai, streaming SQL (Flink SQL và tương tự), xử lý một luồng như một bảng qua SQL, đã trưởng thành, nâng cao khả năng tiếp cận cho các vai trò không phải lập trình viên. Thứ ba, có nỗ lực sôi nổi hướng tới một kiến trúc tách trạng thái-tính toán tách kho trạng thái sang kho lưu trữ đối tượng đám mây để mở rộng tính toán và trạng thái độc lập, và hướng tới hợp nhất luồng và batch trên một định dạng bảng duy nhất (ví dụ một lakehouse như Apache Iceberg). Dù vậy, các con số và phiên bản chi tiết này thay đổi nhanh, nên nên kiểm chứng chúng với tài liệu chính thức mới nhất khi thiết kế.

5. Những điều cần cân nhắc và hàm ý

Thứ nhất, thiết kế bằng cách chuyển sự đánh đổi độ trễ–tính đầy đủ thành yêu cầu nghiệp vụ. Độ trễ cho phép của watermark và kích thước cửa sổ, trước khi là tham số kỹ thuật, là phán đoán nghiệp vụ về "dữ liệu muộn đến mức nào thì còn đáng phản ánh vào kết quả". Nơi tính đầy đủ quan trọng, như quyết toán, hãy đặt độ trễ cho phép rộng rãi và cung cấp một tuyến hiệu chỉnh hậu kỳ; nơi tính tức thời quan trọng, như cảnh báo, hãy đặt nó ngắn và hiệu chỉnh các lượt đến muộn riêng. Dưới góc nhìn của Kỹ sư chuyên nghiệp, cách tiếp cận nên làm là định nghĩa SLA và yêu cầu độ chính xác trước rồi tính ngược ra các tham số.

Thứ hai, làm rõ phạm vi và chi phí của exactly-once. Exactly-once thường bị giới hạn ở "trạng thái engine và một sink cụ thể" và không tự động mở rộng ra toàn bộ hệ thống ngoại vi. Các tích hợp ngoại vi mà thiết kế lũy đẳng là bất khả thi (ví dụ email hay thanh toán hướng ra ngoài) khó ràng buộc trong một giao dịch, nên phải bổ khuyết bằng at-least-once cộng với thiết kế khóa lũy đẳng. Ngoài ra, rút ngắn chu kỳ checkpoint giảm thời gian khôi phục nhưng tăng chi phí phụ trội, nên phải tìm điểm cân bằng giữa mục tiêu thời gian khôi phục (RTO) và thông lượng.

Thứ ba, thiết kế quản lý trạng thái và chiến lược tái xử lý cùng nhau. Để giữ trạng thái theo khóa không phình to vô hạn, hãy đặt sẵn TTL, một state backend (ví dụ RocksDB) và giám sát kích thước trạng thái, và chuẩn bị trước thời gian lưu giữ của broker, một chiến lược tua lại offset, và một cấu trúc tái tạo lũy đẳng cho các view kết quả để dữ liệu quá khứ có thể được tái xử lý an toàn khi logic thay đổi. Thiết kế gắn với triết lý tái xử lý của kiến trúc Kappa ([[lambda-kappa-architecture]]) có thể hạ độ phức tạp vận hành.

Thứ tư, nội tại hóa khả năng quan sát vận hành và xử lý áp lực ngược (backpressure). Độ trễ watermark tăng đột biến, tỷ lệ thất bại checkpoint, độ trễ tiêu thụ (consumer lag), và các khoảng xảy ra áp lực ngược là các chỉ số sức khỏe của một đường ống thời gian thực. Bằng cách kết hợp điều khiển áp lực ngược của reactive streams ([[reactive-streams-backpressure]]) với quan sát phân vùng và offset của hàng đợi thông điệp ([[message-queue-kafka]]), ta phải tự động điều tiết thông lượng khi đột biến và giữ trạng thái không đi lệch khỏi phạm vi có thể khôi phục.

Thứ năm, triển vọng và các công nghệ liên quan. Streaming SQL, các trừu tượng trung lập với engine dựa trên Beam, và hợp nhất luồng–batch dựa trên lakehouse đang hội tụ theo hướng "làm mờ ranh giới giữa thời gian thực và batch". Nhìn cùng với các mô hình nhất quán dữ liệu như nhất quán cuối cùng ([[eventual-consistency]]) và CQRS với event sourcing ([[cqrs-event-sourcing]]), xử lý luồng nên được hiểu không phải như một tính năng đơn lẻ mà như một trục thiết kế xuyên suốt toàn bộ nền tảng dữ liệu thời gian thực.

Tài liệu tham khảo


Tóm tắt một câu: Xử lý luồng nhóm dữ liệu vô hạn thành các đơn vị hữu hạn qua cửa sổ hóa, xác định khi nào chốt phép tổng hợp bằng event time và watermark, và bảo đảm xử lý chính xác một lần qua checkpoint và sink giao dịch — một công nghệ xử lý dữ liệu thời gian thực để thiết kế sự đánh đổi giữa độ trễ, độ chính xác và tính đầy đủ.