BACK
Big Data · Flink · Kafka

DataFlow Pipeline

Real-time data processing pipeline built on Apache Flink and Kafka, supporting high-throughput stream computing and visual monitoring — a core hands-on project for big data coursework.

Role Core Developer
Duration 3 months
Team 4 People
DataFlow Pipeline

Real-Time Streaming Data Platform

DataFlow Pipeline is a real-time data processing platform built on Apache Flink + Kafka. Users configure data processing tasks through a Vue frontend, which are orchestrated into DAG graphs and executed by the Flink engine, with Kafka as the message middleware for efficient data flow.

The platform supports multiple data source integrations (MySQL, Redis, Kafka Topics) with built-in operators for data cleaning, aggregation, and window computation. The Java backend provides RESTful APIs for the frontend, with Redis caching for real-time task status queries. As a core big data course project, it successfully processed million-level daily simulated data streams.

Architecture

Full-stack big data solution with stream computing + visual monitoring

Java
Python
Apache Flink
Apache Kafka
MySQL
Redis
Docker
Vue

Core Capabilities

🔄
Stream Processing Engine
High-performance stream processing engine based on Flink, supporting window computation, state management, and Exactly-Once semantic guarantees
📊
Visual Monitoring
Vue frontend displaying real-time task status, throughput, and latency metrics with task start/stop and configuration management
🛡️
Fault Tolerance and Recovery
Automatic task failure recovery based on Flink Checkpoint mechanism, with Kafka persistent messaging ensuring zero data loss
📈
Multidimensional Analytics
Built-in aggregation, filtering, Join, window and other data processing operators with custom UDF extension support

Performance Metrics

500K+
TPS
99.9%
Availability
100ms
Latency
15+
Jobs

Flink Streaming Job

Core Flink job definition for real-time Kafka data stream processing

StreamJob.java
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
// Define the Flink streaming job
public class StreamJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
// Kafka source with auto-commit
FlinkKafkaConsumer<String> consumer =
new FlinkKafkaConsumer<>(
"sensor-data",
new SimpleStringSchema(),
kafkaProps
);
env.addSource(consumer)
.map(new DataParser())
.keyBy(DataEvent::getDeviceId)
.window(TumblingProcessingTimeWindows.of(
Time.seconds(10)))
.aggregate(new AvgAggregator())
.addSink(new MySQLSink());
env.execute("DataFlow Pipeline");
}
}