这几年带过不少刚接触大数据的同事,也面试过很多候选人,被问得最多的一句话就是:Hadoop、Spark、Flink到底该先学哪个,它们的区别又是什么。这个问题看起来很基础,但真能把三者关系讲清楚的人不多。作为一个从Hadoop 1.x时代一路用过来的人,我打算把这套技术栈的定位、实操要点和踩坑记录一次性整理出来,希望能给正在入坑或者已经踩坑的人一份比较完整的参考。
先说结论:这三个东西不是非此即彼的替代品,而是同一套大数据技术栈里分工不同的三环。Hadoop负责把数据存下来、把集群资源管起来;Spark和Flink负责把数据算起来,其中Spark强在批处理和复杂的离线分析,Flink强在流式处理和秒级响应的实时场景。理解了这个分工,后面所有的选型、学习路径和排查问题就有了解题框架。
这篇文章的内容会覆盖三套件的核心原理与定位、搭建与使用过程中最容易踩的坑(伪分布式格式化、Spark on YARN核数异常、Flink JDBC连接器报错等),以及一套典型的生产架构参考。无论你是准备面试、做大作业、还是刚接手一个集群,应该都能从中找到对应章节直接照着操作。
1. 定位拆解:为什么大数据栈绕不开这三个组件
1.1 Hadoop:最底层的“仓库”与“调度室”
Hadoop是一个生态的统称,核心组件是HDFS、YARN和MapReduce。HDFS解决的是海量文件怎么存的问题,文件被切成块(默认128MB)分布在多台机器的磁盘上,天然具备容错能力;YARN解决的是集群资源怎么分的问题,哪台机器有空闲、分配多少内存和CPU,都由它来统一调度;MapReduce是早期默认的计算引擎,把计算任务拆成Map和Reduce两个阶段来跑。
我在带新人的时候常用一个类比:把HDFS理解成一个超大仓库,货物(数据块)分散存放在很多货架上,即使某个货架塌了,备用的副本也能顶上;YARN就是仓库的管理员,谁要进货、谁要搬货,需要申请多少人力物力,都由管理员分配;MapReduce则是搬货流程,先把大件拆成小件抓起来(Map),再统一整理归位(Reduce)。
虽然MapReduce在速度和开发效率上已经被Spark和Flink远远甩开,但HDFS和YARN在今天依然是大数据底座里最常用的一层。哪怕你计算全用Spark、流处理全用Flink,数据基本还是要落到HDFS上,资源也往往还是YARN来管。这就是为什么很多招聘要求里,“熟悉Hadoop”始终没被删掉。
1.2 Spark:把批处理速度拉到极限的“加工线”
Spark最初就是为了解决MapReduce迭代计算太慢的问题而出现的。MapReduce每一步的结果都要落到HDFS磁盘上,下次再读出来,一来一回I/O开销非常严重。Spark的核心思路是尽量把中间结果留在内存里,通过DAG(有向无环图)把一连串操作串联起来,减少落盘次数。
它的计算单位是RDD(弹性分布式数据集),后来又推出了DataFrame和Dataset。RDD适合底层灵活操作,DataFrame则更适合做结构化数据分析,性能和优化空间都更好。实际项目里绝大多数人用的是DataFrame + Spark SQL的组合,把SQL能力引入进来后,普通业务人员也能快速上手做离线分析。
Spark还能兼顾流计算。它提供了Structured Streaming,本质上还是把流拆成微批来处理。和Flink相比,缺点是端到端延迟会高一些,数据到达与结果产出之间会有几百毫秒甚至几秒的间隔;优点是对接Spark生态非常顺滑,如果你本来就是做离线批处理为主,用一套技术栈顺带解决准实时需求,成本很低。
1.3 Flink:天生为流式数据设计的“流水线”
Flink和Spark最大的区别在于世界观:Spark把一切看成“一批数据”,哪怕流也是切成微批;Flink把一切看成“事件流”,数据是一条一条、永远在流动的。在这个视角下,Flink天然支持高吞吐、低延迟、精确一次(Exactly-Once)语义,而且对事件时间的处理非常成熟。
举个例子:统计过去5分钟每个接口的调用量。如果用离线批处理,只能等5分钟结束之后把日志汇总跑一次;如果用Flink,事件到达的那一刻就开始参与计算,窗口时间到了就立刻输出结果,中间不需要等待,端到端延迟可以做到毫秒到秒级。很多实时大屏、实时风控、实时推荐背后的计算引擎都是Flink。
Flink还有一个杀手锏是Flink CDC(Change Data Capture),可以直接从MySQL等数据库的binlog里读取数据变更,再同步到数仓或消息队列里。这让“数据库实时入湖入仓”的门槛大幅降低,我在后面讲工程化实现时还会再展开。
1.4 三者的分工关系与选型逻辑
所以三者的关系可以这样理解:Hadoop是底座,负责存储和资源管理;Spark是批处理的主力,适合大规模离线计算、复杂ETL、历史数据分析和机器学习预处理;Flink是流处理的主力,适合实时计算、实时数仓、实时风控、数据同步这类对延迟敏感的场景。
选型时可以从两个维度判断:一是数据到达节奏。如果是T+1、每天定时来一批,用Spark就够;如果是日志、订单、传感器指标持续不断地产生,需要尽快响应,就上Flink。二是延迟要求。秒级甚至毫秒级延迟选Flink,分钟级、小时级任务选Spark完全没问题。
这里也要提醒一点,很多团队会用Spark做“准实时”(分钟级)链路,用Flink做“真实时”(秒级)链路,两者并存是非常正常的架构,不存在谁替代谁的问题。学习的时候也不必一上来就纠结“我到底该学哪个”,先把Hadoop的存储和调度搞明白,再用Spark练手做离线分析,最后上手Flink做流处理,这条路走下来会比较顺。
2. Hadoop:存储与资源调度,最容易踩坑的底层
2.1 伪分布式搭建:从安装到Format的完整流程
很多第一次接触Hadoop的人都是在自己的笔记本上用伪分布式模式开始练手。所谓伪分布式,就是在一台机器上同时模拟NameNode、DataNode、ResourceManager、NodeManager这几个角色。虽然和生产环境的真实集群差距很大,但用来理解配置项和启动流程是非常高效的。
我建议新手一定不要直接跳过配置阶段。因为Hadoop的很多坑都出在配置上,比如core-site.xml、hdfs-site.xml、yarn-site.xml里几个关键不变量没写对,后面启动起来各种莫名其妙的问题。流程大概是这样:
- 安装JDK并配置JAVA_HOME环境变量,Hadoop 2.x/3.x依赖Java 8或Java 11;
- 下载Hadoop二进制包并解压,配置
HADOOP_HOME,把$HADOOP_HOME/bin和$HADOOP_HOME/sbin加入PATH; - 修改
etc/hadoop/hadoop-env.sh,显式指定JAVA_HOME,这一步经常被忽略; - 编辑
core-site.xml,设置fs.defaultFS为hdfs://localhost:9000,同时设置临时目录hadoop.tmp.dir,强烈建议指向一个非系统临时目录,否则重启后容易丢元数据; - 编辑
hdfs-site.xml,把副本数dfs.replication设为1,因为伪分布式只有一台DataNode; - 编辑
yarn-site.xml,配置ResourceManager和NodeManager的相关参数; - 配置
mapred-site.xml,把MapReduce的框架指定为YARN。
配置完不要急着启动,还要给NameNode和DataNode的目录授权,然后执行hdfs namenode -format进行格式化。格式化成功后,日志里会出现"In-Memory file system"和“has been successfully formatted”这样的字样。再分别执行start-dfs.sh和start-yarn.sh,用jps检查进程,看到NameNode、DataNode、ResourceManager、NodeManager几个进程都起来,基本就成功了。
2.2 启动格式化失败:新手最容易遇到的第一道坎
格式化这一步看着简单,实际踩坑率非常高。我见过最多的情况是第一次格式化成功,但因为配置改动或者误删数据,想再格式化一次,这时却报错失败了,提示NameNode进程还在运行、端口被占用,或者storage directory ... does not exist、NameNode is not formatted之类。
这里必须先说清楚一个关键机制:格式化本质上是在NameNode本地创建fsimage和edits等元数据文件,同时生成一个current/VERSION文件,里面记录了一个集群ID(clusterID)。DataNode首次启动时也会生成自己的clusterID,然后去和NameNode对齐。如果格式化时DataNode的目录里残留了旧的clusterID,两边就对不上,启动后DataNode会一直报错。
所以常见的排查顺序是:
- 先确认是否有残留进程,执行
jps看一下,有就用stop-dfs.sh停掉,或者直接kill掉相关进程; - 检查
dfs.namenode.name.dir和dfs.datanode.data.dir指向的目录是否为空,如果之前已经格式化过,先把这些目录里的旧文件清空,再重新格式化; - 确认
hadoop.tmp.dir与dfs.namenode.name.dir、dfs.datanode.data.dir之间没有配置冲突; - 清除临时文件后重新执行
hdfs namenode -format。
之前带的一个新人就卡在这里整整两天,最后发现原因很离谱:他第一次格式化用的是root用户启动,后来换成了普通用户,而NameNode数据目录还归属root,权限不足导致格式化只能写一半。这个案例提醒我,Hadoop集群里的目录所有权和用户一致性,比配置本身更容易出问题。
2.3 HA高可用集群与ZooKeeper整合
伪分布式跑通了,接下来应该往真实方向走:高可用(HA)集群。生产环境里NameNode是典型的单点故障,一旦宕机整个HDFS就不可写、不可读。所以Hadoop 2.x之后引入了NameNode HA,用两个节点组成Active/Standby,依赖ZooKeeper来做自动故障切换。
ZooKeeper在这里扮演的角色就是“协调者”:它维护一个ActiveNameNode的锁,Standby节点时刻监听,一旦Active心跳消失,ZooKeeper会通知Standby提升为Active。同时JournalNode负责同步两个NameNode的元数据,保证切换后元数据不丢失。
整合步骤大概是这样:
- 部署至少三个节点的ZooKeeper集群(奇数个),启动之后用
zkServer.sh status确认leader选举正常; - 配置
hdfs-site.xml,开启HA,配置nameservices和多个NameNode的RPC、HTTP地址; - 配置
core-site.xml,把fs.defaultFS指向hdfs://nameservice; - 配置JournalNode地址,并单独格式化ZKFC;
- 按顺序启动JournalNode、ZooKeeper、NameNode和ZKFC,再用
hdfs haadmin -getAllServiceState查看两个NameNode的状态。
这里最容易踩的坑是JournalNode启动顺序。必须先启动JournalNode,再格式化NameNode,否则会报“Unable to determine address of the journalnode for namespace”之类的错误。另外,格式化HA的NameNode时要用hdfs namenode -format,但第二个节点不能直接格式化,而是通过hdfs namenode -bootstrapStandby把元数据同步过去,这一点和单节点很不一样。
2.4 部署策略与面试高频题
聊完HA,顺便说说部署策略。现在团队里比较常见的部署方式有三种:
- 传统手动部署:所有配置文件和启动命令手动维护,适合学习和小规模场景,生产环境维护成本很高;
- Ambari/HDP:可视化安装、监控和管理Hadoop生态组件,界面友好,适合几十个节点以内的集群。Ambari虽然在很多公司已经不再被官方维护,但存量用户量还是很大的;
- Docker/K8s容器化部署:用Hadoop的Docker镜像把各角色容器化,弹性更好,适合测试环境和云原生场景,但对网络、存储有额外要求。Hadoop官方和社区有现成镜像,自己也可以基于基础镜像打一个带自定义配置的版本。
面试题方面,Hadoop最常见的问题无非是这些:HDFS写文件的过程、容错机制中的副本放置策略、一个64MB文件实际占多少存储(答案取决于块大小和副本数)、YARN中Container的内存是如何计算的、MapReduce Shuffle的过程。建议准备面试的人不要只背结论,最好能结合自己实际操作中的案例来讲,比如“我遇到NameNode格式化失败,排查过程是什么”这种真实经历,比干巴巴的背书管用得多。
3. Spark:内存计算与 On YARN 实战
3.1 Spark为什么快:从DAG到内存模型
Spark的核心快在于两点:DAG调度和内存计算。每提交一个Spark任务,Driver端会把整个作业切分成多个Stage,形成DAG。Stage内部的算子会尽量pipeline起来,减少数据在节点之间的Shuffle。能省的一次Shuffle就省掉,这是Spark性能优化的一个重要思路。
内存模型这块,Executors是真正干活的地方。每个Executor有几个核心区域:spark.executor.memory指定整个Executor可用内存,其中Reserved Memory保留给Spark内部使用;User Memory用来存用户数据结构;Storage Memory负责缓存RDD/DataFrame;Execution Memory负责Shuffle、Join、Aggregation等计算过程的临时数据。在Spark 1.6之后,Storage和Execution两块内存可以互相抢占,这样能提高利用率,但也带来一个典型问题:如果缓存的数据量太大,Execution Memory被挤占,会发生频繁的spill甚至OOM。
实际操作中,我习惯先给每个Executor 2G到8G,具体取决于单节点总内存和并发数。公式上可以参考:spark.executor.memory+spark.executor.memoryOverhead<= 节点可用内存;同时一个节点上的Executor数量乘以每个Executor的核数,不能超过节点总的vcore数。这里面的数值不是越大越好,设置太大反而会降低并行度,资源白白浪费。
3.2 Spark on YARN只分配了1个CPU核?原因和排查
“Spark on YARN CPU只能用1个”这是很多人在网上搜的高频问题,我在实际排查中也遇到过。表面上现象是:明明提交了--executor-cores 4,但看YARN上Container分配到的vcore只有1个。根本原因有两种,排查方向完全不同。
第一种是配置层面的问题。打开YARN的capacity-scheduler.xml或fair-scheduler.xml,找到yarn.scheduler.maximum-allocation-vcores,这个值如果被设成1,就表示YARN最多给每个Container一个vcore,无论spark怎么申请都会被限制住。类似的还有yarn.nodemanager.resource.cpu-vcores,这个参数描述的是每台NodeManager节点能提供的总vcore数,如果设成了1或者刚好等于1,资源就放不开。
第二种是Spark侧的计算问题。spark.executor.cores申请的是执行任务的CPU核数,但Spark在on YARN模式下发送给ResourceManager的container请求,最终受制于调度器其他参数。如果节点是1核小机器,或者集群公平调度(FAIR)策略里设置了最大核数限制,也会导致实际分配永远只有1个vcore。
排查步骤我建议这样:
- 先看YARN UI,确认NodeManager实际可用vcore是多少;
- 看
yarn-site.xml里yarn.nodemanager.resource.cpu-vcores和yarn.scheduler.maximum-allocation-vcores两个值; - 看Spark提交命令里
--executor-cores有没有生效,SPARK_CONF_DIR里有没有被别的配置覆盖; - 用
yarn application -list -appStates RUNNING找到application id,再yarn applicationattempt -list <appid>看实际请求到的container资源。
还有一个容易被忽略的点:如果CPU核数和内存不匹配,YARN可能优先保证内存而压缩核数。比如--executor-memory 4G --executor-cores 4,但NodeManager只有4G内存可用,那每个Executor只能用1个核才能保证内存充足。所以调参时一定要让内存和核数的比例匹配上节点实际资源。
3.3 集群搭建与实践中常用的调优参数
Spark集群搭建没有想象中的复杂。Standalone模式下,选一台机器做Master,其他机器做Worker,配置好SPARK_MASTER_HOST,启动后就有一个可用集群。但生产中用得更多的还是on YARN,毕竟集群资源由YARN统一管理,能避免多个计算框架抢资源的问题。搭建Spark on YARN的重点就是把spark-env.sh里的HADOOP_CONF_DIR或者YARN_CONF_DIR指到Hadoop配置目录,提交时用--master yarn --deploy-mode cluster即可。
调优参数这一块,我常用的几项:
spark.executor.memory和spark.executor.cores:决定每个Executor的计算能力,建议先按“节点总内存/(executor内存+overhead)”算并发度;spark.sql.shuffle.partitions:Spark SQL默认的Shuffle分区数,很多人从默认200开始调,可以根据数据量和executor数设置成executor数量 * 核数 * 2到4,实际效果要好一些;spark.dynamicAllocation.enabled:动态资源分配,打开之后Spark会在空闲时释放Executor、任务高峰时自动申请,对节约成本很有帮助;spark.memory.fraction:Storage与Execution合并内存比例,默认0.6,如果任务Shuffle很重,可以保持默认或者调低一点,给User Memory留更多空间。
关于“在某些Spark作业中,executor在YARN上运行时每个container只分配一个vcore”的问题,我在前面单独讲了,但这里要再补一句:很多情况其实是故意设成1vcore的。因为如果一个Executor有4个核,它跑4个并行任务,但每个任务的延迟不同,经常出现一个核一直等数据、另外几个核空转的情况。拆成多个单核Executor反而能提升整体吞吐。所以看到1vcore不要先怀疑是不是配错了,有时候是故意为之。
4. Flink:流式大数据的生产级实现
4.1 Flink Standalone集群搭建与工程化代码
Flink的部署方式分Standalone、on YARN和Kubernetes几种。Standalone简单,适合学习和测试;生产环境我见过的项目大多跑在Flink Kubernetes Operator或on YARN上,但Standalone的搭建逻辑是理解Flink架构的基础。
Standalone搭建流程不复杂:下载Flink包,在conf/flink-conf.yaml里配置jobmanager.memory.process.size和taskmanager.memory.process.size,把TaskManager的地址加入conf/workers,然后启动start-cluster.sh,Web UI默认在8081端口。需要注意的是一台机器如果只配置1个JobManager,那就是单点,生产环境建议至少三台加ZooKeeper(Flink 1.14之前)或者用Kubernetes的Multiple JobManager模式来做高可用。
工程化代码这里想多说几句。很多人写的Flink作业跑个Demo没问题,但放到生产很容易出问题,比如状态后端没有配置、checkpoint没有开、并行度写死在代码里。我现在的工程习惯是:
- 通过
StreamExecutionEnvironment.getExecutionEnvironment获取环境,便于外部提交平台统一注入参数; - 显式开启Checkpoint,间隔60秒左右,设置
CheckpointingMode.EXACTLY_ONCE,并配置state backend为RocksDB; - 所有并行度、Kafka topic、表名等配置都走外部参数或配置文件,绝不硬编码;
- 用
env.configure将Flink原生的配置项从flink-conf透传,这样提交平台和本地测试保持一致。
4.2 Flink SQL的JDBC连接器与SASL认证问题
Flink SQL现在越来越普及,它的JDBC连接器用来读写MySQL、PostgreSQL等传统数据库非常方便。但官网文档里关于连接参数的说明比较简略,实际使用中要补不少细节。
比如用CREATE TABLE定义一张MySQL维表时,连接器配置大致是connector = 'jdbc'、url、table-name、username、password这些,但还有一个隐藏坑:如果MySQL侧连接数很小,而Flink并行度又特别高,会瞬间打满数据库连接池,报“Too many connections”或者“Communication link failure”。解决办法是配置connection-pool-size和connection-pool-max-wait-time,把连接池大小控制在一个合理范围。
热词里提到的flink sql sasl sasl_plaintext,这通常和连接Kafka使用SASL/PLAIN认证有关。Flink SQL连接Kafka时,如果Kafka集群开启了SASL_PLAINTEXT,需要在properties.*里配置安全协议和认证信息。常见写法是:
CREATE TABLE kafka_source ( id BIGINT, name STRING, ts TIMESTAMP(3) ) WITH ( 'connector' = 'kafka', 'topic' = 'user_events', 'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092', 'properties.security.protocol' = 'SASL_PLAINTEXT', 'properties.sasl.mechanism' = 'PLAIN', 'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="admin123";', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json' );这里最常见的问题是JAAS配置的引号或分号写错,导致启动时报“javax.security.auth.login.LoginException: unable to load LoginModule class”。还有一点容易被忽略:properties.sasl.jaas.config里的双引号会被解析掉,写成普通字符串即可,不需要再转义成\"。
4.3 经典链路:Flink消费Kafka写入Elasticsearch
Flink消费Kafka写入ES,可以说是流计算最经典的入门案例。我平时带人演示时都会用这条链路,它几乎覆盖了Flink生产开发的所有关键点:数据源、解析、窗口聚合、Sink和状态管理。
伪代码思路大致是这样的:
- 从Kafka读取JSON字符串;
- 用Flink SQL或者DataStream API做数据清洗和转换,比如过滤掉空数据、补齐字段、做时间窗口聚合;
- 把结果写到Elasticsearch Sink。
在写ES Sink时,最值得注意的问题是主键和幂等性。ES Sink默认是UPSERT模式,需要指定主键,否则每次都是追加文档,会造成重复数据。写入模式一般选STREAMING_UPSERT或STREAMING_BULK,对ES的bulk请求要控制频率和大小,避免把ES集群打挂。
还有一个很隐蔽的坑:Flink和ES的版本兼容性。不同版本的Elasticsearch Connector对ES类型的各种行为不一致,我有一段时间用ES 7.x的beat版本,写数据时总是报MapperParsingException,后来发现是时间字段格式没统一。建议在写ES之前把时间字段统一转成ISO 8601格式字符串,并和ES里的mapping对齐。
4.4 并行度智能扩展与资源消耗最小化
热词里有一条“抛弃并行度设置:Flink智能扩展,资源消耗最小化”,这个观点我比较赞同。很多人提交Flink作业时习惯人为指定一个并行度,比如--parallelism 64,然后长期不调整。但数据量是动态的,白天业务高峰期和深夜低谷期差异巨大,固定并行度要么不够用,要么白白浪费资源。
Flink自带的自动扩缩容能力在1.16之后有了明显增强,尤其是自适应调度器(Adaptive Scheduler)。启用后,作业会根据当前TaskManager资源和积压情况动态调整执行并行度。配合flink-conf.yaml里的这些配置:
jobmanager.scheduler: adaptive restart-strategy: fixed-delay parallelism.default: auto但实际上很多公司并不只依赖Flink自带能力,而是借助外部的Flink Kubernetes Operator来做水平扩缩容。Operator监控TaskManager的CPU和背压指标,超过阈值就自动增加TaskManager副本,空闲时就缩容到最小规模。这个思路比人为拍脑袋定并行度要靠谱得多,也更能节约成本。
我在实际项目里的一般做法是:开发阶段明确并行度,方便排错;上线阶段交给调度器或Operator动态管理,只设置最小和最大并行度范围。这样做之后,资源使用率比原来明显稳了,高峰期没有因为并行度不够导致堆积,低谷期也不会有一堆空闲TaskManager在烧内存。
5. 生产架构与避坑清单
5.1 一套典型的大数据生产架构
到这里,可以把三件套串起来看一个完整的生产链路。我近期参与的一个数据平台项目,简化后的架构大概是这样的:
- 业务系统的MySQL、PGSQL数据库产生变更数据,通过Flink CDC实时捕获;
- Flink CDC把变更数据写入Kafka,作为消息缓冲和削峰;
- Flink消费Kafka数据,做实时清洗、窗口聚合,写入Elasticsearch供大屏实时查询,同时写入HDFS做离线归档;
- 离线任务每天凌晨用Spark读取HDFS上的全量数据,做复杂的ETL和报表计算,结果写回Hive或者MySQL;
- YARN作为Spark和一部分Flink任务的统一资源调度平台;
- 所有元数据和调度任务由统一调度平台管理,异常自动重跑和告警。
这套架构的好处是“流批分离但存储共享”。实时链路秒级产出结果,离线链路讲成本、讲稳定,两条链路的数据都落在同一个HDFS/数仓体系里,不会出现数据孤岛。坏处是要维护的技术栈比较多,对团队的运维要求高。如果是小团队,可以先用Spark做准实时,再逐步引入Flink,不用一步到位。
顺带说一句,大数据技术栈的应用也不只在互联网行业。比如卫星遥感数据处理、轨道动态可视化和覆盖分析这类科学计算场景,同样会依赖HDFS存原始影像数据、Spark做批量几何计算、Flink做实时状态监测。技术和业务场景结合之后,能发挥的空间比很多人想象中大得多。
5.2 常见问题速查表
把这三套件在实战中最容易遇到的问题整理成一张表,方便直接对照排查:
| 场景 | 典型现象 | 常见原因 | 解决办法 |
|---|---|---|---|
| Hadoop格式化 | NameNode ... has been successfully formatted之后启动仍失败 | 临时目录或数据目录权限不一致 | 清空数据目录,确认启动用户一致后重新格式化 |
| Hadoop HA | 第二个NameNode无法成为Standby | 未执行bootstrapStandby | 先启动JournalNode,再hdfs namenode -bootstrapStandby |
| Spark on YARN | 每个Container只有1个vcore | 调度器maximum限制或节点CPU配置过低 | 检查yarn.scheduler.maximum-allocation-vcores与nodemanager.resource.cpu-vcores |
| Spark OOM | Executor频繁GC或抛出OutOfMemory | Execution Memory不足或spill过多 | 调大executor内存,降低shuffle分区倾斜,或调整spark.memory.fraction |
| Kafka SASL | LoginException: unable to load LoginModule | JAAS配置引号或类名写错 | 检查properties.sasl.jaas.config格式,确认类路径 |
| ES Sink数据重复 | ES中出现重复文档 | 未指定主键或没用UPSERT模式 | Sink中指定主键,设置STREAMING_UPSERT |
| Flink背压 | Web UI中作业持续High背压 | 下游写入慢或并行度不足 | 调大下游并行度,优化ES bulk大小,检查Kafka消费位点 |
| HDFS空间不足 | NameNode is in SafeMode | 容量告警进入安全模式 | 清理数据或扩容,确认剩余的容量比例解除SafeMode |
这张表里的每一行都是我或团队同事实际碰到过的。值得强调的是,排查时不要一上来就往深了猜,先看日志里的关键异常,再逐层缩小范围,通常效率最高。
5.3 学习路线与面试建议
如果你是刚入门,我给的建议很直接:先跑通Hadoop伪分布式,理解HDFS和YARN的角色;再用Spark做一两个数据分析案例,比如统计日志中的用户访问量Top10;然后上手Flink,把Kafka到ES这条链路跑起来。这个过程如果能完整走一遍,应付课程设计、毕业设计,甚至初级大数据开发岗的面试都够了。如果你学的就是数据科学与大数据技术专业,更建议把毕设题目落在这套链路上,既有工程深度,又有数据场景,答辩时也好讲故事。
面试时,三个组件的问题通常会这么问:
- Hadoop:NameNode和DataNode的职责?写文件流程?副本怎么放?
- Spark:RDD、DataFrame的区别?Shuffle是什么?如何减少Shuffle?
- Flink:Flink和Spark Streaming的对比?Checkpoint机制?状态存在哪里?精确一次是怎么实现的?
回答这些问题的关键不是背概念,而是有自己的理解。比如问到“Flink和Spark Streaming的区别”,你可以先说思想差异(流vs微批),再结合自己实际跑过的延迟数据或者窗口计算案例,说明为什么某个场景选Flink更合适。面试官想听到的不是标准答案,而是你有没有真正上手过。
最后再说一点个人感受。我从最开始只会在课程里跑MapReduce,到现在能完整搭起一条流批一体的数据链路,最大的体会是:大数据技术栈的学习没有捷径,但一定有高效路径。千万不要陷入“只看文档不跑环境”的陷阱,很多配置项和报错信息只有亲手敲一遍才能记牢。如果这篇文章能帮你少走一点弯路,哪怕只是解决了格式化失败或者vcore只分配1个的问题,我觉得今天这一篇就没有白写。
技术上也是一样,Spark和Flink的版本迭代非常快,新特性和新Bug都在冒。入行之后请保持看官方文档和Release Notes的习惯,别一直停留在教程里的老版本参数上。祝各位早日搭出属于自己的那套集群,下次再被问到Hadoop、Spark、Flink有什么区别时,能用自己的话讲得明明白白。