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 测试方法论

我们设计了以下测试流程:

  1. 数据生成器 :使用 Kafka Producer 以可控速率发送测试数据
  2. 处理程序 :分别实现 Flink 和 Spark Streaming 的处理逻辑
  3. 结果收集 :通过 Prometheus + Grafana 监控系统指标
  4. 压力测试 :逐步增加输入速率直到系统饱和

关键指标定义:

  • 延迟 :从数据进入 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 在真正实时场景中的优势。

更多推荐