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?
In brief
- There are two places to put the model: embedded in the Flink job (no network call per event) or on a remote model server (extra latency in exchange for central model management).
- Async I/O lets one instance handle many requests at once, while capacity caps the number of in-flight requests and creates backpressure when the cap is reached.
- Failure recovery comes from checkpoints, not committed offsets; an exactly-once KafkaSink commits its transactions on checkpoint.
- 1Kafka transactions topicKafka feeds the customer's events into Flink for processing
- 2KafkaSource + watermarkAssigns real event time, tolerates late data, withIdleness for quiet partitions
- 3Async I/O calls the model serverMany requests at once; backpressure when capacity is reached
- 4KafkaSink exactly-onceWrites within a Kafka transaction, commits when the checkpoint completes
- 5Risk score topicThe authorisation system reads the score to make its decision
Each stage answers one of the customer's questions: timing, throughput, and scoring exactly once.
Graphic: FDE Times
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.
env.enableCheckpointing(60_000);
KafkaSource<Txn> source = KafkaSource.<Txn>builder()
.setBootstrapServers("broker:9092")
.setTopics("payments.transactions")
.setGroupId("fraud-scoring")
.setStartingOffsets(OffsetsInitializer.earliest())
.setValueOnlyDeserializer(new TxnDeserializer())
.setProperty("partition.discovery.interval.ms", "300000")
.build();
WatermarkStrategy<Txn> wm = WatermarkStrategy
.<Txn>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withTimestampAssigner((txn, ts) -> txn.eventTimeMillis())
.withIdleness(Duration.ofMinutes(1));
DataStream<Txn> 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.
DataStream<Scored> 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.
The final piece writes the results.
KafkaSink<Scored> sink = KafkaSink.<Scored>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.