Streaming 8 — Chất lượng & Observability luồng

22 thg 7, 2026 2 lượt xem
#observability
#data-quality
#data-engineering
#streaming
#schema-registry

Mô hình tinh thần: vì sao giám sát streaming khó hơn batch

Trong batch, mỗi lần chạy là một job có điểm bắt đầu và điểm kết thúc. Kết thúc xong bạn có một mốc rõ ràng để kiểm tra: đếm số bản ghi, đối soát tổng dư nợ nguồn với đích, chạy dbt tests / Great Expectations rồi mới publish. Nếu sai, job "đỏ" và bạn chặn được dữ liệu bẩn trước khi ai đó dùng.

Streaming không có cái mốc đó. Luồng chạy mãi mãi; không có "job kết thúc" để đặt một checkpoint kiểm tra chất lượng cuối cùng. Bản ghi thứ 1 và bản ghi thứ 1 tỷ đi qua cùng một đường ống đang chạy. Điều này đổi ba thứ căn bản:

  • Không kiểm tra trọn tập (batch-level assertion) được. Bạn không thể nói "tổng số giao dịch hôm nay phải khớp core banking" ngay tại thời điểm streaming, vì "hôm nay" chưa kết thúc. Kiểm tra phải chuyển thành per-record (từng bản ghi) hoặc per-window (theo cửa sổ thời gian).
  • Lỗi tích luỹ âm thầm. Một job batch hỏng thì đỏ ngay. Một consumer streaming tụt lại phía sau vẫn "xanh" — nó vẫn chạy, chỉ là dữ liệu ra ngày càng cũ. Triệu chứng là độ trễ tăng dần, không phải exception.
  • Không dừng để sửa. Bạn không thể tắt luồng thanh toán để vá bug rồi chạy lại như backfill batch. Phải xử lý bản ghi lỗi tại chỗ (dead-letter queue) và cho luồng tiếp tục.

Vì thế observability trong streaming không phải "thêm dashboard" — nó là cơ chế thay thế cho cái mốc kết thúc job mà batch có sẵn. Ta chia làm hai tầng: metric hệ thống (đường ống có khoẻ không) và data quality (dữ liệu đi qua có đúng không). Hai tầng độc lập: hệ thống xanh vẫn có thể tuôn ra dữ liệu sai schema, và ngược lại.


Tầng 1 — Metric hệ thống: đường ống có khoẻ không

Đây là các chỉ số về sức khoẻ vận hành của luồng, độc lập với nội dung dữ liệu.

MetricÝ nghĩaBáo động khi
ThroughputSố bản ghi/giây (hoặc byte/s) qua từng chặngTụt đột ngột (nghẽn) hoặc tăng vọt bất thường
Consumer lagKhoảng cách giữa offset mới nhất producer ghi vào và offset consumer đã đọc — đo số bản ghi bị bỏ lại phía sauLag tăng đơn điệu → consumer không theo kịp
BackpressureTín hiệu chặng downstream chậm đang ép chặng upstream chậm lạiRatio backpressure cao kéo dài
Checkpoint duration / failureThời gian và tỷ lệ thất bại khi lưu trạng thái để exactly-onceDuration phình to hoặc checkpoint fail liên tiếp
Watermark lagChênh giữa event time của watermark và processing time hiện tại — đo độ cũ theo thời gian sự kiệnWatermark "đứng hình" hoặc tụt xa

Consumer lag — hiểu cho đúng

Đây là metric quan trọng nhất và cũng hay bị hiểu sai. Trong Kafka, mỗi partition có một log offset tăng dần. Producer ghi vào cuối log (log-end offset), consumer đọc và commit offset đã xử lý. Consumer lag của một partition là:

lag(partition) = log_end_offset(partition) − committed_offset(consumer, partition)

Lag của consumer group = tổng lag trên mọi partition nó phụ trách. Điểm mấu chốt:

  • Lag đo số bản ghi tồn đọng, không phải thời gian. 10.000 bản ghi lag có thể là 1 giây (luồng nhanh) hoặc 1 giờ (luồng chậm). Muốn biết "dữ liệu cũ bao lâu", nhìn thêm watermark lag / time lag (chênh timestamp).
  • Lag ổn định ở mức thấp là bình thường; lag tăng đơn điệu mới là bệnh — nghĩa là tốc độ tiêu thụ < tốc độ sản xuất, khoảng cách sẽ nới ra mãi.
  • Lag nhìn theo từng partition. Một partition nóng (key lệch) có thể lag nặng trong khi tổng vẫn trông ổn.

Xem lag bằng công cụ có sẵn của Kafka:

# LAG = số bản ghi mỗi partition của group đang bị bỏ lại
kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --describe --group txn-enrichment

# TOPIC        PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
# transactions 0          9,812,443       9,812,461       18
# transactions 1          9,780,102       9,900,050       119,948   <-- partition nghẽn

Chi tiết cơ chế Kafka (offset, partition, consumer group) xem Kafka Streaming. Cross-cutting về observability cấp nền tảng xem DataOps — Observability & sự cố.


Backpressure là gì và xử lý thế nào

Backpressure (áp lực ngược) là hiện tượng: khi một chặng downstream xử lý chậm hơn tốc độ dữ liệu đến, nó phát tín hiệu (hoặc đơn giản là không kéo thêm) khiến chặng upstream phải chậm lại theo, và tín hiệu này lan ngược dần về tận nguồn.

Trực giác: hình dung nước chảy qua chuỗi phễu. Phễu cuối tắc, nước dâng lên phễu trước nó, rồi phễu trước nữa. Backpressure chính là "nước dâng ngược" đó trong đường ống dữ liệu.

Vì sao nó tốt (khi được thiết kế đúng): thay vì upstream cứ đẩy dữ liệu vào một buffer đầy rồi OOM/mất dữ liệu, backpressure buộc cả hệ tự điều tiết về tốc độ của chặng chậm nhất. Flink dùng cơ chế credit-based flow control; Kafka thì tự nhiên có backpressure vì consumer pull theo nhịp của mình — consumer chậm chỉ làm lag tăng, chứ không làm sập producer.

Chẩn đoán: backpressure cao cộng consumer lag tăng → tìm chặng chậm nhất (thường operator có I/O ngoài: gọi API, ghi DB, lookup). Cách xử lý theo thứ tự ưu tiên:

  1. Tăng song song (parallelism) ở chặng nghẽn — thêm partition/subtask. Chú ý key skew: nếu một key chiếm phần lớn dữ liệu thì thêm task không cứu được, phải re-key.
  2. Bỏ I/O đồng bộ trong đường nóng — async I/O, batch các lệnh ghi, cache lookup.
  3. Tăng buffer/bộ nhớ chỉ hoãn được vấn đề, không giải quyết gốc.
  4. Nếu downstream thực sự là trần cứng (DB đích), cân nhắc tách bậc: ghi ra topic trung gian, tiêu thụ theo nhịp DB chịu được.

Tầng 2 — Data quality on-stream

Hệ thống khoẻ không đảm bảo dữ liệu đúng. Trên luồng, kiểm tra chất lượng phải làm per-record hoặc per-window.

Schema validation + Schema Registry & compatibility

Nguồn lỗi phổ biến nhất của luồng là producer đổi schema làm consumer parse hỏng. Giải pháp là tách schema ra khỏi từng bản ghi và quản lý tập trung bằng Schema Registry (ví dụ Confluent Schema Registry) với định dạng có schema như Avro hoặc Protobuf:

  • Producer đăng ký schema, registry trả về một schema ID; bản ghi trên wire chỉ mang ID + payload nhị phân (không nhét cả schema vào mỗi message → tiết kiệm). Consumer lấy schema theo ID để giải mã.
  • Trước khi chấp nhận schema mới, registry kiểm tra compatibility theo luật đã đặt cho subject:
Chế độCho phép đổi gìAi nên nâng cấp trước
BACKWARDXoá field, thêm field có defaultConsumer nâng trước
FORWARDThêm field, xoá field có defaultProducer nâng trước
FULLChỉ thay đổi vừa backward vừa forwardBất kỳ thứ tự
NONEKhông kiểm tra— (rủi ro cao)

BACKWARD (mặc định của Confluent) nghĩa là consumer dùng schema mới vẫn đọc được dữ liệu ghi bằng schema — nhờ vậy bạn nâng cấp consumer trước mà không vỡ. Đây là hàng rào chặn "đổi schema phá vỡ downstream" ngay tại tầng hạ tầng, trước khi dữ liệu bẩn kịp lan.

Dead-letter queue (DLQ) cho bản ghi lỗi

Trên batch, một bản ghi hỏng có thể làm fail cả job và bạn sửa rồi chạy lại. Trên stream, không được để một bản ghi độc làm nghẽn cả luồng ("poison pill"). Mẫu chuẩn: bản ghi không parse được / vi phạm rule sẽ được định tuyến sang một topic riêng — dead-letter queue — kèm metadata (lý do lỗi, offset gốc, timestamp), còn luồng chính chạy tiếp.

# Pseudocode — validate per-record, hỏng thì đẩy DLQ, không chặn luồng
for record in stream:
    try:
        evt = registry.deserialize(record.value)      # schema check
        validate_contract(evt)                         # rule nghiệp vụ
        emit_main(enrich(evt))
    except (SchemaError, ContractViolation) as e:
        emit_dlq({
            "raw": record.value,
            "error": str(e),
            "topic": record.topic,
            "partition": record.partition,
            "offset": record.offset,     # để reprocess đúng vị trí sau khi vá
            "ts": now_iso(),
        })
        metrics.incr("dlq.records", tags={"reason": type(e).__name__})

DLQ phải được giám sát như một metric hạng nhất: dlq.rate tăng vọt là dấu hiệu upstream đổi format hoặc có sự cố nguồn. DLQ cũng là nơi để reprocess sau khi vá — đọc lại từ topic lỗi, không phải khôi phục từ log gốc đã trôi qua.

Data contract

DLQ và schema registry thực thi phần "kỹ thuật", nhưng cái quyết định đúng/sai là data contract — thoả thuận rõ ràng giữa đội sản xuất (producer) và tiêu thụ (consumer) về: schema, ngữ nghĩa từng field, ràng buộc (not-null, miền giá trị, đơn vị, enum), và SLA (freshness, completeness). Contract nên là code (ví dụ file YAML/JSON kèm test) nằm trong CI của producer, để đổi vi phạm hợp đồng thì CI đỏ, chặn từ lúc build chứ không đợi tới lúc dữ liệu chảy. Xem thêm góc quản trị ở Quản lý chất lượng dữ liệu.


Phát hiện anomaly / drift real-time

Schema validation bắt lỗi cấu trúc; anomaly/drift bắt lỗi phân phối — dữ liệu vẫn hợp lệ về schema nhưng "sai một cách bất thường".

  • Anomaly (bất thường tức thời): ví dụ throughput topic transactions rơi 90% trong 5 phút (upstream chết), hoặc tỷ lệ amount = 0 tăng đột biến. Cách làm nhẹ: tính thống kê trượt theo cửa sổ (đếm, trung bình, tỷ lệ null) rồi so với ngưỡng động (ví dụ z-score theo baseline giờ/ngày trong tuần).
  • Drift (trôi phân phối): phân phối một feature dịch dần theo thời gian — ví dụ tỷ trọng giao dịch theo kênh (ATM/POS/online) đổi cấu trúc. Quan trọng cho feature nuôi mô hình fraud: drift âm thầm làm mô hình xuống cấp. Đo bằng khoảng cách phân phối (PSI, KL-divergence) giữa cửa sổ hiện tại và baseline.

Điểm khác biệt với batch: ở đây ngưỡng phải tính trên cửa sổ trượt liên tục, và alert phải phân biệt "biến động ngày lễ" (giảm giao dịch cuối tuần là bình thường) với sự cố thật — nếu không sẽ bị alert fatigue.


SLO cho streaming: freshness và latency

SLO (Service Level Objective) biến "luồng phải nhanh và mới" thành con số đo được và cam kết được.

  • Latency (độ trễ xử lý): thời gian từ khi sự kiện xảy ra (event time) tới khi kết quả xuất hiện ở sink. Nên đo theo phân vị (p50/p95/p99), không phải trung bình — đuôi p99 mới là cái người dùng đau.
  • Freshness (độ tươi): dữ liệu ở đích cũ tối đa bao lâu. Ví dụ SLO: "95% thời gian, số dư khả dụng ở lớp phục vụ trễ ≤ 3 giây so với core." Freshness liên hệ trực tiếp với watermark lagconsumer lag — hai metric hệ thống ở Tầng 1 chính là chỉ báo sớm cho việc sắp vỡ SLO freshness.

Cấu trúc một SLO nên có: SLI (chỉ số đo, ví dụ freshness giây), mục tiêu (p95 ≤ 3s), cửa sổ đánh giá (rolling 28 ngày) và error budget (được phép vượt bao nhiêu %). Khi tiêu hết error budget → dừng thả tính năng mới, ưu tiên vá độ ổn định. Chi tiết vận hành SLO/incident xem DataOps — Observability & sự cố.


Use case thực tế

(Số liệu minh hoạ, không phải số thật NCB.) Luồng làm giàu giao dịch thẻ real-time để chấm điểm fraud: CDC từ core (xem CDC & event-driven) → topic transactions → job enrich (join thông tin khách + merchant) → sink cho engine fraud. Xem tổng quan kiến trúc ở Streaming — Tổng quan.

Sự cố: 09:14, đội core đẩy bản cập nhật thêm field channel nhưng đổi kiểu amount từ long sang string. Diễn biến quan sát qua observability:

  • 09:14dlq.rate topic transactions vọt từ ~0 lên ~1.200 bản ghi/phút; schema registry log INCOMPATIBLE cho subject transactions-value (luật BACKWARD chặn được ở môi trường staging, nhưng producer này ghi thẳng nên bản ghi lỗi rơi vào DLQ ở consumer).
  • 09:16 — engine fraud không nhận đủ dữ liệu; consumer lag partition nóng tăng, freshness SLO (p95 ≤ 3s) bắt đầu vỡ, freshness thực tế lên ~40s.
  • Nhờ DLQ, luồng chính không sập — các bản ghi hợp lệ vẫn chảy, chỉ bản lỗi bị tách ra. Đội on-call rollback schema producer, rồi reprocess DLQ (dùng offset lưu kèm) để bù các giao dịch bị tách.

Bài học: chính hàng rào schema registry + DLQ + SLO alert đã biến một thay đổi phá vỡ thành sự cố khoanh vùng được và phục hồi được, thay vì mất dữ liệu giao dịch âm thầm.

Ghi nhớ

  • Streaming không có "job kết thúc" để soát cuối cùng → kiểm tra chuyển từ batch-level sang per-record/per-window; observability thay thế cho cái mốc kết thúc mà batch có sẵn.
  • Tách hai tầng: metric hệ thống (đường ống khoẻ không) và data quality (dữ liệu đúng không) — chúng độc lập, xanh tầng này không suy ra xanh tầng kia.
  • Consumer lag = số bản ghi tồn đọng (không phải thời gian); bệnh là khi lag tăng đơn điệu. Muốn biết độ cũ theo thời gian thì nhìn watermark/time lag.
  • Backpressure là downstream chậm ép upstream chậm lại, lan ngược về nguồn — nó bảo vệ hệ khỏi OOM; xử lý gốc bằng tăng parallelism và bỏ I/O đồng bộ trong đường nóng, không phải tăng buffer.
  • Schema Registry + Avro/Protobuf + compatibility (BACKWARD…) chặn thay đổi phá vỡ ngay tầng hạ tầng; bản ghi lỗi đẩy vào DLQ để không "poison pill" cả luồng, và để reprocess sau khi vá.
  • Data contract là thoả thuận producer–consumer, nên là code trong CI để vi phạm thì CI đỏ.
  • Anomaly bắt bất thường tức thời, drift bắt trôi phân phối (quan trọng cho feature mô hình) — đều đo trên cửa sổ trượt, cẩn thận alert fatigue.
  • SLO streaming đo bằng freshnesslatency phân vị (p95/p99); consumer lag & watermark lag là chỉ báo sớm cho việc sắp vỡ SLO freshness.

Nguồn tham khảo

  • Streaming Systems — Tyler Akidau, Slava Chernyak, Reuven Lax (O'Reilly) — event time, watermark, độ trễ.
  • Designing Data-Intensive Applications — Martin Kleppmann (O'Reilly) — schema evolution, compatibility, stream processing.
  • Fundamentals of Data Engineering — Joe Reis & Matt Housley (O'Reilly) — data quality, observability, undercurrents.
  • Apache Kafka Documentation — Consumer groups & offsets, monitoring/lag: https://kafka.apache.org/documentation/
  • Apache Flink Documentation — Monitoring back pressure & checkpointing: https://nightlies.apache.org/flink/flink-docs-stable/
  • Confluent Schema Registry Documentation — schema compatibility types (BACKWARD/FORWARD/FULL): https://docs.confluent.io/platform/current/schema-registry/
  • Google Dataflow model paper — Akidau et al., VLDB 2015 (out-of-order, watermark).
  • Google SRE Book — chương Service Level Objectives (SLI/SLO/error budget): https://sre.google/books/

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