BACK
Big Data · Flink · Kafka

DataFlow Pipeline

基于 Apache Flink 与 Kafka 构建的实时数据处理流水线,支持高吞吐量的流式数据计算与可视化监控,是大数据课程的核心实战项目。

Role 核心开发者
Duration 3 months
Team 4 人
DataFlow Pipeline

实时流式数据处理平台

DataFlow Pipeline 是一个基于 Apache Flink + Kafka 的实时数据处理流水线平台。 用户可以通过 Vue 前端界面配置数据处理任务,系统将任务编排为 DAG 图形, 由 Flink 引擎执行流式计算,Kafka 作为消息中间件实现数据的高效流转。

平台支持多种数据源接入(MySQL、Redis、Kafka Topic),内置常用的数据清洗、 聚合、窗口计算等算子。后端采用 Java 开发,提供 RESTful API 供前端调用, 配合 Redis 缓存实现任务状态的实时查询。项目作为大数据课程的核心实践, 成功处理了日均百万级的模拟数据流。

技术架构

流式计算 + 可视化监控的全栈大数据方案

Java
Python
Apache Flink
Apache Kafka
MySQL
Redis
Docker
Vue

核心能力

🔄
流式计算引擎
基于 Flink 的高性能流处理引擎,支持窗口计算、状态管理、Exactly-Once 语义保证
📊
可视化监控
Vue 前端实时展示任务运行状态、吞吐量、延迟等关键指标,支持任务的启停与配置管理
🛡️
容错与恢复
基于 Flink Checkpoint 机制实现任务故障自动恢复,配合 Kafka 的持久化消息保证数据不丢失
📈
多维数据分析
内置聚合、过滤、Join、窗口等多种数据处理算子,支持自定义 UDF 扩展处理逻辑

性能指标

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

Flink 流处理作业

核心 Flink 作业定义,实现 Kafka 数据流的实时处理

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");
}
}