尧图精选

Flink还是Spark?实时计算引擎选型的关键差异与实战逻辑

🕒 发布时间:2026/9/10 13:01:43 📁 来源:尧图网络
前不久跟一个做数据平台的朋友聊天他正在给团队选实时计算引擎从“Spark还是Flink”一路聊到凌晨。他说自己不是不知道两者的区别而是网上对比文章实在太多有的吹Flink真流式天下无敌有的说Spark生态成熟够用就行看完更不知道怎么选了。这种纠结我太理解了。大数据流处理框架选了快十年选型文档攒了一抽屉但真正落地之后你会发现纸面上的性能参数和线上环境的真实表现完全是两回事。这篇文章不打算重复那些官网文档和PPT对比就以我这几年用Flink和Spark做实时数仓、处理Kafka数据流、搞CDC同步的实际经验把两者的核心差异、工程落地细节和选型逻辑一次聊透。不论你是正准备入门流处理、在做技术选型还是单纯想搞明白这两个框架到底什么关系这篇都能给你一个相对完整、可落地的参考框架。1. 先别急着比框架搞懂批处理和流处理的本质差别聊Flink和Spark之前得先把“批处理”和“流处理”这两个概念理清楚。很多人选型选不明白根源就是没搞懂自己的业务到底属于哪一类然后用流处理的框架去做批处理的活或者反过来最后怎么调都不对劲。批处理大家很熟悉了数据落盘后一次性拿过来算。比如凌晨跑昨天的日志分析一个Spark任务读HDFS上的几TB文件跑完输出结果。“批”这个字的核心含义是数据是完整的、有界的处理引擎可以看完全部数据再动手。既然数据都在那儿放着引擎就可以做各种全局优化——谓词下推、列剪裁、动态分区修剪怎么高效怎么来因为等得起。流处理则完全不同。数据像水流一样源源不断涌进来没有明确的终点今天早上九点的数据不会因为引擎没处理完就停下等你。流处理面对的是无界数据它的核心挑战不是“怎么把全量数据算得更快”而是“数据一直在来我怎么保证每一条都处理了、不丢不重、还按正确的顺序输出”。这里有个常见误区很多人以为流处理就是“批处理分得更细一点”Spark Streaming最初就是这么干的。它把数据切成一段段的微批次比如每两秒攒一批然后用批处理的方式算。这种思路的好处是能复用Spark成熟的批处理能力但代价就是延迟被“攒批”这个动作卡住了——永远存在一个批间隔的延迟而且吞吐越高、批次越大延迟越明显。这就像食堂打饭批处理是等菜全上齐了大家一起开饭微批次是每两秒开一桌而Flink这种真流式是来一个客人就上一个菜的流水席菜到即上不等任何人。搞清楚这个底层差异你才能理解为什么同一个业务用两个框架跑出来效果天差地别——不是谁“强”谁“弱”而是它们的计算模型从根上就不是一回事。再往深一层看流处理和批处理对“正确性”的定义也不一样。批处理跑完了发现结果不对把输入重新读一遍重跑就能修但流处理的数据是实时涌来的历史数据可能已经过期、窗口可能已经关闭出了错想对账、回溯、修正的代价要高得多。所以Flink才会把状态管理和Checkpoint当成头等大事来设计——因为它天生就要处理那些“不但算得快还要算得准”的场景。Spark这些年也在往流处理上补课从Spark Streaming的纯微批次到Structured Streaming引入Event Time和Watermark再到2.3版本推出的Continuous Processing模式想在保持批处理优势的同时把延迟打下来。但工程上没有银弹Spark始终是个以批处理为内核的引擎穿上流处理的衣服骨子里的调度模型、内存布局、容错机制还是批处理那一套。理解了这一点后面看技术细节就都通顺了。2. 计算模型的分水岭Flink的流水线式数据流 vs Spark的微批次调度这是Flink和Spark最核心、也最容易被忽略的差异。表面上看两者都能消费Kafka、都能做窗口聚合但数据在引擎内部怎么流动、怎么计算走的是完全不同的两条路。Flink的核心是一个真正的流式数据流引擎。用户写好的逻辑会被编译成一个DAG有向无环图无数条事件记录从这个DAG的源头Source流入经过一个个算子Operator的处理最终从输出Sink流出。关键点在于数据是一条一条在算子之间传递的上游算子的计算结果可以立刻交给下游算子继续算不需要等一个批次攒齐。数据到达、计算、转发这整套动作是以“条”为单位在流水线上并行的。打个比方Flink像一条汽车装配线每个工位只处理自己手头的一辆车车一辆接一辆顺着流水线走整个工厂不存在“把所有车都停下来说大家攒够一批再走”这种操作。Spark就不一样了。Spark Streaming和Structured Streaming默认都是微批次micro-batch模式引擎固定每隔一段时间比如2秒处理一次每个批次内的数据被看成一个小的RDD或DataFrame然后走Spark经典的调度流程——生成作业、切分Stage、Shuffle、执行、产出结果跑完这批再跑下一批。在这种模型下数据流动是“一阵一阵”的仿佛水龙头被装了一个隔膜泵每隔两秒才喷出一股水。因为是批处理引擎每次微批次都需要经历任务调度、序列化、网络洗牌这些开销这决定了Spark的流处理延迟天花板比较难打到秒级以下。还有个很有意思的细节Flink的流水线式执行天然支持算子之间的数据“流式shuffle”也就是说上游算子计算结果边产生边发送给下游整体上的端到端延迟就是一条数据从Source到Sink实际走过的网络和计算时间而Spark的微批次模式里数据先全量落入这个批次的临时存储相当于攒着再进入下一阶段的处理。数据在中间多了一次落地和等待这个动作本身就是延迟的根源。不过微批次模型也不是一无是处。批是一个天然的缓冲遇到数据波动时吞吐可以拉得很稳系统背压的表现也更接近批处理的行事方式不会因为单条数据处理不过来就把整条链路卡死。这是Spark团队当初选这个模型的聪明之处用可接受的延迟换来了高吞吐和工程上的简单可靠。而Flink的高延迟性能建立在一整套复杂的反压机制背压传播和持续调度模型之上工程复杂度和调优门槛明显更高。Structured Streaming后来加了一个Continuous Processing模式号称能做到毫秒级延迟但官方文档都写得很明白它不支持某些操作比如聚合、水印处理的限制很多生产环境真正敢用的极少绝大多数Spark用户跑流处理依然跑的是微批次。所以不要被“支持连续处理”这种字眼迷惑生产环境的真实表现才是唯一标准。说到实际项目里的感受我举个具体例子。做一个实时风控场景业务方要求从交易事件发生到风控规则命中返回结果延迟不能超过500毫秒。用Spark Streaming默认2秒一个批次这需求基本上无法做到延迟会稳定停留在2秒以上这还只是理想情况下的计算延迟。换了Flink之后端到端延迟直接降到200毫秒以内因为中间没有攒批等待的环节。反过来如果只是做一个五分钟维度的实时大屏指标PV、UV、订单量Spark微批次完全够用甚至吞吐更稳定调优也更顺手——这时候硬上Flink就是给自己找麻烦。所以你们会发现技术上谁强谁弱是一个问题你的业务到底需要多低的延迟、心里预期是多少是另一个问题而这两个问题必须放在一起回答才有意义。3. 数据正确性的较量状态管理、窗口水印与精确一次的容错之路聊完计算模型得说说真正见功力的部分状态管理、窗口处理和容错机制。这些才是生产环境里最影响一把手的硬功夫也是网上“Flink比Spark好用”这种说法背后的真实支撑。先说状态State。流处理里的很多计算是“记住之前算到哪再和最新数据合并”的累进过程。比如统计每个用户的累计消费金额得先把这个人历史上所有订单金额存下来再叠加新订单。这种需要跨数据保留的中间结果就是状态。Flink对状态的重视程度可以说是将其视为引擎的基础。它把状态分为Keyed State和Operator State由状态后端负责存储默认的HashMapStateBackend把状态放在堆内存/HeyDB里适合状态量不大的场景RocksDBStateBackend则把状态落到RocksDB的本地存储上几乎无限扩容适合百GB甚至TB级别的状态。更关键的是Flink支持增量Checkpoint——每次只上传变化的状态部分大大提升了超大状态下的备份效率。我在实际项目里用RocksDB存过规模比较大的去重用户ID几亿个Key的状态做Checkpoint完全没有压力这对Spark来说是很难想象的事。Spark Structured Streaming的状态存储要朴素得多。它早期主要靠HDFS的Write-Ahead Log预写日志来保存状态每次状态更新都要写一份日志状态太大时会非常吃力。这也是为什么Spark做有状态的流计算比如大窗口累计、长时间会话聚合时一旦状态规模上来就会遇到明显的性能瓶颈和调优地狱。后来Spark 3.x引入了状态存储插件接口允许接入RocksDB之类的实现但相比Flink原生对状态一等的设计成熟度和细粒度管理还是有肉眼可见的差距。然后是窗口和乱序数据。实时数据从源头到引擎网络抖动、生产者积压都会造成数据乱序比如Kafka里消息的写入顺序和事件实际发生顺序可能并不一致。处理乱序数据的经典机制是Watermark水印它在Flink里的设计非常精巧水印表示“在这个时间戳之前的数据都已经到了”引擎遇到水印就知道可以安全地触发窗口计算了同时允许用户配置allowedLateness来控制窗口关闭后还能等多久迟到的数据。Flink的窗口操作是目前流引擎里最完整的滚动窗口、滑动窗口、会话窗口还有自定义窗口配合Process Function可以在窗口结束时拿到该算的全部上下文。而Spark Structured Streaming虽然也支持Watermark和事件时间但它的水印和窗口语义是在微批次基础上实现的批与批之间水印推进策略相对粗粒度窗口计算的类型和灵活性也明显不如Flink丰富。遇到那种需要复杂会话切分、多级窗口叠加的业务用Flink写起来是“顺着思路来”用Spark就得各种绕路。最后说说一段式“精确一次”Exactly-Once这件事。听起来像是数据库事务才能做到的事但这几年流处理引擎确实把这个概念推到了工程前沿。Spark的Structured Streaming在做端到端精确一次时依赖两层引擎内部通过WAL和状态管理保证重启后不重复计算写入外部系统比如Kafka、HDFS时要依赖外部Sink本身支持幂等写或事务性写入来避免结果重复落盘。这套方案不是不能用但对Sink的要求比较高很多时候需要自己实现幂等逻辑。Flink的思路更系统它的Checkpoint机制基于Chandy-Lamport分布式快照算法在数据流里按周期插入Barrier屏障每个算子收到Barrier后把状态做分布式快照同时阻断该通道上的数据处理以此保证所有算子在同一个时间点上冻结。配合两阶段提交——预提交和真正提交——Flink可以实现端到端的精确一次语义Kafka作为Sink时是标准做法这也是它稳坐实时数仓核心位置的原因之一。不过要泼一盆冷水精确一次不是免费的午餐。开启两阶段提交意味着Checkpoint正常完成的时候外部Sink才会收到确认并提交事务如果Checkpoint太频繁外部存储上的未提交临时数据量会很大内部性能开销也会上涨。所以实际工程里大家往往根据场景在“至少一次”和“精确一次”之间灵活切换——不是所有的写入都需要精确一次很多统计分析类业务“至少一次幂等清洗”就完全够用了。我自己的习惯是对账类、金融交易类、金额累加类场景必须精确一次而指标监控、流量看板、推荐日志这种场景用至少一次加幂等键性能好很多还省心。4. 工程落地实录部署、内存模型、调优与Cache的深坑技术能力再强最终要落到工程上跑得稳才行。这部分我讲讲两个框架在真实环境里的部署运维细节和调优经验都是我踩过的坑换来的。4.1 部署模式与集群资源Spark的部署方式相对多元本地模式、Standalone、YARN、Mesos还有Kubernetes都支持生产环境用YARN模式的居多。有个很常见的疑问“Spark on YARN提交是不是只需要一个Spark客户端就行了”答案是对的只要一个安装了Spark客户端包含spark-submit脚本和客户端jar包的机器通过网络连接到YARN集群的ResourceManager就可以提交任务。任务提交后由RM为应用分配一个Container来启动ApplicationMaster再由AM去申请Executor的Container。关键点在于客户端版本最好和集群版本对齐且客户端到集群的网络要通。很多新手卡在这一步往往就是版本号差了一个小版本提交时API还能用但到了调度阶段就开始出各种不可名状的怪问题。Flink的部署模式也在快速发展有YARN模式有Standalone模式而这些年Kubernetes已经成了重点方向尤其是Flink Kubernetes Operator出现后用声明式的方式来管理Flink集群和任务生命周期跟云原生的契合度比Spark高不少。Flink on K8s还有一个好处是资源隔离比YARN更干净容器粒度可控出问题影响面小。但成本是运维复杂度上去了——K8s本身那一套网络、存储、权限模型不是开玩笑的团队没有容器化基础的话老老实实跑YARN反而更稳。4.2 内存模型的差异化理解两个框架的内存模型我都调过风格完全不一样搞不清楚的话线上任务分分钟OOM给你看。Spark的Executor内存是出了名的“多人合租房”总内存分为执行内存Execution Memory用于Shuffle、Join、Aggregation等计算操作、存储内存Storage Memory用于缓存RDD/DataFrame、还有留给用户代码的用户内存和系统开销。Spark 1.6之后用统一内存管理器来管理执行内存和存储内存两者之间的边界不固定一方内存吃紧时可以向另一方借贷。调Spark内存核心要理解这一层一堆人以为只用调整spark.executor.memory就行结果Shuffle内存不够时任务一直Spill到磁盘性能掉的没法看其实更该关注的是spark.memory.fraction和spark.memory.storageFraction的比例。Flink的TaskManager内存管理方式更类似“按房间分租”分为JVM Heap内和Heap外的Framework内存、Task内存、网络缓冲内存等其中Task内存里又区分托管内存Managed Memory用于RocksDB状态后端时尽量给够和普通JVM内存。调Flink内存有个经典误区堆外内存总是不够用。默认情况下网络缓冲内存和托管内存都要占堆外空间如果数据量大、并发高堆外内存设小了会直接报Direct buffer memory或Unable to allocate heap off-heap memory错误。我一般习惯把taskmanager.memory.task.off-heap.size和网络内存提前规划好留足余量而不是等到报错再去加。4.3 一条经典链路的调优记录Flink消费Kafka写入Elasticsearch这是实时业务里出现频次最高的链路之一我拿一次实际调优来举例。最初版本是Flink从Kafka读取业务日志经过ETL清洗、按某个维度聚合然后通过Elasticsearch Connector写入ES。上线后观察到一个很典型的毛病ES索引写入延迟高并且Flink侧反压Backpressure持续飘红Kafka消费Lag越拉越大。刚开始以为ES集群能力不够扩容了之后效果还是不理想后来排查才发现根源在Flink作业的并行度设置和批量写入参数上Flink写入ES时如果每个并行实例的批量写参数没调默认对内还是单条请求模式反而没有发挥ES批量接口的优势。后来把bulk.flush.max.actions调到1000、bulk.flush.max.size.mb调到5MB、bulk.flush.interval.ms调到2000ms之后写入效率立刻翻倍。其次是并行度没对齐Kafka Topic有12个分区Flink Source并行度设了12但Sink并行度还是默认的1等于12个Source通道的数据最后都压到一个线程写ES怎么可能不堵。把Sink并行度调到和Source一致后反压肉眼可见地降了下来。最后是索引模板和mapping的坑ES那边某个字段采用了非常重的自定义分词器写入时CPU飙高后来把不需要分词的字段改成keyword类型、调整refresh_interval从1s改为30s整个链路才算彻底稳下来。这三个问题单独看都是小事但它们串起来恰好体现了流处理链路调优的系统性Source、算子、Sink、外部系统瓶颈任何一端没配合好整条链路就会被最慢的那一环拖死。4.4 Flink SQL、CDC与常见异常排查这几年Flink SQL的成熟速度非常快很多原来要写几百行DataStream代码的逻辑几张SQL就搞定了。尤其是Flink CDC出现后实时数据同步这事被彻底平民化了——以前同步MySQL到数仓要用Canal监听Binlog再写程序消费现在用Flink CDC组件声明一个source表放进SQL里直接读配合Flink SQL还能直接做全增量一体同步非常方便。不过Flink CDC也踩过不少坑最有代表性的一个旧版本Flink CDC默认只做全量同步不会自动切到增量监听需要靠配置开启而且某些复杂DDL变更比如加列CDC机制并不会自动把新列映射到目标端需要重建或手工修改同步逻辑。另外Flink CDC连接MySQL时权限问题也很经典账号必须要有SELECT, REPLICATION SLAVE, REPLICATION CLIENT或对应版本的新权限名等Binlog读取权限否则任务一启动就报Access denied。Flink SQL的另一个高频异常是JDBC连接器的连接问题常见报错是Failed to send data to Kafka或Could not find suitable driver之类的信息。这类问题绝大多数不是Flink的锅而是依赖冲突Flink自带的JDBC驱动版本和Mysql驱动版本不匹配或者是flink-connector-jdbc_2.12和mysql-connector-java放在lib下的目录方式不对。经验是外部依赖一律打包进作业的JAR里不要随手丢进Flink的lib目录去“碰运气”版本冲突会造成各种无法从报错信息直接找到原因的诡异问题。Spark那边的经典坑是Using Sparks default log4j profile: org/apache/spark/log4j-defaults.properties这类日志提示。这通常不是错误但经常误导新手以为Spark出问题了。更实际的价值是当你在YARN日志里看到这个信息时意味着Spark还在用默认log4j配置没有加载你自己提供的日志配置这时需要在spark-submit时使用--files log4j.properties并配合--conf spark.driver.extraJavaOptions-Dlog4j.configurationfile:log4j.properties来强制指定。这种配置层面的小坑往往比框架本身的功能问题更让人摸不着头脑。5. 到底怎么选型结合业务场景、数据规模与团队能力的决策框架到了最关键的部分——回到开头那个朋友的困惑到底选Flink还是选Spark。我的立场一贯是没有绝对的好坏只有适不适合自己的业务场景。判断维度可以从下面几个角度切入第一看业务延迟要求。如果你的业务对延迟非常敏感比如毫秒级的实时推荐、实时风控、在线特征计算那Flink几乎是唯一现实的选择。Spark微批次模型无论怎么优化遇到“实时”两个字都显得不够干脆。反过来如果延迟能接受秒级到分钟级——比如实时报表、运营看板、数仓分层里的DWD/DWS实时汇总Spark完全够用运维还更熟。第二看业务形态和计算需求。如果你的核心场景偏“批”而不是“流”比如离线ETL、大规模数据科学计算、机器学习训练前的大规模处理那Spark的生态MLlib、GraphX、庞大的DataSource家族是Flink短期内无法比拟的。Flink的强项在实时数仓链路里的持续处理、复杂事件处理、状态流计算上。我见过不少团队一上来就全量上Flink做所有实时任务结果发现一些偏批的维度计算在Flink里写起来极其别扭、调试成本极高最后只好把任务拆开批的归Spark流的归Flink。第三看团队的技术存量。传统Java后端团队转Flink相对平滑因为Flink本身就是Java/Scala生态API设计更贴近数据流编程的习惯。一个熟悉流式理论和Java的工程师上手Flink的周期远比想象中短。反过来如果团队已经用Spark做了好几年数仓有大量现成的Spark SQL脚本和调优经验那新增流处理场景时优先考虑Structured Streaming能把团队学习和运维成本降到最低——毕竟换来一个“Spark也能做流处理”的渐进式路线比推倒重来学一套新引擎划算得多。第四看生态和数据源环境。现在很多公司的实时数据都走Kafka、Pulsar这类消息中间件两边都支持得很好。但如果你需要做数据库CDC比如MySQL/Oracle/PG到数仓的实时同步那Flink CDC的成熟度和易用性目前是明显领先的Spark的CDC方案往往要自己拼装复杂度较高。另一个很实用的点如果下游主要落地到ClickHouse、Doris、StarRocks这类OLAP引擎两边都有完善连接器但如果下游是Iceberg、Hudi、Delta Lake这类湖格式Spark的湖仓集成尤其配合Hudi会更顺滑Flink生态对湖格式的支持这几年也在跟进但在复杂场景培养度和Spark还有些差距。第五看资源和可观测性投入。Flink的作业调试、Checkpoint调优、反压定位相比Spark要更精细但工具链也更完整——Web UI里各项指标详尽配合Prometheus监控可以做得很漂亮。Spark流处理的调优更偏向Spark传统那一套上手门槛相对低但遇到状态大、延迟要求高的场景时排查路径往往更绕。把这些维度落成一个简单的决策思维模型我个人总结是这样一句判断话术如果项目的主要矛盾是“延迟要低、状态要准、逻辑复杂”选Flink如果主要矛盾是“吞吐要大、批流一体、生态丰富”选Spark。当然大规模场景里两者完全可以共存一条实时数仓链路里用Flink做实时同步和核心实时计算用Spark做T1的离线补偿和批量数据对账互相补位。这种“FlinkSpark双引擎”的架构在我近年接触的中大型公司里已经是主流做法了。最后一个实操建议选型的最终拍板不要只拉着架构师开会听汇报而是拿一段最少可用的真实业务数据分别用Flink和Spark跑一遍端到端POC概念验证把延迟、吞吐、恢复时间、运维成本这四项指标测出来。纸面上的对比看得再多不如POC的测量结果来得可靠——这比任何博客和经验分享都更贴合你自己的业务。我这几年做实时化改造的体会是最有效的路径是“先确定业务核心目标再倒推引擎选择”而不是“听别人说某框架流行就先学它”。流处理是手段业务价值才是目的。把这条想透你的选型就是站在自己真实的问题之上做的决策。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →