# Kafka and Flink for FDEs: wiring a customer's data stream into a real-time scoring model

> Getting a model to run in a notebook is the easy part. Customers will ask harder questions: will messages be lost, will anything be scored twice, and what happens to the pipeline when the model slows down?

Bản gốc: https://fdetimes.net/en/guides/kafka-flink-real-time-model-scoring/

Picture a customer that runs a payment gateway. The data science team hands you a fraud-scoring model that performs well on historical data, and the platform team points you to a Kafka topic where transactions arrive continuously. Every transaction needs a risk score before the authorisation system makes its decision.

Banks and payment gateways are the typical users of Flink calling a remote model to analyse transactions in real time, according to Confluent. In that architecture Kafka is the familiar layer that feeds data into Flink, while Flink handles processing and calls the model.

For an FDE, the hard part is rarely the model. It is the three questions the customer will ask after the first incident: are messages lost, is any transaction scored twice, and what happens when the model is slow? The answers to all three sit in a few lines of configuration that you need to understand thoroughly.

## Should the model live inside the job or outside it?

Kai Waehner, who writes extensively about Kafka, divides deployments into two styles. The embedded style puts the model directly inside the stream processing application, so each scoring call costs no network round trip. The remote style sends a request to a model server over RPC, an API or HTTP and waits for the response.

| | Embedded | Remote |
|---|---|---|
| Where the model runs | Inside the Flink job | On a separate model server |
| Network call per event | No | Yes, over RPC, API or HTTP |
| Model management | Tied to the job: changing the model means redeploying the job | Centralised in one place |
| The cost | Less flexibility | Added latency |

The "model management" entry for the embedded column is an inference from how embedding works: the model sits inside the job, so swapping model versions means redeploying the whole Flink job.

Confluent describes remote inference as the option for high-throughput systems, trading latency for flexibility. On a customer site, the first question to ask is: who owns the model, and how often do they retrain it?

If the data science team wants to ship new versions without touching the Flink job, choose remote. The rest of this article follows that route.

## A fraud-scoring pipeline, one piece at a time

The first piece reads the data and assigns time to it.

```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 needs a `WatermarkStrategy` with two parts: a `TimestampAssigner` that pulls the real transaction time from the data, and a `WatermarkGenerator` that decides how long to wait for late data. `forBoundedOutOfOrderness` with 20 seconds says you accept transactions arriving up to that late, a figure you should agree with the customer rather than guess.

The `withIdleness` line saves you from a subtle bug. If a Kafka partition has no data, it can hold the watermark of the entire job in place, and every time-based computation will wait indefinitely. Marking the input as idle after one minute solves exactly that.

The `partition.discovery.interval.ms` line is also worth explaining to the customer. KafkaSource discovers new partitions automatically when the topic scales out, without restarting the job, and by default checks every 5 minutes, which is the 300000 milliseconds above.

The second piece calls the model. How you make that call determines how many transactions each instance can score per second.

```java
DataStream scored = AsyncDataStream.unorderedWait(
txns,
new ScoreWithModelServer(),   // RichAsyncFunction using an asynchronous HTTP client
500, TimeUnit.MILLISECONDS,   // per-request timeout
100);                         // capacity: maximum in-flight requests
```

The Async I/O operator sends requests to the model server and receives responses without blocking. That lets one parallel instance handle many requests at once. Take a hypothetical: the model server responds in 50ms.

Called synchronously inside a `map`, each instance can score only 20 transactions per second. With a capacity of 100, the theoretical ceiling rises to 2000, provided the model server can keep up. Do not turn 2000 into a promise to the customer: the real ceiling is set by the model server's capacity and the job's parallelism, so measure before you commit.

Capacity has a second effect. When in-flight requests hit the cap, Async I/O creates backpressure back towards Kafka instead of building an unbounded backlog in memory.

**Điểm mấu chốt:** Backpressure is not a bug: it is how the pipeline tells you the model is overloaded.

The final piece writes the results.

```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);
```

In `EXACTLY_ONCE` mode, KafkaSink writes all messages within a Kafka transaction and commits only when a checkpoint completes. That is why the `enableCheckpointing` line at the top of the job is not optional. The `transactionalIdPrefix` must be unique, so name it after the job and its version.

## In what order should you build it on site?

Do not assemble all three pieces in one session. Start with a job that only reads the topic and prints it, to confirm the deserializer reads the data correctly and the job has access to the topic before adding anything else. Then add watermarks and check in the Flink UI that the watermark advances on every partition.

Next, replace the real model with a stub that returns a fixed score with simulated latency. Measure throughput at several capacity levels and record the results, because they are your evidence when negotiating with the model server's operations team over how many requests per second they must handle. Only once the numbers are stable should you connect the real model and turn on exactly-once.

## Failures that break the pipeline without anyone noticing

The first is looking at committed offsets in Kafka and assuming they are the recovery mechanism. KafkaSource does commit offsets when a checkpoint completes, but the Flink documentation states clearly that it does not rely on those offsets for fault tolerance. Committed offsets exist only to monitor consumer lag; the real state lives in the checkpoint.

The second is running two jobs, such as an old version and an experimental one, with the same `transactionalIdPrefix`. The uniqueness requirement exists for a reason, so put the environment name and version into the prefix from the start.

The third is forgetting `withIdleness` on a topic whose partitions swing between quiet and busy, for example overnight transactions. The symptom is that time-windowed features stop producing results while the job still shows "green". The fourth is seeing backpressure and raising capacity indiscriminately, when the right move is to ask whether the model server is overloaded.

## How to put this skill on your CV

Do not write "experience with Kafka and Flink". Write one line with numbers from your own exercise: moved from synchronous model calls to Async I/O, throughput rose from X to Y, exactly-once via checkpoints. Pair it with a small repo whose README answers the customer's three questions: are messages lost, is anything scored twice, what happens when the model is slow.

Customers are not buying a model. They are buying a commitment that every transaction will be scored exactly once, on time. If you can point to the line of configuration behind each word of that commitment, they will trust you.

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

- Run Kafka and Flink locally, write a fake model server that returns a random score after 50ms, then compare throughput between synchronous calls and Async I/O with a capacity of 100.
- Create a topic with 4 partitions, push data into only 3 of them and watch the watermark stall, then add withIdleness and see it move again.
- Write a one-page runbook for the customer that answers three questions: are messages lost, is anything scored twice, what happens when the model is slow, with each answer pointing to the relevant line of configuration.

## 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/)
