Hadoop、Spark、Flink怎么选?从批处理到流处理的技术演进与实战权衡
选型之前先看清三个引擎各自站在什么位置上个月我帮一家做供应链系统的公司做技术选型评审数据团队负责人一上来就问我我们要不要放弃Hadoop全量切到Flink我反问他为什么想切他说因为领导看到社区里都在说“Hadoop已死Flink才是未来”。这个回答我听过太多次了。过去几年里“XX已死”的句式几乎每年换一个主角——先是MapReduce被Spark“拍死”接着Spark Streaming被Flink“拍死”。可现实是HDFS到现在依然是很多企业数据底座的一部分Spark承担着大量离线ETL和机器学习任务Flink则在实时数仓里越来越吃香。这三个引擎并不是一代淘汰一代的关系而是各自占据了不同的计算边界。这篇想把Hadoop、Spark、Flink的选型逻辑讲透它们底层到底差在哪儿各自擅长什么哪些业务场景该用哪个以及我过去几年做选型时反复权衡的那几个问题。适合正在搭大数据平台的技术负责人、准备转行数据开发的工程师以及被各种面试题里Hadoop/Spark/Flink搞得头疼的同学。1. 选型不是比新老先回答三个问题再说1.1 从磁盘批处理到内存计算再到流优先三个时代的不同答案把时间倒回2006年左右Google发表MapReduce论文之后Hadoop迅速成为大数据的事实标准。那时候的数据处理核心思路是“先把数据切块放到HDFS上再用MapReduce任务把结果算出来”。MapReduce的设计目标是吞吐优先它能用成百上千台机器处理PB级数据但代价是每一步中间结果都要落盘所以一个复杂任务跑几十分钟甚至几小时很正常。2014年前后Spark开始大规模流行它的核心贡献是把中间结果尽量留在内存里。同样是跑一个多阶段的复杂作业Spark可能比MapReduce快一个数量级。再加上Spark SQL把门槛降到“会写SQL就行”很多人都觉得Spark就是大数据的终极答案。再往后业务对实时性的要求上来了。风控要毫秒级识别异常操作大促要实时盯流量和订单监控要秒级看到指标变化。批处理无论如何都满足不了这种时效性所以Flink这种“天生以流为核心”的引擎开始登上主舞台。注意Flink不是“处理得更快的批处理”它在架构设计上和批处理完全是两条路。1.2 选型先回答的三个问题我一般会把选型决策压缩成三个问题尤其适合第一次搭数据平台的团队你的数据能容忍多久的延迟分钟到小时级别批处理就够了Spark和Hive都能胜任秒级到分钟级考虑Spark Structured Streaming这类准实时方案毫秒到秒级才需要上Flink这样的真流处理。你的数据是有界的还是无界的这里的“有界”指的是数据在某个时间点就结束了比如昨天的订单表“无界”指的是数据永远在产生比如App里每时每刻的点击日志。无界数据天然适合流式处理框架因为你需要“持续计算”而不是“跑完就结束”。团队能长期维护哪种技术栈这是最容易被忽视的问题。Hadoop、Spark、Flink全都涉及Java/Scala体系但上手成本差异很大。只写SQL的团队可能一开始用Hive或Spark SQL更顺手有Java/Scala底子的团队才有余力把Flink的流式任务真正玩明白。1.3 “新引擎替代旧引擎”是最大的误读不少读者会把“Hadoop”和“MapReduce”画等号然后得出“Hadoop过时了”的结论。但Hadoop生态里不只是MapReduce还包括HDFS分布式文件系统、YARN资源调度、ZooKeeper协调服务。哪怕你完全不用MapReduce只用Spark和FlinkHDFS依然可以充当底层存储YARN依然可以是任务调度平台。换句话说Hadoop这个“底座”并没有消失消失的只是“必须用MapReduce写作业”这个习惯。我见过很多团队喊着“去Hadoop”结果只是把MapReduce换成了Spark SQLHDFS换成对象存储或者保留原样。所以不要被“谁替代谁”这种口号带偏先看业务需要什么再看引擎擅长什么。2. 从运行原理看本质区别磁盘、DAG与状态流2.1 Hadoop MapReduce为什么慢以及它“慢但稳”的价值用一个最经典的WordCount来感受MapReduce的运行方式输入文本被切成一个个分片map阶段把每个词拆出来打上“词, 1”的标签输出先写到本地磁盘shuffle阶段系统把相同key的记录拉到一个reduce节点上并排序reduce阶段再做累加。中间结果反复落盘是MapReduce慢的根本原因。但换个角度看这种设计换来了极强的稳定性每一步都有中间文件和容错兜底哪怕某个节点挂了数据也可以从上一个阶段恢复。直到今天MapReduce仍然在某些超大规模离线场景里有存在价值比如数据量极大、又不能把所有中间结果都塞进内存的任务。但在绝大多数实际项目里Spark已经能承接这类离线批处理的活儿而且快得多。所以现在很多企业里的MapReduce作业主要是一堆历史遗留任务新任务基本不会再写了。2.2 Spark的RDD、DAG与惰性求值内存红利到底是怎么来的Spark的核心抽象是RDD弹性分布式数据集。你的数据被切分到多个分区分布在集群的不同机器上。开发者对RDD调用一系列转换算子比如map、filter、joinSpark不会立刻执行这些操作而是先构建出一张由操作构成的有向无环图DAG——我喜欢把它想成一张“菜谱”上面记录了每一步该做什么但真正开火做饭要等“行动算子”出现也就是需要输出结果的那一刻。这种“惰性求值”的设计配合“血缘”机制是Spark能比MapReduce快的关键。血缘记录了一个分区是怎么由源头数据一步步转换来的。某个分区算到一半节点挂掉了Spark可以根据血缘从原始数据重新算出来不需要像MapReduce那样把每一步中间结果都落盘。再加上引擎会把多个窄依赖的算子尽量在一个阶段内串起来执行减少shuffle次数性能自然就上来了。Spark的另一个红利是内存计算对迭代式算法极友好。机器学习里很多算法要反复扫描数据调整参数如果每轮迭代都要读写磁盘那基本没法训练。Spark把训练数据缓存在内存里迭代效率直接成倍提升。这也是为什么到目前为止Spark MLlib在机器学习流水线里仍然比Flink更适合。但内存不是无限的数据量超过executor内存时就会溢写到磁盘所以后面我要专门说Spark内存配置的坑。2.3 Flink的Checkpoint、状态与事件时间它凭什么敢叫真流式Flink一上来就按“事件驱动”设计。数据流里的每个事件到了算子就立刻触发计算不存在“攒一批再算”的概念。但这只是“流式计算”的表面Flink真正强大的是三个配套能力状态管理、Checkpoint和时间语义。先看状态。如果一个算子在持续接收同一个用户的登录事件要统计这个用户连续登录几天那“这个用户之前登录了几天”就属于算子内部的中间状态。Spark Streaming的微批模型天然弱化状态概念你需要把状态放到外部存储里反复对账而Flink把状态做成引擎一级的能力直接保存在StateBackend里读写都在本地速度很快。再看Checkpoint。分布式系统里节点随时可能挂掉Flink通过给“算子状态 数据源的位置”定期做分布式快照实现故障后的精确恢复。这个概念听起来简单但“流”是无穷无尽的你把状态快照刷到哪、怎么保证快照之间数据不重不漏全都要靠barrier对齐这类机制来实现。这也是Flink代码学习曲线比较陡的原因之一。最后是时间语义。Flink支持事件时间Event Time、处理时间Processing Time和摄取时间Ingestion Time配合Watermark机制处理乱序数据。比如一个传感器事件晚到了10秒用处理时间统计就会算错窗口但用事件时间配合Watermark可以在一定容忍范围内把它归到正确的窗口里。Spark Structured Streaming虽然也支持事件时间但底层仍然是微批调度在事件驱动和低延迟场景下始终不如Flink“原汁原味”。3. 一张对比表收拢核心差异外加几个容易搞混的边界3.1 从延迟、状态、SQL、生态看三者的真实位置维度Hadoop以MapReduce为核心SparkFlink核心计算模型批处理批处理 微批流真流处理典型延迟分钟到小时秒到分钟毫秒到秒中间结果存储磁盘内存优先可溢写磁盘流处理数据在管道中流转状态管理不内置依赖外部组件弱基本靠外部存储强状态管理内置StateBackend时间语义无完整概念处理时间为主支持事件时间事件时间 / 处理时间 / 摄取时间SQL能力Hive SQLSpark SQL成熟离线数仓常用Flink SQL流批一体近几年进步很快机器学习生态弱Spark MLlib成熟较薄弱典型运行环境YARNYARN/Kubernetes/StandaloneYARN/Kubernetes/Standalone从这张表能清楚地看到Spark最舒服的位置在“需要跑复杂批处理、又要兼顾一些准实时需求的场景”Flink最舒服的位置在“低延迟、高状态、强一致性要求的实时链路”Hadoop则更多以HDFS、YARN这些基础设施身份出现而不是“写MapReduce作业”这件事。3.2 被误读最多的四个边界问题“Hadoop已死”不等于“HDFS已死”。HDFS作为分布式存储到今天仍是很多平台数据湖的地基只是新架构里它和对象存储比如MinIO、云上OSS/S3经常共存计算层让位给了Spark和Flink。“Spark Streaming不等同于Flink”。Spark Structured Streaming在2.x以后已经很不错但微批调度的本质决定了它在窗口粒度、事件时间和状态一致性上会比Flink吃力。实时场景不要轻易拿Spark Streaming硬顶Flink的位。“Flink也扛得住批处理不代表你该把离线全切给Flink”。Flink 1.x开始推流批一体思路是对的但生产环境里批处理生态如Hive元数据、Spark SQL各种成熟语法、成千上万的历史调优参数不是说换就换。“Hive还没死”。很多公司把Hive当数仓的SQL入口底层引擎可能换成了Spark或者Tez但Hive Metastore这个元数据服务反而越来越重要Spark SQL和Flink SQL都会去连它。所以“Hadoop生态过时”这个说法根本不成立。4. 从真实的搜索需求看三者的典型使用方式4.1 集群搭建和安装配置类需求入门者的必经之路回头看大家的搜索记录“hadoop安装与配置”“ubuntu 安装hadoop”“hadoop伪分布式搭建”“hadoop集群搭建”“hadoop的docker镜像”这一类占了很大的比例。这背后的需求很明确新人入坑大数据第一关就是环境搭建。我的建议是不要贪多求快。按“单机伪分布式 → 三节点集群 → 容器化部署”的顺序走。伪分布式有两个作用一是验证HDFS读写流程二是搞懂配置文件的依赖关系比如namenode和datanode怎么互相发现、ZooKeeper参与HA时起什么作用。伪分布式跑通之后立刻上三节点集群因为单机环境里你根本体会不到“数据被切到多台机器、shuffle跨节点传输”是什么滋味。有条件的话再用Docker镜像起一个包含HDFS、YARN、Hive的小集群这个过程基本就是在帮你复习“Hadoop生态由哪些组件构成”。热搜词里有“hadoop和zookeeper整合实战”很多人顺手就把ZooKeeper忽略掉其实NameNode高可用和好多分布式协调场景都离不开它。4.2 Spark做用户行为分析离线数仓里的中坚力量“用户复购率 spark 脚本 redshift”这个热搜词特别有代表性。用户复购率这类问题数据往往存到数据仓库里比如Redshift然后你用Spark脚本把用户订单明细读出来按用户ID聚合统计“多次下单的用户占比”或者“间隔多少天再次购买”。这种场景天然适合Spark SQL或DataFrame API数据是“有界的”历史数据计算一次可能要扫描几亿行但对结果时效性没那么敏感。我见过很多团队把“跑批”的任务从Hive往Spark迁移之后执行时间缩短一大截但过程中也踩过不少坑。比如Spark默认并行度和实际文件数量不匹配导致executor忙闲不均比如某个用户的订单量极端大group by的时候出现数据倾斜。这些都是离线分析的经典问题后面避坑部分我会展开。4.3 Flink CDC和Flink SQL实时数仓绕不开的组合“flink cdc”“flink sql client sql gateway”“flink 数据血缘”“flink的jdbc连接器异常”这些热搜词背后是一群真正在做实时数仓的工程师。什么是Flink CDC简单说让Flink直接监听MySQL、PostgreSQL这类数据库的binlog把增删改操作变成流式事件加上Flink SQL做过滤、关联、聚合实现“业务库变更 → 数仓实时同步 → 下游消费”的端到端链路。这套东西对传统数据同步方式比如每天全量拉取或定时增量拉取是一个很大的升级原来你看到的数据普遍有半小时甚至一天的延迟现在可以做到秒级。而且Flink SQL把门槛压得很低很多不写Java的同学也能用SQL完成清洗和聚合逻辑再配合SQL Gateway之类的组件整个团队可以把流式任务当成“SQL脚本”来管理。当然Flink CDC绝不是“装上就能用”。MySQL的binlog格式要设置成ROW多个Flink作业读同一个MySQL实例时server-id不能冲突DDL变更怎么同步到下游连接池被耗尽怎么排查。热搜词里那个“flink的jdbc连接器异常”大概率就是这类问题引起的。这些我都会在第六部分专门讲。5. 真正的生产组合拳HDFS打底Spark和Flink各司其职你以为选型是“三选一”但真实生产环境里它们往往是“三合一”。5.1 离线链路HDFS/Hive打底Spark SQL处理T1很多公司的离线数仓仍然是这样的结构业务库每天定时同步到Hive ODS层再经过DWD、DWS层的清洗汇总最终输出报表和指标。这套链路里HDFS负责海量数据的低成本存储Hive Metastore负责表格元数据计算引擎换成Spark SQL调度用Azkaban或DolphinScheduler。这条链路的优势是稳、成本低、生态成熟。Spark SQL可以直接读Hive表语法和传统SQL几乎一样团队不用专门学一套新的编程模型。对于“今天出昨天的经营报表”这种需求它比任何流式计算都合适。5.2 实时链路Flink CDC进KafkaFlink SQL清洗入仓当业务开始要“今天实时看今天的营收”“异常订单立刻报警”时离线链路就不够用了。这时候需要一条实时链路常见骨架是业务库 → Flink CDC → Kafka → Flink SQL → Doris/ClickHouse/StarRocks → 应用端。Flink CDC监听binlog把变更事件打到Kafka实时链路里的Flink SQL做数据清洗、维表关联、多流Join结果写入OLAP引擎供查询和展示。Kafka在这里起到削峰填谷和消息缓冲的作用不至于让下游一波动就丢数据。为什么要搞两条链路因为离线链路成本低、适合复杂大查询在线链路成本高、适合低延迟的小查询。两条链路存储的都是同一份业务数据只是建模粒度和刷新频率不同。早期很多团队用Lambda架构就是这个思路只是当时实时引擎是Storm或Spark Streaming。现在换成Flink原理相通。5.3 什么时候不需要Flink不是所有公司都要立刻上Flink。如果实时指标只有一两个用Spark Structured Streaming定时跑个微批也能交差如果数据量没到每秒几十万事件用Canal监听binlog然后写到消息队列再让一个定时任务去消费也撑得住。判断标准很简单当你的实时任务开始涉及复杂状态、精确一次、事件时间窗口或者一个简单的实时任务消耗了大量运维精力时才是认真考虑Flink的节点。否则为了“跟上技术潮流”而引入Flink只会让团队陷入checkpoint失败、状态后端调优、版本兼容这些无底洞。6. 选型后最容易翻车的五个实操坑6.1 伪分布式不等于真分布式环境搭建只是第一步很多人在自己电脑上搭了个hadoop伪分布式或spark local模式就跑任务、调代码感觉一切顺利。可真上了集群问题全来了网络传输导致的shuffle成本、多节点间的资源竞争、节点故障带来的任务重试这些在单机上都体验不到。我的建议是学完单机后尽早接触一个至少三台机器的小集群用真实数据规模跑一跑。不需要多高的配置4核16G内存起步就够了。环境搭建本身不是目的目的是让你对“分布式的瓶颈到底在哪”有真实感知。6.2 Spark内存参数不调集群再大也救不了OOMSpark性能优化里“内存”是绕不开的坎。常见的坑是给executor配了过大的内存反而因为GC变慢拖垮任务或者driver内存太小收集结果时直接OOM更常见的是没理解堆内和堆外内存的关系导致缓存数据挤占了shuffle所需的内存。一个靠谱的起步配置是每个executor 4~8核内存16~32G同时留出20%左右的内存给系统开销动态资源分配要打开让集群按负载自动伸缩。不要抄一个网上配置就到处用要结合你的数据量做压测。然后是熟悉几个关键参数spark.executor.memory、spark.executor.cores、spark.sql.shuffle.partitions、spark.memory.offHeap.enabled这几个参数写对了大部分性能问题能缓解一多半。另外数据倾斜比内存不够更隐蔽。group by或join时某个key的数据量特别大会导致少数executor忙、其他全空闲。遇到这种问题不要急着加资源先定位是不是几个热点key在作怪该加盐加盐该广播广播效果往往立竿见影。6.3 Flink的“精确一次”是端到端承诺别被一个引擎误导Flink官方说的exactly-once是指“在Flink引擎内部基于checkpoint机制保证状态一致性”。但一条实时数据要经过Source比如读取Kafka、Flink内部计算、Sink比如写入MySQL或Kafka如果Source端不记录消费位点Sink端不实现幂等或事务性写入那数据依然可能丢或重。真实项目中我见过太多人以为“用了Flink就能精确一次”结果下游出现重复数据才回过头来排查。正确的认知是Flink负责中间状态的精确一次两端需要配合。比如Kafka Source会保存offsetSink端Kafka Connector支持两阶段提交JDBC Sink需要靠幂等写或者事务表来保证最终一致。选型时如果业务对一致性要求极严比如金额相关一定要把整条链路的一致性模型画清楚。6.4 版本兼容矩阵不看分分钟踩JDBC连接器的坑热搜词里“flink的jdbc连接器异常”这个搜索词背后几乎都是同一个故事你下载了一个Flink版本再引入一个Connector结果运行时报出ClassNotFoundException或者奇怪的序列化错误去GitHub搜了一圈发现是版本不匹配。Flink生态里有两套版本特别容易混一套是Flink主版本一套是各个Connector自己的版本。比如Flink CDC的某些版本要对应特定Flink版本MySQL CDC和PostgreSQL CDC支持的参数也各不相同。项目初始化时最稳妥的做法是先选定一个Flink主版本再去官网查这张版本兼容矩阵把Connector版本固定住整个项目统一管理依赖不要各自用各自的最新版。还有个很典型的坑是JDBC驱动冲突。Flink自带的老版本MySQL驱动和你项目里新引入的mysql-connector-java打架导致连接报错。遇到这种情况优先看是不是驱动被shade进作业jar包里或者flink的lib目录里混入了重复驱动。排查思路一般是这样先确认Connector和Flink版本兼容→再确认驱动版本一致→再检查时区、SSL参数是否匹配程序运行环境→最后看是不是连接池耗尽导致“看起来像连接异常”。按这个顺序大多数“JDBC连接器异常”都能在半小时内定位。6.5 数据倾斜和小文件两个容易被热搜词掩盖的硬骨头数据倾斜不只是Spark的问题Flink里也会出现。比如按用户维度聚合时某个大用户的事件量占了一大半导致下游算子背压Flink的窗口聚合遇到热点key也会造成某个subtask长期瓶颈。解决思路有两条先看业务上能不能拆key把一个热点key加随机后缀拆开再分别聚合不行的话再考虑加资源但治标不治本。Hadoop生态里还有个老生常谈的小文件问题。HDFS默认块大小128MB如果一个任务生成了几万个小文件NameNode内存会被元数据压垮读取性能也大幅下降。根源经常出在“分区字段太多、写入过于频繁”。治理手段包括写入前合并小文件、定期用任务做文件合并、把流式写入的目标改成对象存储再做批式落Hive。这个问题在实时链路里尤其突出Flink写Hive或写S3时如果不对文件滚动策略做设置一天下来就是几十万个碎片文件。7. 选型之外我最后想说的三条实在话7.1 入门路线可以更朴素如果你还在纠结学Hadoop还是Spark还是Flink我的建议是用最朴素的方式把链路跑通。先装一个Hadoop伪分布式跑通HDFS上传下载再装一个Spark用一份百万级的数据写几条SQL完成同比环比、分组TopN这类分析。然后把MySQL的一个订单表用Flink CDC同步到另一张表观察数据怎么实时流动。这三步做完你对“离线”和“实时”的差别会有比看十篇对比文章更深刻的体感。7.2 中小团队先离线后实时别一上来就全家桶我见过不少中小团队一上来就把Hadoop Spark Flink K8s全套铺开结果运维压力立刻超过业务收益。如果你团队只有两三个人又刚刚开始做数据平台不妨先让“离线链路跑起来”把数仓分层、调度、元数据管理这些基本功做扎实。等业务方明确提出“我要看到实时指标”时再引入Flink一上来就盯住一两个核心场景比如实时大屏或实时告警千万不要把实时链路铺得太宽。7.3 已有Hadoop资产不要轻易推翻如果你们现在已经有了一套稳定的Hadoop集群HDFS里存着好几个P的数据Hive Metastore里管着几百张表我的建议很直接别为了“跟上技术潮流”就推翻重来。HDFS可以继续当底座Spark继续跑离线批处理Flink作为新增的实时引擎接进来。你不一定要“切换技术栈”你只需要“在合适的位置新增一个组件”。最后说一个我的个人习惯。做选型评审的时候我从来不看哪个框架的GitHub Star多、哪个框架又在社区刷了屏我只问业务三个问题你的数据延迟底线是多少你的团队能养活多少套组件你的核心指标值不值得为低延迟多付出两倍运维成本。把这三个问题想明白Hadoop、Spark、Flink的答案其实已经写在你的业务里面了。