Streaming 9 — Real-time trong Ngân hàng

22 thg 7, 2026 2 lượt xem
#banking
#aml
#data-engineering
#streaming
#real-time
#fraud-detection

Khi độ trễ là tiền và là rủi ro

Đây là bài khép lại series Streaming & Real-time Data chuyên sâu. Bảy bài trước xây dựng nguyên lý — event time, windowing, state, exactly-once, CDC. Bài này gom tất cả lại quanh một ngành mà độ trễ không chỉ là trải nghiệm người dùng mà là tiền thật và rủi ro pháp lý: ngân hàng.

Trong batch, "chậm một giờ" nghĩa là báo cáo cập nhật muộn một giờ — khó chịu nhưng không chết. Trong ngân hàng real-time, "chậm 200ms" có thể là ranh giới giữa việc chặn được một giao dịch gian lận và việc tiền đã rời khỏi tài khoản không thể thu hồi. NAPAS chuyển khoản 24/7 hoàn tất trong vài giây; một khi ghi Có cho người thụ hưởng, giao dịch gần như không đảo ngược. Cửa sổ để "suy nghĩ" của hệ thống chống gian lận vì thế chỉ còn tính bằng mili-giây, đồng bộ ngay trong luồng cấp phép (authorization).

Lưu ý sandbox: môi trường minh hoạ dùng PostgreSQL, không có Kafka/Flink cluster. Toàn bộ mã và cấu hình bên dưới là minh hoạ để đọc hiểu, không chạy trực tiếp trong sandbox. Các con số độ trễ/ngưỡng là ví dụ minh hoạ, không phải hằng số production.

Bốn bài toán real-time cốt lõi của ngân hàng

Trước khi vẽ kiến trúc, hãy phân biệt bốn nhóm bài toán vì chúng có ràng buộc độ trễ và mô hình xử lý rất khác nhau.

Bài toánĐộ trễ mục tiêu (minh hoạ)Kiểu quyết địnhRàng buộc nặng nhất
Phát hiện gian lận (fraud) trong cấp phép50–300ms, đồng bộ trong luồngChặn / cho qua / thử thách (OTP)Độ trễ cực thấp, không được chặn nhầm quá nhiều
Thanh toán tức thời (instant payment)Vài giây end-to-endGhi Nợ/Có, xác nhậnĐúng-một-lần cho tiền, HA
Cảnh báo biến động số dưVài giây, bất đồng bộGửi thông báoKhông mất/nhân đôi thông báo
Giám sát AMLGiây → phút, near-real-timeSinh alert/case điều traKhông bỏ sót, truy vết đầy đủ

Điểm chung: tất cả đều là luồng sự kiện giao dịch phát ra từ core banking và các switch (thẻ, chuyển khoản). Điểm khác: fraud và payment nằm trên đường đi quan trọng (critical path) của giao dịch — nếu hệ real-time chậm/chết thì giao dịch bị treo; còn cảnh báo số dư và AML nằm cạnh đường đi (side path) — có thể chậm hơn mà không làm treo giao dịch. Phân biệt này quyết định cách bố trí kiến trúc.

Kiến trúc tham chiếu end-to-end

Mô hình tinh thần: nguồn giao dịch → xương sống sự kiện → bộ máy xử lý luồng → quyết định + lưu trữ phân tích.

Đọc sơ đồ theo bốn tầng:

  • Nguồn. Core banking phát sinh bút toán; switch thẻ và cổng NAPAS phát sinh yêu cầu cấp phép. Với core, cách sạch nhất để lấy sự kiện mà không đụng vào DB giao dịch là CDC — đọc từ transaction log qua Debezium, biến mỗi thay đổi row thành một event. Với switch, sự kiện authorization được đẩy thẳng vào Kafka.
  • Xương sống. Kafka là log bền vững, phân vùng theo khoá (thường là số tài khoản/số thẻ) để đảm bảo mọi sự kiện của cùng một tài khoản đi tuần tự vào cùng một partition — điều kiện tiên quyết để tính velocity đúng thứ tự.
  • Bộ máy. Flink (hoặc Spark Structured Streaming cho các luồng near-real-time) chạy song song rule engine, tính feature từ state có khoá và chấm điểm ML. Kết quả hợp nhất thành một quyết định.
  • Đầu ra. Quyết định fraud quay lại luồng authorization đồng bộ trong ngân sách độ trễ; alert AML/số dư đi bất đồng bộ; đồng thời mọi sự kiện đổ vào lakehouse (medallion bronze/silver/gold) để huấn luyện lại mô hình và báo cáo tuân thủ.

Đây chính là kiến trúc Kappa (xem Lambda vs Kappa ở bài tổng quan): một đường luồng duy nhất phục vụ cả real-time lẫn (qua replay từ Kafka/lakehouse) tính lại lịch sử, thay vì duy trì song song hai code path batch và stream.

Phát hiện gian lận: rule + ML trên cùng một luồng

Chống gian lận real-time hầu như luôn là hệ lai (hybrid), không phải chọn một:

  • Rule engine — luật tường minh, giải thích được, phản ứng tức thì với mẫu đã biết: "quẹt thẻ ở hai quốc gia trong 10 phút", "hơn 5 giao dịch trong 1 phút", "số tiền > 50 triệu ngoài giờ hành chính". Rẻ, minh bạch, dễ audit — quan trọng khi phải giải trình với cơ quan quản lý.
  • ML scoring — mô hình (gradient boosting, hoặc mạng nơ-ron) chấm một điểm rủi ro liên tục từ hàng chục feature, bắt được mẫu tinh vi mà luật cứng bỏ lỡ. Bù lại khó giải thích hơn và cần feature chất lượng.

Mấu chốt kỹ thuật của phần ML là feature real-time phải khớp với feature lúc huấn luyện — nếu không sẽ dính training/serving skew. Ví dụ feature "số giao dịch của thẻ này trong 5 phút qua" phải được tính y hệt ở hai nơi: lúc train (trên dữ liệu lịch sử trong lakehouse) và lúc serve (trên state của Flink). Đây là lý do người ta dùng feature store hoặc, tối thiểu, chia sẻ cùng một định nghĩa cửa sổ giữa batch và stream.

Các feature real-time điển hình đều là hàm của cửa sổ theo event time trên state có khoá theo tài khoản/thẻ:

  • đếm số giao dịch trong cửa sổ trượt 1/5/60 phút (velocity),
  • tổng/độ lệch số tiền so với trung bình 30 ngày,
  • số merchant/quốc gia khác nhau trong cửa sổ,
  • khoảng cách địa lý giữa hai lần quẹt liên tiếp (impossible travel).

Cửa sổ velocity — trái tim của phát hiện chuỗi bất thường

Velocity là tần suất và cường độ giao dịch trong một khoảng thời gian ngắn. Một tài khoản bình thường quẹt vài lần một ngày; khi thẻ bị đánh cắp, kẻ gian thường "vét" nhanh — nhiều giao dịch nhỏ liên tiếp để test thẻ, rồi vài giao dịch lớn. Cửa sổ velocity bắt đúng dạng bùng nổ đột ngột này.

Sơ đồ trên minh hoạ cửa sổ trượt (sliding window) dài 2 phút: mỗi sự kiện mới cập nhật bộ đếm và tổng tiền trong 2 phút gần nhất; các sự kiện cũ hơn 2 phút tự động rời khỏi cửa sổ. Khi count hoặc sum vượt ngưỡng, hệ ra quyết định chặn ngay trong luồng cấp phép.

Diễn đạt cùng logic bằng SQL luồng (phong cách Flink SQL) — chú ý HOP là cửa sổ trượt, và dùng event time txn_time chứ không phải thời điểm xử lý:

-- Velocity: đếm giao dịch & tổng tiền trong cửa sổ trượt 2 phút,
-- trượt mỗi 20 giây, theo từng số thẻ. Event time = txn_time.
SELECT
    card_id,
    window_start,
    window_end,
    COUNT(*)                         AS txn_count,
    SUM(amount)                      AS txn_amount,
    COUNT(DISTINCT merchant_country) AS distinct_countries
FROM TABLE(
    HOP(
        TABLE transactions,
        DESCRIPTOR(txn_time),        -- cột event time đã gắn watermark
        INTERVAL '20' SECOND,        -- bước trượt
        INTERVAL '2' MINUTE          -- độ dài cửa sổ
    )
)
GROUP BY card_id, window_start, window_end
HAVING COUNT(*) > 4                  -- quá nhiều giao dịch, hoặc
    OR SUM(amount) > 50000000        -- tổng tiền vượt 50tr / 2 phút, hoặc
    OR COUNT(DISTINCT merchant_country) > 1;  -- đa quốc gia bất thường

Bảng transactions phải được khai báo cột txn_timeevent time với watermark (ví dụ WATERMARK FOR txn_time AS txn_time - INTERVAL '30' SECOND) để xử lý dữ liệu đến trễ đúng cách: giao dịch từ POS mất sóng đến muộn 20 giây vẫn được xếp vào đúng cửa sổ theo lúc quẹt thật, không phải lúc Kafka nhận.

Với các luật cần state phức tạp hơn HOP (ví dụ "impossible travel" cần nhớ vị trí và thời điểm của giao dịch trước đó), người ta rơi xuống tầng DataStream API và dùng KeyedProcessFunction với ValueState/MapState + timer — đúng mô hình đã mô tả ở bài state & exactly-once.

Thanh toán tức thời: đúng-một-lần cho tiền

Instant payment (NAPAS 24/7, chuyển khoản nhanh) đặt ra ràng buộc khắc nghiệt nhất: một lệnh chuyển tiền phải được ghi đúng một lần. Xử lý thiếu → khách mất tiền; xử lý thừa → ghi Nợ hai lần. Exactly-once ở đây không phải "cho đẹp" mà là bắt buộc.

Ba lớp bảo vệ, kết hợp:

  1. Khoá idempotency ở nghiệp vụ. Mỗi lệnh chuyển mang một transaction_reference duy nhất (NAPAS cấp). Hệ hạ nguồn khử trùng lặp theo khoá này: nếu đã thấy ref rồi thì bỏ qua bản lặp. Đây là tuyến phòng thủ vững nhất vì không phụ thuộc đảm bảo của tầng vận chuyển.
  2. Exactly-once trong engine + sink giao dịch. Flink với checkpoint + two-phase commit sink đảm bảo một sự kiện chỉ tác động state và output đúng một lần kể cả khi job restart. Với sink là DB, dùng transaction; với sink là Kafka, dùng transactional producer.
  3. CDC thay vì dual-write. Đừng để dịch vụ vừa ghi core DB vừa tự publish Kafka (dual-write dễ lệch khi một trong hai lỗi). Thay vào đó, ghi vào DB rồi CDC phát sự kiện từ log — sự kiện và trạng thái DB không bao giờ lệch nhau.

Về HA: đường thanh toán nằm trên critical path nên phải chịu được lỗi node. Kafka replication factor ≥3, Flink checkpoint xuống lưu trữ bền vững, và job cấu hình restart tự động — mất một task manager thì khôi phục từ checkpoint gần nhất mà không mất/nhân đôi giao dịch.

Cảnh báo số dư và giám sát AML — nhánh bất đồng bộ

Cảnh báo biến động số dư là bài toán "nhẹ" nhất về logic nhưng nhạy về trải nghiệm: khách kỳ vọng nhận thông báo trong vài giây sau khi tài khoản thay đổi. Nguồn tự nhiên là chính luồng CDC của bảng tài khoản: mỗi sự kiện ghi Nợ/Có → tra ngưỡng đăng ký của khách → đẩy thông báo. Vì nằm ở side path, nó không được phép làm chậm giao dịch; và cần khử trùng lặp để không gửi hai lần cùng một thông báo khi job replay.

Giám sát AML (phòng chống rửa tiền) là near-real-time chứ không cần mili-giây, nhưng đòi hỏi state dài hạn và không bỏ sót. Các mẫu điển hình:

  • Structuring/smurfing — chia nhỏ nhiều giao dịch dưới ngưỡng báo cáo để né. Bắt bằng cửa sổ tổng theo ngày/tuần trên state có khoá theo khách.
  • Velocity dòng tiền — tiền vào rồi ra ngay (pass-through), hoặc chuỗi chuyển qua nhiều tài khoản trung gian.
  • Đối chiếu danh sách trừng phạt/PEP theo thời gian thực khi mở giao dịch quốc tế.

Luồng AML tạo alert đưa vào hàng đợi điều tra (case management), đồng thời lưu toàn bộ bằng chứng vào lakehouse để truy vết. Về khung tuân thủ, KYC/CTR/SAR và ngưỡng báo cáo, xem bài Compliance trong series ngân hàng. Vì AML cần nhìn cả lịch sử dài, đây là nơi ranh giới stream/batch mờ đi: cửa sổ ngắn chạy trên luồng, còn phân tích mạng lưới quan hệ (graph) thường chạy batch trên lakehouse rồi feed ngưỡng ngược lại luồng.

Ràng buộc xuyên suốt và cách các khái niệm đã học ghép vào

Gói lại toàn series, ba ràng buộc chi phối mọi lựa chọn thiết kế trong ngân hàng:

  • Đúng-một-lần cho tiền. Bất cứ chỗ nào chạm số dư đều cần idempotency + exactly-once. Không thoả hiệp.
  • Độ trễ thấp có ngân sách rõ. Fraud trong authorization có "budget" mili-giây; vượt là chặn nhầm hoặc bỏ lọt. Kiến trúc phải đo và bảo vệ ngân sách này (p99, không chỉ trung bình).
  • Tính sẵn sàng cao (HA). Đường tiền không được downtime; side path có thể degrade nhưng critical path phải trụ.

Bản đồ khái niệm → ứng dụng ngân hàng:

Khái niệm (bài)Ứng dụng ngân hàng ở bài này
Event time & watermarkXếp giao dịch vào đúng cửa sổ theo lúc quẹt thật, chịu được POS đến trễ
Windowing (tumbling/sliding/session)Cửa sổ velocity trượt cho fraud; cửa sổ ngày cho AML structuring
State & checkpointFeature real-time theo thẻ; nhớ vị trí giao dịch trước (impossible travel)
Exactly-once (2PC / idempotent sink)Ghi tiền đúng một lần trong instant payment
CDC / event-drivenLấy sự kiện từ core mà không dual-write; nguồn cho cảnh báo số dư
Kappa / stream-table dualityMột đường luồng, replay để huấn luyện lại và tính lại lịch sử

Use case thực tế

Bối cảnh (minh hoạ): NCB triển khai chống gian lận thẻ real-time trên luồng authorization của switch. Mỗi ngày ~2 triệu giao dịch thẻ, đỉnh giờ trưa ~120 giao dịch/giây. Ngân sách độ trễ cho quyết định fraud: p99 < 200ms đồng bộ trong cấp phép.

Thiết kế: authorization đẩy sự kiện vào Kafka (khoá = số thẻ, 24 partition, RF=3). Một job Flink giữ keyed state theo thẻ, tính đồng thời: (1) rule velocity HOP 2 phút như SQL ở trên, (2) 18 feature real-time nạp vào mô hình gradient boosting đã train offline trên lakehouse. Quyết định hợp nhất: rule chặn cứng nếu vượt ngưỡng; nếu không, điểm ML > 0.9 → thử thách OTP, > 0.98 → chặn.

Kết quả minh hoạ (không phải số đo thật): so với hệ batch cũ chạy mỗi 15 phút, thời gian phát hiện giảm từ ~15 phút xuống ~150ms; một chuỗi 5 giao dịch "test-rồi-vét" trong 90 giây bị chặn ngay ở giao dịch thứ 5 thay vì phát hiện sau khi tiền đã đi. Cùng lúc, mọi sự kiện đổ về lakehouse (bronze → silver → gold) để tuần nào cũng huấn luyện lại mô hình và sinh báo cáo AML định kỳ.

Ghi nhớ

  • Ngân hàng real-time có bốn nhóm bài toán với ràng buộc khác nhau: fraudpayment trên critical path (đồng bộ, mili-giây/giây); cảnh báo số dưAML trên side path (bất đồng bộ, giây/phút).
  • Kiến trúc tham chiếu: core/switch → CDC/Kafka → Flink (rule + feature + ML) → quyết định về authorization + sink lakehouse — về bản chất là kiến trúc Kappa.
  • Fraud = rule + ML, không chọn một: rule minh bạch/audit được, ML bắt mẫu tinh vi; feature real-time phải khớp feature train để tránh serving skew.
  • Cửa sổ velocity (sliding/HOP theo event time trên keyed state) là công cụ chủ lực bắt chuỗi giao dịch bùng nổ bất thường.
  • Đúng-một-lần cho tiền là bắt buộc: idempotency theo transaction_reference + exactly-once (2PC) + CDC thay cho dual-write.
  • Ba ràng buộc chi phối mọi thiết kế: exactly-once cho tiền, độ trễ thấp có ngân sách p99 rõ, và HA cho đường tiền.
  • Series khép lại ở đây: mọi khái niệm (event time, window, state, exactly-once, CDC) đều tìm được chỗ đứng cụ thể trong một hệ ngân hàng real-time.

Nguồn tham khảo

  • Streaming Systems — Tyler Akidau, Slava Chernyak, Reuven Lax (O'Reilly). Nền tảng event time, watermark, windowing.
  • Designing Data-Intensive Applications — Martin Kleppmann (O'Reilly). Log, exactly-once, stream-table duality, dual-write vs CDC.
  • Fundamentals of Data Engineering — Joe Reis & Matt Housley (O'Reilly). Kiến trúc pipeline, batch vs streaming.
  • Apache Flink — tài liệu chính thức, Windows & Stateful Stream Processing: https://nightlies.apache.org/flink/flink-docs-stable/
  • Apache Kafka — tài liệu chính thức: https://kafka.apache.org/documentation/
  • Debezium — tài liệu chính thức (CDC): https://debezium.io/documentation/
  • "Questioning the Lambda Architecture" — Jay Kreps (O'Reilly Radar, 2014).
  • NAPAS — Công ty CP Thanh toán Quốc gia Việt Nam, dịch vụ chuyển nhanh 24/7: https://napas.com.vn/

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