Back to Blog
System Architecture

Real-Time Data Processing: Flink + Kafka Best Practices

📅 2026.01.20 ⏱️ 10 min 👤 Eric Pan

Why Streaming?

In the era of big data, traditional batch processing can no longer meet business requirements for real-time analytics. User behavior analysis, real-time risk control, and IoT data monitoring all require data-to-insight latency controlled at the second or even millisecond level. The combination of Apache Flink and Apache Kafka provides a mature technical solution for building high-throughput, low-latency, Exactly-Once semantic real-time data processing pipelines.

This article shares practical experience from course projects and internships using Flink + Kafka to build real-time data processing systems, covering architecture design, core concepts, code implementation, and production deployment.

Stream processing is not merely faster batch processing. It is a different data-processing paradigm: process data as it arrives, rather than waiting for a batch to accumulate.

Overall Architecture

A typical Flink + Kafka real-time data processing architecture includes the following core components:

The core advantage of this layered architecture is that each layer can be independently scaled and replaced. For example, when supporting new data sources, you only need to add new Kafka Producers without modifying Flink job code.

Kafka Topic Design

Kafka Topic design directly impacts system throughput and fault tolerance. Here are our best practices:

A key design decision: should large messages (like complete user behavior events) go directly into Kafka? The answer is to avoid it. Kafka handles large messages inefficiently — store large payloads (over 1KB) in object storage (S3/MinIO), with Kafka messages containing only reference pointers.

Flink Real-Time Aggregation

Below is a complete real-time user behavior aggregation example, demonstrating how to read from Kafka, group by user, and compute behavior statistics within 5-minute tumbling windows:

UserBehaviorAggregator.java
public class UserBehaviorAggregator {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60000); // 60s checkpoint
// Read user activity events from Kafka
KafkaSource<UserEvent> source = KafkaSource.<UserEvent>builder()
.setBootstrapServers("kafka:9092")
.setTopics("user.events")
.setGroupId("flink-behavior-agg")
.setValueOnlyDeserializer(new UserEventDeserializer())
.build();
// Aggregate over five-minute tumbling windows
env.fromSource(source, WatermarkStrategy
.<UserEvent>forBoundedOutOfOrderness(Duration.ofSeconds(10))
.withTimestampAssigner((e, t) -> e.getTimestamp()),
KafkaSourceOffsets.initializer())
.keyBy(UserEvent::getUserId) .window(TumblingEventTimeWindows.of(Time.minutes(5))) .aggregate(new BehaviorCountAggregate()) .sinkTo(kafkaSink); } }

Key configuration notes: forBoundedOutOfOrderness(10s) allows data to arrive up to 10 seconds late; enableCheckpointing(60000) takes a state snapshot every 60 seconds, ensuring Exactly-Once semantics.

Exactly-Once Semantics

Exactly-Once is the core challenge of stream processing systems. Flink achieves end-to-end Exactly-Once through Checkpoint + two-phase commit protocol:

Configuration essentials:

Exactly-once processing has a performance cost. In our tests, enabling it reduced throughput by about 15%, but for finance and risk-control workloads, data accuracy matters much more than throughput.

State Management

Flink's state management capability is its core advantage over simple message forwarders. Proper State usage enables complex stream processing logic:

State backend selection:

We use the RocksDB + incremental Checkpoint combination, reducing Checkpoint time from 45 seconds (full) to 3-5 seconds (incremental), significantly reducing Checkpoint interference with normal processing.

Performance Tuning

Below are our Flink + Kafka performance tuning experiences from actual projects:

Before and after tuning comparison:

Summary

The Flink + Kafka combination provides an industrial-grade solution for real-time data processing. Core advantages:

Of course, this solution also has challenges like a steep learning curve and high operational complexity. Start with small-scale PoCs, gradually accumulate experience, then expand to production. For teams, stream processing isn't a "quick win" technology, but an engineering practice requiring sustained investment and optimization.