资讯动态

Spark实时模式vs Flink流处理:延迟对比与迁移实战

发布时间:2026/10/9 8:39:24 来源:尧图企业网站定制
1. 先给结论Spark 的实时到底是不是伪命题前两年有个朋友问我能不能把 Spark 的实时模式用起来把 Flink 从实时王座上拉下来当时社区里流行两种声音一种说 Spark Structured Streaming 已经出了 Continuous Processing 模式正式进入实时领域了另一种说 Flink 的 Checkpoint 一致性在架构上就领先Spark 再怎么改都追不上。这两种说法其实都只说对了一半真正判断还是要落到你自己的业务场景和延迟需求上。先说我的结论。对于输入一条日志、几十毫秒内就要看到结果这种强实时场景Apache Spark 的 Real-Time Mode 目前并不具备生产级竞争力它的 Continuous Processing 至今仍是实验特性支持的操作极其有限硬上只会自找麻烦。但如果你能接受 310 秒级别的端到端延迟又希望吞吐拉满、能用熟悉的 Spark 生态做批流一体处理那微批次架构完全够用甚至比 Flink 更适合你。这篇文章我会把为什么讲透微批次到底卡在哪、Flink 的毫秒级延迟来源是什么、两边的一致性模型有什么区别。后面还会结合 Flink JDBC 连接器异常、Spring Boot 整合 Flink 这两个高频踩坑点把从 Spark 迁移到 Flink 路上最容易被绊倒的地方一一拆开。2. Real-Time Mode 的真相微批次架构到底卡在哪2.1 微批次是怎么工作的理解了调度模型就理解了延迟来源Spark 的实时模式准确说是基于微批次的近实时模式时间粒度的下限就是一个批次间隔。Spark StreamingDStream时代的做法很直接设一个批次间隔比如 3 秒数据流被切成一段段 RDD每个 RDD 交给 DAG 调度器执行。你可以想象一个收费站车明明是一辆接一辆开过来的但收费站必须凑满 3 秒的车辆才放行一次所以单辆车过站的时间下限就是等满一个批次窗口。到了 Structured StreamingAPI 层升级成了 DataFrame 的增量查询支持 Event Time、Watermark、窗口代码写起来跟批处理几乎一样但执行层还是那套微批次引擎。你写trigger(ProcessingTime(3 seconds))就是 3 秒一个批次写trigger(Once())就是跑一次开一个批次本质是变相批处理。真正接近实时的是 Spark 2.3 引入的trigger(Continuous(1 second))。它的思路是让 source 与 sink 直连一组持续运行的 executor 任务数据不切批次就能被消费任务心跳以 epoch 为单位推进。问题在于这个模式到今天都还标着 experimental只支持select、where、map、filter这类无状态操作聚合、Watermark、窗口、Sort 统统不支持。换个直白说法几乎所有真实世界的实时计算需求按用户维度统计 PV/UV、风控特征提取、异常事件流式检测这个模式全都做不了。微批次的延迟来源可以归结为三点。第一是任务调度开销每个批次都要重新生成 Job、重新分配 Task哪怕这个批次只有几百条数据调度成本也一分不少第二是批次边界等待上游必须等整个批次的数据到齐才能进入下游算子第三是 Shuffle 持久化的代价跨节点 Shuffle 要落盘合并这部分在低延迟路线里是纯负担。这三笔开销叠在一起端到端延迟就很难压进秒级以下。2.2 实测视角Spark Structured Streaming 的延迟表现我在自己的测试环境里3 台 worker、每台 8 vCPU/16GKafka 作为 source下游是控制台和 MySQL sink跑过一轮对比。批次间隔设 1 秒时从 Kafka 消息被生产到 MySQL 出现结果平均端到端延迟在 2.5 秒上下P95 能到 5 秒左右高峰期还会冲到 8 秒把批次间隔放宽到 5 秒P95 基本稳定在 810 秒。这个数字受集群资源、数据量、sink 写库速度影响很大我给的是常见范围不是权威基准。值得注意的细节是端到端延迟并不是简单地等于批次间隔 处理时间。Spark 的宏观调度存在批次排队如果某个批次因为数据倾斜或 GC 卡住后续批次会被整体堵住。延迟的波动性其实是微批次相对难解决的问题因为它天然自带积压缓冲区。你调小批次间隔到 500 毫秒调度开销占比反而会升高导致吞吐明显下降这就是微批次在延迟和吞吐之间的结构性矛盾。2.3 打破微批次壁垒目前只打了一个缺口回到标题里那个很有感染力的说法——打破微批次壁垒。从技术演进的视角看Continuous Processing 确实把执行模型往前推了一步证明 Spark 团队想让引擎具备真正的流式执行能力。但从工程落地视角看这个缺口非常窄无状态操作能做有状态计算做不了简单 ETL 能做复杂窗口和聚合做不了任务少的时候能跑任务一多、状态一多立刻打回原形。而且有一个容易被忽略的成本连续模式牺牲了微批次的吞吐优势。批量化是 Spark 高吞吐的根本原因凑批之后可以做向量化执行、批量写出、更优的压缩。一旦改成逐条处理单条数据都要走完整的算子链路吞吐能力反而下降。这等于用自己最强的吞吐优势去换一个还没补齐的延迟能力性价比很低。3. Flink 的实时护城河毫秒级延迟从哪来3.1 流式计算的核心数据到了就处理而不是凑批再处理Flink 的整个执行模型就是为流而生。算子之间通过有界缓冲bounded buffer直接衔接数据从 source 读出来立刻进入下游算子中间没有等一个批次凑齐的环节。它可以做到每条记录独立驱动计算这就是它和 Spark 微批次最本质的分歧点——一个按条处理一个按批处理。很多人以为 Flink 快只是因为它不等批其实还有一个关键设计是流水线执行。上游算子处理完的数据可以立即传给下游下游不用等上游全部完成。这跟 CPU 流水线很像指令不需要全部执行完才出结果而是每个阶段只要拿到自己需要的数据就往下传。延迟自然就被压缩到毫秒级。Flink 的反压机制也值得一提。它的 Credit 协议让下游有足够能力接收数据时通过信用额度通知上游上游才继续发送数据不会在算子之间无限堆积。相比之下Spark 微批次对背压的感知是滞后的积压往往通过批次处理时间变长来体现等你发现的时候堆积已经发生了。3.2 端到端 Exactly-Once 的路线差异快照 vs 重放两边都宣称支持 Exactly-Once但实现代价完全不同。Spark 能保证 Exactly-Once主要是因为批次天然是事务边界一批数据要么完整处理完要么失败后整批重放source 端配合 Kafka offset 就可以做到不丢不重。这个思路非常批处理重放的粒度也是整个批次。Flink 采用异步分布式快照Chandy-Lamport 算法的变体。它定期在数据流中插入 BarrierBarrier 流经每个算子时把该算子的状态做一次快照全部算子快照完成即生成一个 Global Checkpoint。如果出现故障就回滚到最近一次完成的快照并重放那之后的输入数据。关键差异在于Flink 可以做部分重放只重放快照之后的数据恢复粒度细、恢复时间短Spark 恢复的是整个批次进度批次越大恢复代价越大。这带来的直接体验差异是Flink 在保证 Exactly-Once 的同时能把延迟保持在低位而 Spark 的批次模型天然要在一致性和延迟之间做一个更明显的权衡。当然Flink 的 Checkpoint 也不是零成本它需要存储状态快照状态大了恢复也会变慢需要配合增量 Checkpoint、RocksDB 状态后端来优化。3.3 状态管理与时间语义Flink 比 Spark 多走了几步实时计算离不开状态。Flink 的状态是显式的一等公民Keyed State 提供了 ValueState、ListState、MapState可以通过ValueStateDescriptor声明还能设置 TTL。状态后端可以选择内存或 RocksDB可以支撑 GB 级甚至 TB 级的状态数据。Spark 的微批次模型并不直接暴露状态概念Structured Streaming 里要做有状态计算本质上是用 StateStore 保存中间聚合结果表由 Watermark 清理过期状态。这种设计对简单聚合够用但要实现复杂的会话管理、模式匹配、或者精细化状态生命周期控制就会很别扭。时间语义上Flink 原生了 Event Time Watermark 三种窗口滚动、滑动、会话allowedLateness允许迟到数据在窗口关闭后的一段时间内继续更新结果。Structured Streaming 虽然也支持 Event Time但窗口逻辑相对固定连续模式下更是连 Watermark 都用不了。如果你的业务强依赖乱序事件处理和复杂窗口Flink 是明显更顺的选择。3.4 两张引擎关键能力对照下面这个表格是我做选型时常用的对照维度给出了我心目中两者各自的定位| 维度 | Spark Structured Streaming | Flink 流处理 | |-------------------|---------------------------------------|---------------------------------| | 执行模型 | 微批次默认、连续处理实验 | 真正的流式执行 | | 端到端延迟 | 秒级通常 3~10 秒 | 毫秒级通常数百毫秒以内 | | 吞吐 | 高批量化友好 | 高但单条链路开销略大 | | Exactly-Once | 靠批次重放实现 | 靠异步分布式快照实现 | | 状态管理 | StateStore 中间聚合表够用但受限 | Keyed State、多状态后端、TTL | | 事件时间/窗口 | 支持但窗口类型相对固定 | 原生强支持滚动/滑动/会话窗 | | 运维复杂度 | 跟随 Spark 集群生态熟悉 | 独立 JobManager 体系概念更多 | | 典型使用场景 | 吞吐优先的批流一体、数仓增量 ETL | 延迟敏感、状态复杂、事件驱动计算 |这个表不是用来做二选一的裁决而是帮你在面对具体需求时快速定位延迟容忍度、状态复杂度、团队技术基础这三个条件会直接决定答案。4. 从 Spark 切到 Flink 的实操实录连接器异常与框架整合两大坑4.1 Flink JDBC 连接器异常从驱动类加载到连接泄漏切换引擎时最先撞上的往往不是核心 API而是各种连接器问题。Flink JDBC 连接器最常见的第一类异常是驱动类找不到报错长这样java.lang.ClassNotFoundException: com.mysql.cj.jdbc.Driver。很多人本地调试没问题一打包提交到集群就挂原因是 Flink 的类加载默认是 Parent-First你依赖的 JDBC 驱动没有打进最终的 fat jar。排查思路分三步先检查依赖的scope是不是provided如果父工程里已经声明了子模块打包时很容易丢其次看assembly或者shade配置是否把mysql-connector-j排除掉了最后在运行节点上看一眼 Flink 的lib目录确认驱动是否需要在集群层面统一放置。我建议在提交命令里加上-C file:///path/to/mysql-connector.jar临时验证驱动问题解决后再固定到构建配置里。第二类高发问题是连接数被打爆。Flink JDBC sink 每个并行子任务会维护自己的数据库连接并行度设为 32就是 32 个连接常驻如果下游数据库连接池上限只有 20直接报连接超时。一个务实技巧是降低 sink 并行度或者改用批量重写提交的方式控制写入频率。另一个非常隐蔽的坑是时区参数MySQL 高版本连接串如果不配serverTimezoneAsia/Shanghai会出现日期偏移导致的脏数据这类错误在日志里不容易直接看出来表现为数据错乱而不是连接失败。第三类是事务型连接器的记录锁问题。用 JDBC sink 做 Upsert 时如果没有合理的唯一键设计批量的INSERT ... ON DUPLICATE KEY UPDATE在并发场景下容易产生死锁。我在实践中的做法是把实时明细落到一个单独的写入缓冲表再通过定时任务批量合并到业务表避免实时链路直接和在线事务抢锁。纯粹追求低延迟时表的写入模型要单独设计不能照搬离线数仓的覆盖写思路。4.2 Spring Boot 整合 Flink不要掉进 fat jar 依赖的坑第二个高频话题是 Spring Boot 整合 Flink。很多人习惯把所有事情放在一个 Spring Boot 工程里在SpringBootApplication启动类里直接new StreamExecutionEnvironment然后execute()本地跑通了一部署就出问题。最典型的是依赖冲突。Spring Boot 2.x 自带 Jackson 2.x、Logback、Guava 28而 Flink 1.17 以后的运行时用 Guava 30两边在 fat jar 里撞车之后会出现各种莫名其妙的NoSuchMethodError或序列化异常。Netty 和 Jackson 的版本冲突也会在提交到 Flink 集群时爆发。这个问题的根因是 Flink 集群用 Parent-First 类加载策略加载了 Spring Boot fat jar 里的依赖让 Flink 自己内部的类版本被覆盖。正确做法不是去改 Flink 的类加载顺序改完后面还有更多麻烦而是在工程结构上做隔离。我采用过的方案是拆两个模块job-module只依赖 Flink API 和业务逻辑不引入 Spring Bootapp-module负责 Spring Boot 启动通过 REST 或配置文件把运行参数传给 job 模块Job 模块打包成可提交的 fat jar。Job 里要用到配置项时优先用 Flink 的Configuration和ParameterTool传参而不是依赖 Spring 的Value这样既保留 Spring Boot 的工程化能力又避免框架层面的污染。还有一个容易忽略的点不要在 Spring Boot 启动时同步阻塞执行 Flink 的execute()。Job 提交后任务运行在独立集群上你本地的 Spring 容器退出并不会影响集群里的任务但同步execute()会让启动接口一直挂着造成任务看起来没起来的假象。正确姿势是提交后立即返回任务 ID再通过 Flink REST API 查询状态。4.3 从微批次思维迁移到流式思维的三个关键转变从 Spark 过来的人代码习惯需要切三个视角。第一个是 Offset 管理视角。Spark Structured Streaming 的 Kafka source 会在批次结束时提交 offset你可以在代码里看到清晰的控制点Flink 的 Kafka offset 由 Checkpoint 统一管理开启 Checkpoint 后消费位点跟随快照滚动。如果你习惯手动提交 offset这个动作在 Flink 里是多余的反而可能干扰一致性。关掉 Checkpoint 后 Flink 的位点提交行为又会变成至少一次这个差异必须提前想清楚。第二个是窗口计算视角。Spark 的窗口通过withWatermark(ts, 10 seconds).groupBy(user, window(ts, 1 minute))这类 API 隐式完成Flink 则要求你显式调用keyBy(...).window(TumblingEventTimeWindows.of(Time.minutes(1)))还要自己指定WatermarkStrategy。很多人刚上手时把 Spark 那段代码照搬过来结果发现事件时间没生效或者窗口一直不触发多半就是 Watermark 没配对。第三个是状态管理视角。Spark 里你不太会感知状态存储这个概念中间结果自动由 StateStore 管理Flink 里每个有状态算子都要声明状态描述符和状态后端状态 TTL、状态清理策略都需要自己规划。从 Spark 切过来的人前期最不适应的就是状态竟然是要买声明的。但正是这种显式化让 Flink 能够在超大状态上保持高效的访问和恢复。5. 选型建议与踩坑后的个人经验5.1 什么情况继续用 Spark什么情况果断上 Flink如果你所在的团队已经深度绑定 Spark 生态数据链路是离线批处理加实时增量 ETL下游是数仓或 BI 报表延迟容忍度在 10 秒以上那 Structured Streaming 是性价比最高的选择。它不需要额外维护一套独立流计算集群同一套 Spark 引擎既能跑凌晨的全量表加工也能跑白天的增量任务运维心智负担小。还有一点别忽略Spark 对 Hive 血缘、Iceberg/Delta Lake 这类湖格式的整合非常顺批流一体这条路有天然优势。如果业务要求秒级甚至毫秒级响应比如实时风控、实时推荐特征、在线监控告警或者需要复杂的会话窗口、事件模式匹配、高并发状态读写就果断上 Flink。这类场景里微批次的延迟抖动和批次积压风险会直接把业务效果拖垮。Flink 的运维复杂度的确更高但相比业务受损的代价这个成本值得付。5.2 我踩过的选型失误希望你别重复说一个我自己的真实教训。之前做过一个实时大屏项目业务方最初说延迟要求是1 分钟内可用我心想这个要求 Spark 微批次完全够用就选了 Structured Streaming。结果上线后业务方看到其他竞品做到了秒级刷新立刻把需求改成了3 秒内必须看到数据。微批次在那个延迟条件下面临巨大压力后来整改成 Flink Kafka Redis 架构重写了窗口逻辑和状态管理才把需求兜住。这次经历给我的反思是选型时要看的不是当前需求而是需求在演进的路上会不会往更低延迟的方向走。另外一个小建议如果你想在团队内做技术预研、但不想立刻上 Flink可以先拿同一个业务场景在两个引擎上各跑一版用文档记录 P95 延迟、吞吐、峰值积压、恢复时长这几个指标。这类一手数据比任何博客文章都有说服力也能帮你提前发现运维上要补的功课。5.3 一个实在的收尾建议从 Spark 微批次到 Flink 流式不是一次简单的框架切换而是对实时这两个字的理解升级。Spark 把实时处理变成了一系列快速执行的迷你批次Flink 则真正把计算变成了对每条数据的连续加工。选择哪一边取决于你要的产品体验是接近实时还是真正实时。我个人现在的习惯是把延迟要求拆成两个数字端到端 P95 容忍上限和峰值积压容忍量拿这两个数字去做引擎选型比听任何颠覆、王座这类技术叙事都可靠得多。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑