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:
- Data Source Layer (Kafka Producers) — Business systems, log collectors (Flume/Filebeat), and IoT devices write data to Kafka Topics
- Message Buffer Layer (Kafka) — Serves as the data bus, providing persistent storage, decoupling producers and consumers, and supporting parallel reads by multiple consumers
- Compute Layer (Flink) — Consumes data from Kafka, executes window aggregation, stream joins, CEP, and other computation logic, writing results back to Kafka or external storage
- Storage Layer — Persists computation results to MySQL, Redis, Elasticsearch, etc., for downstream application queries
- Display Layer — Real-time Dashboards (Grafana) or business systems consume computation results for visualization and alerting
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:
- Partition Strategy — Partitions = max(target throughput / single partition throughput, consumer instances). We typically start with 6-12 partitions and dynamically adjust based on actual load
- Replication Factor — Production environments recommend replication.factor=3 with min.insync.replicas=2, ensuring no data loss on single-node failure
- Message Retention — Hot data Topics set to 24-hour retention, cold data Topics to 7-day retention, supporting Flink job replay and recomputation
- Message Compression — Uses LZ4 compression, achieving the best balance between throughput and CPU overhead
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 Core Concepts
Understanding Flink's core concepts is the foundation for building reliable stream processing applications:
- Event Time vs Processing Time — Event Time is when data was produced; Processing Time is when it's processed. Production environments should always use Event Time with Watermarks for out-of-order data
- Watermark — Measures Event Time progress. When Watermark exceeds window end time, Flink triggers window computation. Properly setting Watermark delay is key to balancing accuracy and latency
- Checkpoint — Flink's fault tolerance mechanism, periodically persisting job state snapshots to HDFS/S3. On job failure, recovery from the latest Checkpoint achieves Exactly-Once semantics
- State — Flink provides rich state backends (RocksDB, Heap), supporting Keyed State and Operator State, forming the foundation for complex stream processing logic
Choosing a watermark is a balancing act. Too small an allowance drops late data; too large an allowance increases latency. Start with the actual degree of out-of-order arrival and find the right value through experiments.
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:
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:
- Source Side — Kafka Source records current consumption offset during Checkpoint, resuming from the Checkpoint offset on recovery
- Operator Side — Flink's state backend (e.g., RocksDB) persists all state to distributed storage during Checkpoint
- Sink Side — Uses Kafka Sink's transactional write mode (KafkaProducer + transactions), only committing messages after Checkpoint completion
Configuration essentials:
- Set
CheckpointingMode.EXACTLY_ONCE - Kafka Sink uses
DeliveryGuarantee.EXACTLY_ONCE - Kafka Consumer sets
isolation.level=read_committed - Checkpoint interval shouldn't be too long (recommended 30s-60s), otherwise recovery time will be excessive
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:
- ValueState — Stores a single value, suitable for counters, accumulators, and similar scenarios
- ListState — Stores element lists, suitable for caching the most recent N records
- MapState — Stores Key-Value mappings, suitable for deduplication, joins, and similar scenarios
- ReducingState / AggregatingState — Supports incremental aggregation, avoiding traversal of all elements on each window trigger
State backend selection:
- HeapStateBackend — State stored in JVM heap memory, suitable for small state volumes (under GB level), with fast read/write
- RocksDBStateBackend — State stored on local disk (RocksDB), supports incremental Checkpoint, suitable for large state volumes (TB level) in production
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:
- Parallelism Settings — Flink operator parallelism should align with Kafka partition count. Source parallelism = Kafka partition count; operator parallelism adjusted based on CPU intensity
- Network Buffers — Increase
taskmanager.memory.network.fraction(default 0.1) to improve Shuffle efficiency; recommended value is 0.2 - Backpressure Handling — Use Flink Web UI's backpressure monitoring to locate bottleneck operators. Common causes: slow downstream writes, oversized windows, improper state backend configuration
- Kafka Optimization — Increase
fetch.min.bytesandfetch.max.wait.msto reduce small-batch fetch network overhead
Before and after tuning comparison:
- Throughput — From 50K events/s to 200K events/s
- End-to-End Latency — P99 from 2.5s to 800ms
- Checkpoint Duration — From 45s to 3-5s (incremental Checkpoint)
- Resource Utilization — CPU utilization from 40% to 70%
Summary
The Flink + Kafka combination provides an industrial-grade solution for real-time data processing. Core advantages:
- True Stream Processing — Record-by-record processing instead of micro-batching, with lower latency
- Exactly-Once Semantics — End-to-end data accuracy guarantee
- Rich State Management — Supports complex stateful computation logic
- Strong Fault Tolerance — Checkpoint mechanism ensures job recoverability
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.