Flink 1.20 LTS 与 Spark 3.5 流处理对比:3个核心场景下的延迟与吞吐实测
Flink 1.20 LTS 与 Spark 3.5 流处理对比:3个核心场景下的延迟与吞吐实测
在实时数据处理领域,Flink 和 Spark Streaming 一直是两大主流框架。随着 Flink 1.20 LTS 版本的发布和 Spark 3.5 的更新,两者在流处理能力上都有了显著提升。本文将通过三个典型业务场景的实测数据,对比分析这两个框架在延迟和吞吐量方面的表现,为技术选型提供数据支撑。
1. 测试环境与方法论
1.1 测试环境配置
我们搭建了一个包含 10 个节点的测试集群,每个节点配置如下:
| 组件 | 规格配置 |
|---|---|
| CPU | Intel Xeon Platinum 8480C |
| 内存 | 256GB DDR5 |
| 存储 | 2TB NVMe SSD |
| 网络 | 25Gbps RDMA |
| 操作系统 | Ubuntu 22.04 LTS |
软件版本:
- Flink 1.20.0 (Standalone 模式)
- Spark 3.5.0 (YARN 模式)
- Kafka 3.6.0 (作为数据源)
- Zookeeper 3.8.3
1.2 测试方法论
我们设计了以下测试流程:
- 数据生成器 :使用 Kafka Producer 以可控速率发送测试数据
- 处理程序 :分别实现 Flink 和 Spark Streaming 的处理逻辑
- 结果收集 :通过 Prometheus + Grafana 监控系统指标
- 压力测试 :逐步增加输入速率直到系统饱和
关键指标定义:
- 延迟 :从数据进入 Kafka 到处理完成的时间差
- 吞吐量 :每秒成功处理的消息数量
- 资源利用率 :CPU、内存、网络的使用率
2. 窗口聚合场景对比
窗口聚合是流处理中最常见的操作之一,我们测试了滚动窗口(Tumbling Window)下的表现。
2.1 测试代码实现
Flink 实现 :
DataStream<Event> events = env
.addSource(new FlinkKafkaConsumer<>("input-topic", new EventSchema(), props))
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.aggregate(new CountAggregate())
.addSink(new KafkaSink<>("output-topic", new CountResultSchema(), props));
Spark 实现 :
val stream = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", brokers)
.option("subscribe", "input-topic")
.load()
.select(from_json($"value".cast("string"), schema).as("event"))
.groupBy(window($"event.timestamp", "10 seconds"), $"event.userId")
.count()
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", brokers)
.option("topic", "output-topic")
.start()
2.2 性能测试结果
| 指标 | Flink 1.20 | Spark 3.5 | 差异 |
|---|---|---|---|
| 平均延迟(ms) | 23 | 152 | -85% |
| 最大吞吐(万/s) | 78 | 65 | +20% |
| 99分位延迟(ms) | 45 | 210 | -79% |
| CPU利用率 | 68% | 82% | -17% |
关键发现 :
- Flink 的延迟表现显著优于 Spark,特别是在高负载情况下
- Spark 的微批处理模型在吞吐量上接近 Flink,但资源消耗更高
- 当窗口大小增加到 1 分钟时,两者差距缩小
3. 状态更新场景对比
状态管理是复杂流处理应用的核心需求,我们测试了带有状态计算的用户行为分析场景。
3.1 状态处理实现差异
Flink 状态后端配置 :
state.backend: rocksdb
state.checkpoints.dir: hdfs:///flink/checkpoints
state.savepoints.dir: hdfs:///flink/savepoints
state.backend.rocksdb.ttl.compaction.filter.enabled: true
Spark 状态存储 :
spark.conf.set("spark.sql.streaming.checkpointLocation", "/spark/checkpoints")
spark.conf.set("spark.sql.streaming.stateStore.providerClass",
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
3.2 性能对比数据
测试场景:维护每个用户的最近 10 次行为记录
| 状态规模 | Flink P99延迟 | Spark P99延迟 | 状态恢复时间 |
|---|---|---|---|
| 10万用户 | 28ms | 105ms | 12s vs 45s |
| 100万用户 | 52ms | 230ms | 25s vs 98s |
| 1000万用户 | 115ms | 420ms | 48s vs 210s |
典型问题记录 :
- Spark 在状态超过 500MB 后出现明显的 GC 停顿
- Flink 的增量检查点机制在大型状态下优势明显
- 测试中 Spark 出现了 3 次状态丢失需要手动恢复的情况
4. 复杂事件处理(CEP)场景
CEP 用于检测数据流中的复杂模式,是风控和监控系统的核心需求。
4.1 测试模式设计
我们实现了一个信用卡欺诈检测规则:
模式序列:
1. 同一卡号在10分钟内
2. 出现3笔以上交易
3. 交易金额均超过信用额度的80%
4. 且交易地点跨度过大
Flink CEP 实现 :
Pattern<Transaction, ?> pattern = Pattern.<Transaction>begin("start")
.where(new SimpleCondition<>() {
@Override
public boolean filter(Transaction value) {
return value.getAmount() > creditLimit * 0.8;
}
})
.timesOrMore(3)
.within(Time.minutes(10));
Spark 实现 : 由于 Spark 原生不支持 CEP,我们使用连续多个 SQL 查询模拟:
-- 第一步:筛选大额交易
CREATE TEMP VIEW large_trans AS
SELECT * FROM transactions WHERE amount > credit_limit * 0.8;
-- 第二步:窗口内计数
SELECT card_no, COUNT(*)
FROM large_trans
GROUP BY card_no, window(timestamp, '10 minutes')
HAVING COUNT(*) >= 3;
4.2 性能对比
| 复杂度 | Flink 吞吐量 | Spark 吞吐量 | 延迟差异 |
|---|---|---|---|
| 简单模式(3事件) | 45万/s | 28万/s | 1.6倍 |
| 复杂模式(5事件) | 32万/s | 14万/s | 2.3倍 |
| 超复杂模式 | 18万/s | 4.5万/s | 4倍 |
模式检测准确性 :
- Flink 准确检测到所有测试用例
- Spark 在事件乱序超过 30 秒时出现漏报
- 对于跨多个窗口的模式,Spark 实现复杂度呈指数增长
5. 选型决策指南
基于实测数据,我们总结出以下决策树:
是否需要亚秒级延迟?
├── 是 → 选择 Flink
└── 否:
├── 是否需要复杂状态管理?
│ ├── 是 → 选择 Flink
│ └── 否:
│ ├── 是否已有Spark集群?
│ │ ├── 是 → 考虑 Spark
│ │ └── 否 → 选择 Flink
└── 主要处理批分析?
├── 是 → 选择 Spark
└── 否 → 选择 Flink
特殊场景建议 :
- IoT 设备监控 :Flink 的低延迟特性更适合
- 夜间报表生成 :Spark 的批处理模式更经济
- 金融风控系统 :必须选择 Flink 以保证实时性
- 数据湖 ETL :两者均可,取决于现有技术栈
在实际项目中,我们曾遇到一个电商大促场景,Flink 在流量突增 10 倍时仍保持稳定,而 Spark 需要动态调整批处理间隔才能跟上数据速率。这也印证了 Flink 在真正实时场景中的优势。
更多推荐
所有评论(0)