尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Spark与Flink核心区别详解:架构、实时性、编程模型与选型指南

Spark与Flink核心区别详解:架构、实时性、编程模型与选型指南 1. 一个跑批老兵眼中的Spark和Flink做了这么多年数据开发Spark和Flink这两套东西几乎是大数据领域绕不开的两座大山。我自己从Spark 1.6时代就开始用后来因为实时业务需要又从零啃Flink期间踩过的坑、写错的代码、调不过的参都能凑一本数据工程师灾难回忆录了。这篇文章不是官方文档的复读机而是想把这两套框架放在一起从架构理念、实时性、编程模型、部署运维、真实场景选型这几个维度掰开揉碎讲清楚。不管你是刚入行的新人还是被面试官追着问Spark和Flink区别的求职者抑或是正在为项目做技术选型的老手我相信这篇文章都能给你一些参考价值。先说一个最核心的结论方便你有个整体认知Spark 本质上是批处理引擎用微批的方式模拟流处理Flink 本质上是流处理引擎把批处理当作有界流来处理。流批一体这个口号Spark喊得响Flink做得更彻底。这句总纲会贯穿全文所有细节差异基本都能从这句总纲里推出来。2. 架构与设计理念的根本差异2.1 Spark的计算驱动和Flink的数据驱动先讲架构层面的东西。Spark的核心抽象是RDD弹性分布式数据集它把数据切分成一个个partition在集群上并行计算。整个过程是任务驱动的一个Job被拆成多个StageStage之间通过Shuffle连接每个Stage内部是一个个Task这些Task由Driver节点调度到Executor上执行。DAG有向无环图是Spark最核心的执行计划概念每次行动操作都会触发一次完整的DAG执行。这里有个关键特征Spark的存储和计算是分离的。RDD本身不保存数据它只是描述数据如何从源头计算出来的蓝图。真正落地的数据要么在HDFS上要么在内存缓存中要么在外部存储里。Spark的执行是延迟的只有遇到action算子比如count、saveAsTextFile时前面的transformation才会真正执行。这种设计让Spark非常适合复杂的多阶段批处理任务比如ETL、数仓分层、机器学习特征工程。Flink则完全不一样。它的核心抽象是DataStream一切皆流。Flink的架构是数据驱动的数据一到算子算子立刻处理处理完立刻发送给下游算子整个过程像一条流水线数据持续不断地流过每一个算子。Flink也有DAG但它更准确的说法是数据流图每个节点是算子边是数据流通道。Flink的StreamGraph会进一步优化成JobGraph然后分发到TaskManager上执行。这两者最直观的区别Spark像是在工厂里做批次加工一批原料到了统一进行切分、打磨、包装然后运走Flink像是流水线作业原料一件一件地进来经过每个工位马上被处理下一件紧接着就来了。2.2 为什么Spark把流切成微批Flink却坚持真流这是一个特别值得聊的话题。Spark Streaming时代注意是老的Spark Streaming不是后来的Structured Streaming它的思路是把连续不断的数据流按照时间间隔切成一个个小批次比如2秒一个batch然后调度Spark批处理作业去处理每个小批次。这就是微批micro-batch的本质。微批的好处是显而易见的批处理的所有优化手段都能直接用容错机制简单有中间结果落盘代码执行确定性高吞吐量非常高。但缺点也很致命——延迟被限制在批次间隔以上你设置2秒一批那延迟至少是2秒而且批次边界的对齐问题会带来一定的数据倾斜和延迟抖动。Flink从设计第一天就没走这条路。它用的是连续流执行模型数据到了算子就处理不需要等待凑够一批。Flink的流水线在TaskManager之间是端到端的网络传输上游算子处理完一条数据立即序列化发送给下游算子中间不需要落盘因此延迟可以做到毫秒级别。这套设计源自Flink的德国血统柏林理工大学等机构发起的项目它从一开始就把流处理当作头等公民而不是批处理的附庸。我做个类比你就明白了Spark Streaming像是公交车固定时间发车乘客需要等车Flink像是出租车伸手即停随时出发。公交车的优势是票价便宜、载客量大出租车的好处是随时可走、路径灵活。类似的道理Spark适合吞吐量优先、对延迟不敏感的场景Flink适合延迟敏感、需要实时响应的场景。2.3 容错机制的截然不同血缘恢复 vs 分布式快照容错这块值得单独掰扯因为面试问到的概率极高也是实际运维中差异最大的地方。Spark的容错依赖血缘Lineage机制。RDD的每个transformation都会记录在血统中一旦某个分区的数据在计算过程中丢失比如Executor宕机缓存在内存中的数据丢了Spark会从源数据重新执行这一部分transformation来恢复它。这套思路很简单也很有用但它的恢复粒度是整个分区恢复时间取决于数据量和计算链路的长度。如果有一条特别长的血缘链比如20个transformation之后数据丢了恢复的成本会非常高。Flink的容错用的是分布式快照基于Chandy-Lamport分布式快照算法。它会周期性由checkpoint interval配置在数据流中插入屏障barrier把整个计算状态做一次全局快照保存到外部存储如HDFS、S3、RocksDB。一旦发生故障Flink从最近一次成功的快照恢复状态同时回放这段时间内的数据。这个机制可以让Flink实现端到端的精确一次Exactly-Once语义配合Kafka这类支持消息回放的Source可以保证数据不丢不重。这里要补充一个实操经验Flink的checkpoint设置不是越大越好也不是越小越好。太小导致频繁做快照对性能影响大太大导致故障恢复时间变长。我一般建议生产环境设置为30秒到60秒之间状态比较大的场景配合增量checkpoint用。你看网上那些Flink数据血缘的热搜词很多人是在问Flink里怎么追踪数据血缘关系。实际上Spark和Flink都有相应的机制但Flink的checkpoint机制天然保留了状态和数据流之间的关系对数据审计和血缘追溯有天然优势这也是我之前在做一个数据合规项目时首选Flink的原因之一。3. 实时性与处理模型这是两者最大的分水岭3.1 延迟对比秒级和毫秒级不是同一个量级如果只记住一个数字来说明两者的区别那就是延迟。Spark Streaming微批模式的典型延迟在1到10秒级别取决于batch interval即使后续的Structured Streaming优化了很多延迟也仍然在100毫秒到1秒这个区间而且它本质上仍然是靠微批实现的。Flink的流处理延迟是毫秒级的在标准网络环境下端到端延迟可以做到几十毫秒以内。这个差异在实时风控场景中是致命的。比如银行卡盗刷检测如果延迟2秒盗刷交易可能已经完成了如果延迟50毫秒系统有足够时间拦截交易。我之前参与过一个反欺诈项目最初用Spark Structured Streaming做实时特征计算结果因为延迟压不住而被迫切换到Flink。切换完以后延迟从秒级降到了几百毫秒以内模型效果马上就上来了。这个项目直接让我对实时性三个字有了更深刻的体感。当然不是说Spark不行Spark的优势是吞吐量和批量计算能力。同样是3TB的HDFS文件做聚合计算Spark比Flink要快不少因为Spark的批处理优化得更极致但说到每秒处理上百万条Kafka消息并且每条消息延迟都要求在100毫秒以内Flink的架构优势就体现出来了。3.2 时间语义和窗口计算Flink的杀手锏时间语义是流处理中最容易搞晕、也最影响正确性的东西。我先解释一下事件时间、处理时间和摄入时间这三个概念。处理时间Processing Time数据到达处理引擎时的系统时间。你什么时候处理就记什么时间。事件时间Event Time事件实际发生的时间通常在消息体里自带比如用户点击按钮的时刻。摄入时间Ingestion Time数据进入流处理系统的时间是事件时间到处理时间之间的折中。Spark在早期的StreamingDStream时代只支持处理时间这是它被诟病最多的地方之一。直到Structured Streaming才引入事件时间支持但实现上对乱序数据的处理能力依然有限更多是依赖watermark的周期性推进和状态清理。Flink从第一天就把事件时间当作一等公民。它提供了完整的watermark机制来处理乱序数据支持在事件时间上做窗口聚合滚动窗口、滑动窗口、会话窗口还能处理延迟数据side output给人用。这是什么概念就是说Flink能真正理解业务事件发生的先后顺序而不是只看数据到达系统的时间。举个例子用户0点下单付了款但消息通过网络延迟2分钟后才到Kafka。如果按处理时间算这个事件会被归入2分钟后的统计窗口按事件时间算它会正确地归入0点那一分钟的交易统计。这在实时报表、实时大屏、异常检测中都是核心能力。我之前做一个实时订单数据大屏刚开始用Spark结果发现订单归属时间段经常错位后来换成Flink用事件时间加水印才彻底解决。当时我们加了一个规则watermark延迟一分钟给乱序数据留缓冲实测下来准确率大幅提升。3.3 批处理单向流的有界流思想前面说过Flink把批处理看作有界流这个思想很有深意。在Flink中读一个文件其实就是读取一条有终点的数据流读完数据流自然结束。处理逻辑上不需要区分我是在做批处理还是在做流处理同一套API可以既处理有界数据又处理无界数据。这就是Flink流批一体的真正含义。Spark反过来它的根本是批处理流处理是在批处理框架上做的扩展。具体到API层面Structured Streaming把流抽象成不断增长的无界表每次微批任务其实就是一次小的批处理任务。写起来很顺手但底层执行机制终究是攒一批算一批和Flink的持续计算有本质差异。实际项目中这个差异会影响什么最典型的是状态管理和精确一次语义。Flink天然支持有状态流处理状态可以是Keyed State按Key维度保存的状态可以跨事件保存中间结果比如计算每小时内每个用户的累计消费金额这种需求Flink的API做起来非常顺手。Spark要想实现类似功能需要借助外部存储如Redis来手动管理状态或者用updateStateByKey这类算子但性能和规模都有限。这里我想插一句如果你在面试或者实际项目中遇到状态管理这个话题可以记住这个结论——Spark是无状态批次计算模型状态需要外部存储协助Flink是原生的有状态流计算模型状态内置且支持容错。这句话基本可以终结80%关于两者差异的讨论。4. 编程模型与API生态写起代码来感觉完全不同4.1 Spark SQL/DataFrame的推拉式魅力Spark在API设计上的成功是不可否认的。RDD时代其实挺反人类的你要写很多底层代码。但DataFrame/Dataset API出现后Spark的使用门槛大幅下降写起来非常像写SQL但又保留了代码的灵活性。我个人的感受是Spark让大数据开发离SQL工程师越来越近。你只要把数据读进来然后用类似SQL的语法df.groupBy(col1).agg(sum(col2))就能完成聚合分析。Spark Catalyst优化器会根据你的写法自动优化执行计划你不用关心底层怎么跑的。还有SparkSQL直接支持纯SQL语法比如spark.sql(SELECT * FROM t WHERE ...)这大大方便了从传统数据库迁移过来的团队。在Spark ETL脚本这个热搜场景下Spark的表现更是无可挑剔。你可以把十几个数据源读进来做清洗、过滤、join、聚合、窗口计算最后直接写回目标表整个流程用DataFrame API可以实现得非常简洁。我用Spark写了不下几百个ETL脚本稳定性和性能都很好基本上不需要太多人工干预。对于常规的T1批处理任务Spark依然是当之无愧的首选。4.2 Flink SQL实时数仓的香饽饽Flink的DataStream API本身就是一大优势但真正让我觉得Flink已经能跟Spark掰手腕的地方是Flink SQL。Flink SQL发布以来社区热度一路飙升现在还经常看到Flink SQL相关的热搜词上榜这说明了它的实用价值。Flink SQL允许你直接用SQL语法处理无界流数据。你可以定义一个Source表比如映射到Kafka主题再定义一个Sink表比如映射到MySQL或ClickHouse然后把两者用一个INSERT INTO SELECT语句连接起来Flink就会持续不断地消费Kafka数据、执行计算、写入目标存储。这个过程是持续运行的你不用手动触发、不用等批次结束完全是流式的。如果你搜过Flink SQL相关的资料会看到很多类似这样的操作先建Kafka维表然后和流表做joins或者用窗口函数做实时聚合最后写入指标系统。这些都是实时数仓里的经典玩法。Flink SQL还有一个很强的点——它支持维表关联。实时数据流到之后你需要去MySQL或者HBase里查维度信息比如用户姓名、商品分类来做关联。Flink SQL的Temporal Table Join可以让你实时查询外部维表并自动处理维表数据变化。这个功能在实时数仓场景里极其常用Spark SQL在批处理模式下也能做类似的事但在流模式下就弱了不少。4.3 DataFrame和Flink Table API的体验对比一个例子为了让你有更具体的感受我用一个案例来对比。假设我们要统计每个用户每天的总消费金额消费数据从Kafka进来。SparkStructured Streaming写法大致是// Spark Structured Streaming 微批处理 val df spark .readStream .format(kafka) .option(kafka.bootstrap.servers, localhost:9092) .option(subscribe, orders) .load() .selectExpr(CAST(value AS STRING) as json) .select(from_json($json, schema).as(data)) .selectExpr(data.user_id, data.amount, data.ts) .withWatermark(ts, 1 minutes) .groupBy($user_id, window($ts, 1 day)) .agg(sum(amount).as(daily_total)) .writeStream .format(console) .outputMode(update) .start()Flink SQL写法大致是-- Flink SQL 连续流处理 CREATE TABLE orders ( user_id BIGINT, amount DECIMAL(10,2), ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 1 MINUTE ) WITH ( connector kafka, topic orders, properties.bootstrap.servers localhost:9092, format json ); CREATE TABLE daily_sum ( user_id BIGINT, total_amount DECIMAL(10,2), window_start TIMESTAMP(3) ) WITH ( connector print ); INSERT INTO daily_sum SELECT user_id, SUM(amount), TUMBLE_START(ts, INTERVAL 1 DAY) FROM orders GROUP BY user_id, TUMBLE(ts, INTERVAL 1 DAY);你看Flink SQL写起来跟写标准SQL几乎一模一样学习和迁移成本极低。这也是为什么Flink在实时数仓领域的渗透率越来越高很多从Oracle、MySQL背景转过来的人见到Flink SQL的第一反应都是原来实时计算还能这么写。4.4 机器学习与图计算的生态差异说完数据处理再简单说下生态。Spark最大的生态优势之一是MLlib机器学习库和GraphX图计算这是Spark能覆盖批处理机器学习全链条的核心竞争力。你可以用Spark做特征工程用MLlib直接训练模型比如逻辑回归、随机森林、ALS推荐然后用Spark Structured Streaming做模型的上线预测打分一整套流程不需要切换框架。Flink在机器学习方面相对薄弱虽然有FlinkML项目但成熟度和社区活跃度都不如Spark MLlib。如果您要做在线学习、实时特征计算加模型推理Flink可以配合其他工具来做例如使用modelserver或外部推理服务但生态的完整度确实不如Spark。这也是很多公司在技术栈统一的考虑下继续选择Spark的原因——一个框架搞定ETL、数仓、机器学习团队学习和维护成本都低。如果你公司的业务以离线为主实时只是补充那Spark肯定更适合如果实时业务本身占大头那你可能需要在Flink之外再配一套其他机器学习工具。5. 部署、运维与常见坑生产环境的真实体验5.1 集群部署Spark简单Flink也不难但需注意细节这个部分网上各种安装教程很多我简单说下两者在大数据集群里的部署差异。Spark的部署模式有Local、Standalone、YARN、Mesos、Kubernetes几种。生产环境最常用的是Spark on YARN因为大部分公司的大数据集群已经部署了Hadoop直接复用YARN资源调度器就行了。Spark任务的提交方式也很简单spark-submit --master yarn --deploy-mode cluster然后等它跑完。在spark集群搭建这个热搜词下面你能找到大量教程思路基本都是先装Hadoop再配Spark环境变量再启动Master和Worker。总体来说Spark的部署链路比较成熟踩坑概率不高。Flink的部署方式包括独立集群Standalone、YARN、Mesos、Kubernetes。生产环境我推荐Flink on YARN或者如果公司已经上了K8s可以直接用Flink on Kubernetes。Flink on YARN有个特别方便的特性——per-job cluster模式每个Flink作业都动态申请一个完整的Flink集群作业结束自动释放资源多个作业互不影响。用flink run -m yarn-cluster -yn 3这样的命令就能提交作业。有个容易被忽略的坑Flink的jobmanager和taskmanager的内存配置。如果你在集群里部署Flink不设置jobmanager.memory.process.size或taskmanager.memory.process.size默认值可能与集群实际资源不匹配尤其是容器化环境下容易OOM。我建议你在提交作业之前先确认自己给TaskManager分配的内存和CPU核数然后显式设置这几个参数jobmanager.memory.process.size: 2g taskmanager.memory.process.size: 4g taskmanager.memory.managed.size: 2g # 状态后端用5.2 两者的状态管理与运维复杂度状态管理是Flink运维中最需要关注的维度。Flink的状态分为Keyed State和Operator State存储在后端StateBackend里可以是内存、RocksDB或者混合模式。生产环境状态一般不放在内存容易OOM我普遍推荐用RocksDB StateBackend它把状态持久化到本地磁盘天然支持增量checkpoint能扛住很大的状态量。这里有一个我踩过的坑。状态大小会随着运行时间不断增长特别是按用户维度做窗口聚合时如果用户的维度很大比如几亿用户状态文件动辄几十GB。如果你没有设置状态过期时间TTL那么状态永远不清理最终磁盘爆满或恢复超时。所以生产环境一定要设置State TTL类似这样StateTtlConfig ttlConfig StateTtlConfig .newBuilder(Time.hours(72)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build();Spark这边没有状态管理的概念它的容错依靠血缘和Cache的自动清理机制。你只需要关心Executor的内存设置和Shuffle的调优比如spark.executor.memory、spark.shuffle.memoryFraction等参数。部署和运维的日常维护成本相对更低。这也是不少团队在能不用实时就不用实时的原则下倾向于Spark的原因之一。5.3 常见连接器异常与排查实录JDBC和HDFS写代码的都知道连接器永远是踩坑重灾区。热搜词里面出现的flink的jdbc连接器异常我太有共鸣了。Flink JDBC连接器用起来有几个经典大坑第一个坑是连接数耗尽。Flink的并行度如果很高比如50每个并行度都会定期往目标数据库写数据默认的连接池可能瞬间被占满。解决办法一是调大数据库的连接数上限二是在Flink Sink端合理设置SinkFunction的重试和批量写入配置。对MySQL写入我建议用org.apache.flink.connector.jdbc这个官方连接器并设置sink.buffer-flush.max-rows为1000左右如果你的JDBC目标是Doris这类数据库可以考虑使用官方适配器。第二个坑是数据类型不匹配。网上有句报错信息挺出名flink type is datev2, but arrow type is dateday。这是Flink和下流系统Doris之间日期类型映射不一致导致的。简单说Doris的DateV2在Flink读入时被映射成了某种类型但是在JDBC/Arrow传输过程中转换失败。解决方案也简单在下发建表SQL时显式把日期字段指定为DATE而不是DATETIME或者在Flink侧用CAST把类型收窄/转宽。这种报错本身不可怕怕的是你不会看日志去推断类型不一致的问题。排查思路我总结下来就三步先看Source端的schema定义、再看目标端建表语句中的字段类型、最后对比两者之间的类型映射表。第三个坑是flink 一定要hdfs这个问题。很多人疑惑Flink不部署HDFS行不行。答案是Flink本身不强制依赖HDFS但如果你要开checkpoint生产环境几乎必开那必须有一个支持持久化的文件系统。HDFS是最常见的选择但S3、OSS、GCS、甚至本地文件系统也能用。唯一要注意的是如果你用RocksDB做状态后端还要选对状态存储路径不然任务一重启状态全丢。我之前有个项目为了省成本一开始没上HDFS用的本地文件系统做checkpoint结果JobManager一重启checkpoint全丢了数据从头开始消费。后来老老实实接上HDFS一劳永逸。Spark这边也有很多连接器问题比如HDFS连接超时、Spark SQL连接JDBC目标数据库报错等。排查经验基本上遵循先网络、再认证、再参数、最后类型。如果你遇到Spark读取MySQL时报错大概率是MySQL驱动版本与Spark自带的JDBC版本冲突解决方法有两个一是升级MySQL Connector/J到最新版二是在启动脚本里显式指定driver类名。5.4 资源消耗分析谁更吃内存说到资源这是所有大数据项目绕不开的预算问题。Spark在批处理场景下的内存消耗并不小——每个Executor要预留一部分内存给存储缓存一部分给执行Shuffle、Join等。尤其是做大量Join操作的时候Shuffle产生的临时文件会占不少磁盘和CPU。社区里关于Spark内存的讨论这么多也说明大家普遍被这个问题困扰过。Flink的资源消耗取决于你开了多少并行度、状态有多大、是否开RocksDB。Flink的checkpoint本身会占额外的网络和磁盘I/O如果你的checkpoint间隔过短、状态又大资源开销会非常明显。我用过40个并行度的Flink任务开着RocksDBcheckpoint间隔30秒最顶峰的时候光状态存储就要占几百GB磁盘。从性价比角度看同样的逻辑在批处理场景用Spark跑通常比用Flink省资源因为Flink为了支持流式计算会有额外的记录序列化开销和网络传输开销。但在纯流处理场景里Flink的单位资源吞吐量是明显优于Spark Streaming老DStream的和Structured Streaming打平或略优。6. 真实场景下的选型建议和面试答案6.1 怎么选一句话版和详细版先给一句话版本延迟敏感型、事件时间敏感型、需要精确一次语义、需要持续运行的有状态计算选Flink批处理、ETL、复杂SQL分析、机器学习全链路、吞吐优先且延迟能容忍到秒级选Spark。然后说详细版的考量维度。我见过很多团队在技术选型时争论不休其实吵来吵去都是没抓住问题的本质。下面这4个问题问完答案基本就出来了。第一个问题你的数据是批的还是流的如果数据每天、每小时落地一次用批处理方式算那Spark天然合适如果数据是持续不断产生的流Kafka、Pulsar而且你需要立刻计算那就Flink。第二个问题你对延迟的容忍度是多少5分钟那Spark也可以3秒你要认真考虑Flink1秒以内别犹豫了直接上Flink。延迟这个指标可不是简单的体验问题它直接关系到业务价值能否实现。第三个问题你的数据形态和计算模式的复杂度如何如果你的核心场景是把N张表做join后做聚合分析那Spark SQL的Catalyst优化器比Flink SQL在批处理查询上的优化要成熟得多如果你的核心场景是每条Kafka消息做多层规则判断、状态更新、窗口统计、实时输出那Flink碾压Spark。第四个问题你的团队能力和离线/实时技术栈现状如果团队已经深度会用Spark且离线数仓已经很稳定实时需求只是锦上添花那用Spark Structured Streaming就够了没必要为了用Flink去重构全套技术栈。反过来如果实时业务是公司增长的核心那即使团队要花一个月去学习Flink也是值得的投资。6.2 从绝密100个Spark面试题看考点面试怎么答网上那个绝密100个Spark面试题熟背100遍的说法我看了都想笑——如果背题能解决面试问题那面试官的价值何在但既然面试确实会考我给你整理几道最核心的Spark/Flink区别题以及建议回答的方向。Spark和Flink的核心区别是什么不要只说一个是批一个是流。建议分三点第一执行模型上Spark是微批/批处理模型Flink是连续流模型第二延迟上Spark秒级Flink毫秒级第三状态管理上Flink原生支持有状态流处理配checkpoint实现精确一次Spark的状态管理依赖外部存储或微批的重算。为什么Flink能实现毫秒级延迟答案核心是Flink是持续处理模型数据到达即处理无批次等待同时通过Distributed Snapshot机制做状态快照不需要像微批那样等一个完整批次算完再做checkpoint。还有网络传输层面Flink的资源管理是流水线式的TaskManager之间可以复用网络缓冲池避免频繁创建和销毁连接。Flink的窗口和Spark Streaming的窗口有什么区别这个值得好好回答。Spark Streaming的窗口是基于微批的窗口边界是批次编号的倍数窗口数据就是多个批次数据的叠加Flink的窗口是基于时间语义的你可以定义事件时间、处理时间、摄入时间并能通过watermark处理乱序数据。用一句话Spark的窗口是按批次打包Flink的窗口是按时间切片。哪些场景Spark和Flink可以互相替代如果按固定批次消费Kafka数据批量拉取再处理且SQL逻辑不复杂那Spark和Flink都能用如果业务接收秒级延迟、且不需要精确的Event Time语义Spark完全够用。技术选型没有绝对的最优只有适不适合。6.3 场景案例实时数仓、数据分析、用户复购最后分享几个我实际做过的场景来说明。场景一实时用户行为分析推荐Flink。客户要实时看到每个页面的PV/UV延迟要求小于1分钟。我们当时用了Flink SQL直接消费Kafka里的埋点数据用TUMBLE窗口做分钟级聚合写到Doris/ClickHouse前端大屏直接读。整个过程Flink SQL只写了不到100行上线以后延迟大概3秒左右包括Kafka到Flink到OLAP的端到端耗时客户非常满意完全没有Spark的参与。场景二离线用户复购率分析推荐Spark。客户要算过去30天用户的复购率数据是历史订单表量级在几十亿行。这种肯定是批处理用Spark最简单。当时我们写了一个Spark ETL脚本从Hive读用户订单按user_id聚合出购买次数大于等于2的用户数/总用户数作为复购率跑完写回目标表大约30分钟跑完非常稳定。场景三金融实时风控强烈推荐Flink。线上交易系统的风控引擎要求每笔交易在100ms内判定是否可疑。这种情况Spark Streaming的微批压根就顶不住直接用Flink DataStream API写规则引擎结合事件时间做窗口统计状态用RocksDB存用户历史行为特征配合Tidb或HBase做维度存储端到端延迟控制在50ms左右。这种场景如果你选错了框架项目可能就直接黄了。7. 数据血缘、SQL支持与功能演进趋势7.1 数据血缘Flink原生的血缘能力数据血缘这个热搜词也值得单独说说。在大数据治理、合规审计中你需要知道每一张报表、每个指标的底层数据从哪来、经过了哪些加工。Spark中数据血缘是通过RDD的Lineage实现的但它是物理层面的血缘更多用于容错恢复不适合直接用作业级和数据级治理。Flink则不同因为checkpoint机制天然会记录数据流的状态迁移加上Flink SQL的Statement SET、EXPLAIN等工具你可以追踪一条数据从Source到Sink的完整路径。很多在线下数据平台、元数据管理系统通过解析Flink作业的JobGraph和ExecutionPlan来生成数据血缘。如果你做的是To B项目客户现场往往有强合规要求数据血缘能力可以直接变成你的销售亮点。7.2 SQL功能完整度与持续演进Spark在批处理SQL领域耕耘多年Catalyst优化器和Tungsten执行已经是相对成熟的技术对复杂SQL的优化支持比如谓词下推、列裁剪、动态分区裁剪都很出色。Flink SQL虽然起步晚一些但发展速度极快而且它特有的流式SQL能力是Spark不具备的——你可以用普通的INSERT INTO语法实现持续不断的写入而不是一次性任务。我之前做个一个实时大屏项目写了一条类似这样的Flink SQL从Kafka读取交易流水关联MySQL维表商户表按5分钟窗口聚合成各商户交易金额然后写入ClickHouse。整个过程全部用SQL定义完成业务方看着都很惊讶因为以前这种实时统计至少需要写几百行Java代码。这就是Flink SQL的魔力也是它热度持续攀升的原因。从长期趋势看Flink SQL的功能会越来越接近Spark SQL两者都在往统一流批的方向走。对开发者来说抽象层越来越好用底层能力则各有侧重Spark还是更擅长吃内存拿吞吐Flink更擅长保状态扛延迟。7.3 Databricks的Spark与Apache Flink的社区走向还要提一句开源社区。Spark最大的推手是Databricks这是一家商业化公司很多Spark核心功能比如Delta Lake、MLflow是Databricks在推但Apache Spark项目本身是Apache基金会的顶级项目开源社区非常活跃。Flink这边现在是Ververica原Data Artisans在做商业化支持Apache Flink本身也是顶级项目社区活跃度在流处理领域尤其高。从招聘市场看现在要求同时掌握Spark和Flink的岗位越来越多。如果你在做职业规划我的建议是先用Spark把批处理功底打牢再用Flink深入流处理。这两者的适用场景有差异但核心的数据处理思维分区、Shuffle、JOIN、窗口是相通的。学会了底层原理换框架只是换个API的事情。8. 生产环境踩坑备忘录从安装到调优8.1 集群装好了不代表能用几个必看的配置不管你是照着spark安装与使用还是flink安装配置到部署的教程来操作装完之后都不建议直接上线。有几个关键配置一定要检查。Spark方面我一般会优先确认这几个参数spark.sql.shuffle.partitions默认200如果你的数据量小200个分区太浪费调低到50数据量大要调高到500以上。spark.executor.memory和spark.executor.cores这俩是最影响性能和稳定性的不要贪多。每个Executor的并发task数和内存要匹配否则GC频繁任务全卡死。spark.dynamicAllocation.enabled如果想自动伸缩Executor可以打开但要注意和ResourceManager的配额冲突。Flink方面我推荐先在本地以flink run命令跑通一个简单的WordCount再上集群。上集群前重点看这几个配置taskmanager.numberOfTaskSlots单台机器上能跑多少slot不是你机器核数越多越好要留点资源给操作系统和网络组件。parallelism.default默认并行度建议通过提交作业时用-p参数显式指定而不是依赖配置文件。restart-strategy作业失败后自动重启策略建议用failure-rate比如5分钟内最多重启3次。state.checkpoints.dircheckpoint保存路径用hdfs://namenode:8020/flink/checkpoints。8.2 常见报错排查速查表含热搜问题这里把几个高频问题以表格形式整理出来有遇到类似情况的可以直接对照排查。报错或现象可能原因排查思路Spark Executor OOMExecutor内存不够或Shuffle内存占比不合理查看YARN日志确认task的内存消耗调整spark.executor.memoryOverheadSpark Job卡在某个Stage数据倾斜少数task处理了大量数据用--conf spark.sql.shuffle.partitions400调大分区或者对Key加盐提交Flink作业时YARN队列资源不足YARN队列的队列容量满了换队列或者申请更多资源配额Flink checkpoint持续超时状态过大或网络I/O瓶颈调大checkpoint interval启用RocksDB增量checkpoint检查是否背压Flink任务数据重复消费checkpoint失败后从旧checkpoint恢复确认Source端是否开启Exactly-OnceKafka source要配置setStartFromLatest()或精确一次模式JDBC连接器写入MySQL卡死连接池满了或SQL事务长时间不提交调大连接池、减少batch size、检查MySQL锁等待flink type is datev2, but arrow type is datedayFlink和Doris/目标系统日期类型映射不一致在Flink侧用CAST统一类型在目标表DDL里明确字段类型8.3 一次Flink生产事故复盘checkpoint全没了最后分享一个真实教训这个事故我至今记得。项目背景一个Flink作业从Kafka消费用户行为日志经过状态计算后写到Elasticsearch。运行了一个多月都正常某天突然JobManager宕机重启后作业恢复但ES中的指标从某个时间点开始出现重复累加。排查过程花了很长时间。后来发现根因是Flink作业虽然配置了RocksDB和checkpoint但checkpoint目录指向的是本地磁盘而不是共享存储HDFS/OSS。JobManager宕机后新启动的JobManager无法访问旧TaskManager的本地状态文件于是只能从最近一个全局checkpoint恢复实际是空整个状态被清空Kafka从最新位置开始消费导致之前积累的统计全部丢失。从那以后我给自己定了一条铁律Flink生产环境的state backend存储路径必须放到共享存储上。无论是HDFS、S3还是OSS一定要保证多个节点都能访问。这个教训也写在了我们团队的技术规范里后面再没出现过类似的坑。9. 最后一点实践心得想给还在纠结选型的朋友一个建议不要听别人说Flink流处理天下第一就盲目切换也不要觉得Spark过时了——这两种说法都太肤浅了。我在实际中做的最多的项目是Spark做离线数仓、Flink做实时链路两套框架配合起来用覆盖了绝大多数业务场景。如果你现在有一个实时需求要评估可以直接跑一个小demo来对比用Kafka生产100万条数据分别用Spark Structured Streaming和Flink都消费一遍统计端到端延迟、吞吐量、资源消耗。实测数据比任何口水战都有说服力。我当年就是这样做完对比后才坚定了在实时项目中用Flink的决定——虽然Spark够用但Flink的精确一次语义和事件时间支持让我少写了至少一半的纠错代码。另外如果你还是学生或者初入行者我的建议是两条腿走路先掌握Spark的DataFrame/SQL再做Flink的DataStream/SQL。这两套API的思维其实底层相通等你理解了分布式计算里的分区、Shuffle、Join、窗口这些概念后具体框架只是一个顺手工具。真正值钱的是你对数据处理流程的理解而不是背了多少面试题。
返回列表