# Kafka và Flink cho FDE: ghép dòng dữ liệu của khách vào model chấm điểm thời gian thực

> Model chạy tốt trên notebook là chuyện dễ. Khách sẽ hỏi bạn những câu khó hơn: có mất message không, có chấm trùng không, model chậm thì pipeline sẽ ra sao.

Bản gốc: https://fdetimes.net/vi/bach-khoa/kafka-flink-du-lieu-streaming-cho-du-doan-thoi-gian-thuc/

Thử hình dung khách của bạn là một cổng thanh toán. Đội data science giao cho bạn một model chấm điểm gian lận cho kết quả đẹp trên dữ liệu lịch sử, còn đội platform chỉ cho bạn một topic Kafka nơi giao dịch đang đổ về liên tục. Mỗi giao dịch phải có điểm rủi ro trước khi hệ thống duyệt kịp ra quyết định.

Ngân hàng và cổng thanh toán chính là nhóm điển hình dùng Flink gọi model từ xa để phân tích giao dịch theo thời gian thực, theo mô tả của Confluent. Trong kiến trúc đó, Kafka là lớp quen thuộc đưa dữ liệu vào Flink, còn Flink lo xử lý và gọi model.

Với một FDE, phần khó ít khi nằm ở model. Phần khó là ba câu khách sẽ hỏi ngay khi có sự cố đầu tiên: có mất message không, có giao dịch nào bị chấm hai lần không, model chậm thì chuyện gì xảy ra. Câu trả lời cho cả ba nằm trong vài dòng cấu hình mà bạn cần hiểu kỹ.

## Model nên sống trong job hay ngoài job?

Kai Waehner, người viết nhiều về Kafka, chia cách triển khai thành hai kiểu. Kiểu embedded nhúng model thẳng vào ứng dụng xử lý luồng, nên mỗi lần chấm điểm không tốn một lượt gọi mạng. Kiểu remote gửi request tới model server qua RPC, API hoặc HTTP rồi chờ response.

| | Embedded | Remote |
|---|---|---|
| Model chạy ở đâu | Bên trong job Flink | Trên model server riêng |
| Gọi mạng mỗi sự kiện | Không | Có, qua RPC, API hoặc HTTP |
| Quản lý model | Đi cùng job: đổi model là deploy lại job | Tập trung ở một chỗ |
| Cái giá phải trả | Ít linh hoạt hơn | Thêm độ trễ |

Dòng "quản lý model" của cột embedded là suy luận từ chính cách nhúng: model nằm bên trong job, nên muốn thay phiên bản model thì bạn phải deploy lại cả job Flink.

Confluent mô tả remote inference là lựa chọn cho hệ thống thông lượng cao, chấp nhận độ trễ để đổi lấy sự linh hoạt. Ở site khách, nên hỏi trước tiên: ai sở hữu model và bao lâu họ retrain một lần?

Nếu đội data science muốn đưa phiên bản mới lên mà không phải chạm vào job Flink, hãy chọn remote. Phần còn lại của bài đi theo hướng đó.

## Một pipeline chấm điểm gian lận, từng đoạn một

Đoạn đầu tiên là đọc dữ liệu và gắn thời gian cho nó.

```java
env.enableCheckpointing(60_000);

KafkaSource source = KafkaSource.builder()
.setBootstrapServers("broker:9092")
.setTopics("payments.transactions")
.setGroupId("fraud-scoring")
.setStartingOffsets(OffsetsInitializer.earliest())
.setValueOnlyDeserializer(new TxnDeserializer())
.setProperty("partition.discovery.interval.ms", "300000")
.build();

WatermarkStrategy wm = WatermarkStrategy
.forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withTimestampAssigner((txn, ts) -> txn.eventTimeMillis())
.withIdleness(Duration.ofMinutes(1));

DataStream txns = env.fromSource(source, wm, "payments");
```

Flink cần một `WatermarkStrategy` gồm hai phần: `TimestampAssigner` lấy thời điểm giao dịch thật từ dữ liệu, và `WatermarkGenerator` quyết định chờ dữ liệu đến trễ bao lâu. `forBoundedOutOfOrderness` với 20 giây nói rằng bạn chấp nhận giao dịch đến muộn tới chừng đó, một con số bạn nên thống nhất với khách thay vì tự đoán.

Dòng `withIdleness` cứu bạn khỏi một lỗi khó thấy. Nếu một partition Kafka không có dữ liệu, nó có thể kìm watermark của cả job đứng yên, và mọi phép tính theo thời gian sẽ chờ mãi. Đánh dấu input là idle sau một phút giải quyết đúng chuyện đó.

Dòng `partition.discovery.interval.ms` cũng đáng giải thích cho khách. KafkaSource tự phát hiện partition mới khi topic scale-out mà không cần restart job, và mặc định cứ 5 phút kiểm tra một lần, tức là đúng con số 300000 mili giây ở trên.

Đoạn thứ hai là gọi model. Cách gọi ở đây quyết định mỗi instance chấm được bao nhiêu giao dịch mỗi giây.

```java
DataStream scored = AsyncDataStream.unorderedWait(
txns,
new ScoreWithModelServer(),   // RichAsyncFunction dùng HTTP client bất đồng bộ
500, TimeUnit.MILLISECONDS,   // timeout mỗi request
100);                         // capacity: số request đang chờ tối đa
```

Toán tử Async I/O gửi request tới model server và nhận response theo kiểu non-blocking. Nhờ vậy một instance song song có thể xử lý nhiều request cùng lúc. Hãy làm một phép tính giả định: model server trả lời trong 50ms.

Gọi đồng bộ trong một `map`, mỗi instance chỉ chấm được 20 giao dịch mỗi giây. Với capacity 100, trần lý thuyết lên tới 2000, miễn là model server chịu nổi. Đừng biến con số 2000 thành lời hứa với khách: trần thật do sức chịu của model server và số instance song song của job quyết định, nên hãy đo trước khi cam kết.

Capacity còn có tác dụng thứ hai. Khi số request đang chờ chạm trần, Async I/O tạo backpressure ngược lên phía Kafka thay vì tích một backlog vô hạn trong bộ nhớ.

**Điểm mấu chốt:** Backpressure không phải lỗi: đó là cách pipeline báo rằng model đang quá tải.

Đoạn cuối cùng là ghi kết quả.

```java
KafkaSink sink = KafkaSink.builder()
.setBootstrapServers("broker:9092")
.setRecordSerializer(new ScoredSerializer("payments.fraud-scores"))
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("fraud-scoring-v1")
.build();

scored.sinkTo(sink);
```

Ở chế độ `EXACTLY_ONCE`, KafkaSink ghi mọi message trong một transaction Kafka và chỉ commit khi checkpoint hoàn tất. Đó là lý do dòng `enableCheckpointing` ở đầu job không phải chuyện tùy chọn. `transactionalIdPrefix` phải là duy nhất, nên hãy đặt tên gắn với job và phiên bản.

## Tự dựng ở site khách theo thứ tự nào?

Đừng ghép cả ba đoạn trong một buổi. Hãy bắt đầu bằng một job chỉ đọc topic và in ra, để chắc rằng deserializer đọc đúng dữ liệu và job có quyền truy cập topic trước khi thêm bất cứ thứ gì khác. Sau đó thêm watermark, nhìn vào Flink UI xem watermark có tiến lên trên mọi partition không.

Bước kế tiếp là thay model thật bằng một stub trả điểm cố định với độ trễ giả lập. Đo thông lượng ở vài mức capacity khác nhau và ghi lại, vì đó là bằng chứng để bạn thương lượng với đội vận hành model server về số request mỗi giây họ phải chịu. Chỉ khi các con số đã ổn mới nối model thật và bật exactly-once.

## Những lỗi khiến pipeline hỏng mà không ai hay

Lỗi đầu tiên là nhìn offset đã commit trong Kafka rồi tin rằng đó là cơ chế khôi phục. KafkaSource có commit offset khi checkpoint hoàn tất, nhưng tài liệu Flink nói rõ nó không dựa vào offset đó để chịu lỗi. Offset đã commit chỉ để giám sát độ trễ của consumer; trạng thái thật nằm trong checkpoint.

Lỗi thứ hai là chạy hai job, chẳng hạn bản cũ và bản thử nghiệm, với cùng một `transactionalIdPrefix`. Yêu cầu prefix duy nhất có lý do, nên hãy đưa tên môi trường và phiên bản vào prefix ngay từ đầu.

Lỗi thứ ba là quên `withIdleness` trên một topic có partition lúc vắng lúc đông, ví dụ giao dịch ban đêm. Triệu chứng là các feature theo cửa sổ thời gian ngừng ra kết quả dù job vẫn "xanh". Lỗi thứ tư là thấy backpressure rồi tăng capacity vô tội vạ, trong khi việc nên làm là hỏi xem model server có đang quá tải không.

## Ghi kỹ năng này vào CV thế nào

Đừng viết "có kinh nghiệm Kafka, Flink". Hãy viết một dòng có số đo từ chính bài tập của bạn: chuyển từ gọi model đồng bộ sang Async I/O, thông lượng tăng từ bao nhiêu lên bao nhiêu, exactly-once qua checkpoint. Kèm theo đó là một repo nhỏ có README trả lời đúng ba câu hỏi của khách: mất message không, chấm trùng không, model chậm thì sao.

Khách không mua một model, họ mua lời cam kết rằng mỗi giao dịch sẽ được chấm đúng một lần, kịp lúc. Nếu bạn chỉ ra được đúng dòng cấu hình đứng sau từng chữ trong lời cam kết đó, khách sẽ tin bạn.

**Thử ngay tuần này:**

- Dựng Kafka và Flink trên máy, viết một model server giả trả điểm ngẫu nhiên sau 50ms rồi so thông lượng giữa gọi đồng bộ và Async I/O với capacity 100.
- Tạo một topic 4 partition, chỉ đẩy dữ liệu vào 3 partition rồi quan sát watermark đứng yên, sau đó thêm withIdleness và xem nó chạy lại.
- Viết một trang runbook cho khách trả lời ba câu: mất message không, chấm trùng không, model chậm thì sao, mỗi câu chỉ vào đúng dòng cấu hình liên quan.

## Nguồn

- [Using Apache Flink® for Model Inference: A Guide for Real-Time AI Applications](https://www.confluent.io/blog/using-flink-for-model-inference-a-guide-for-realtime-ai-applications/)

- [Real-Time Model Inference with Apache Kafka and Flink for Predictive AI and GenAI](https://www.kai-waehner.de/?p=6771)

- [Kafka | Apache Flink](https://nightlies.apache.org/flink/flink-docs-stable/docs/connectors/datastream/kafka/)

- [Generating Watermarks | Apache Flink](https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/event-time/generating_watermarks/)

- [Asynchronous I/O for External Data Access](https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/operators/asyncio/)
