简介:本资源为大数据面试高频考点整理合集,面向准备Java后端、大数据开发及大厂技术岗位的求职者。文档系统梳理CAP理论、BASE理论、ACID与BASE对比、2PC两阶段提交、MySQL Replication主从复制及数据库恢复与备份策略等核心主题,既讲清概念,也点明分布式场景下的取舍思路与典型应用,可帮助读者建立完整知识框架,并在面试中更有条理地作答。资源为1个pdf文件,整体约38.52MB,内容密度较高,适合集中浏览与考前速查。目前已有258人学习下载,可作为大数据方向面试冲刺阶段的高性价比补充资料。
1. 面试题 PDF 的价值不在答案,而在标准答案背后的框架
《大数据面试题(含答案).pdf》这类资料在大数据求职圈里流传极广,但它最大的坑恰恰是“含答案”。很多人把 PDF 里的答案背下来就去面试,结果被追问两轮就露馅。原因很简单:面试官看的不是你能不能说出 Flink 的 Checkpoint 机制是什么,而是看你有没有自己的理解层次——能不能从机制讲到适用场景,再从场景反推参数配置,最后落到线上踩过的坑。答案背得再好,也只是复述,不是理解。
这份 PDF 真正有用的地方,是它把散落的知识点按模块收敛成了一个相对完整的复习框架:Hadoop 生态、Hive 数仓、Spark 内存计算、Flink 实时计算、Kafka 消息队列、数仓建模、数据倾斜调优、系统设计场景题。一个五年以上经验的工程师拿到这份文档,正确的用法不是背诵,而是对照目录做“自我提问”:每个知识点,我能不能讲清原理、适用边界、常见误区和调优手段?能讲清,说明这个模块过关了;讲不清,就去补那一块的底层源码或实测数据。本文就按这份 PDF 最常见的组织方式,帮你把核心模块的答题骨架拆开,每个模块给出可直接套用的回答结构、参数示例和追问方向,让你面对面试官时不是“背答案”,而是“讲方案”。
2. 离线数仓面试题:从数据分层到数据倾斜的标准应答框架
离线数仓是大数据面试的必考模块,Hive 和 Spark SQL 相关题目几乎每场必出。PDF 里这类题的答案通常很长,但记住一个原则:面试官想听的永远是你“为什么这么设计”,而不是“做了什么”。下面按高频考点拆解。
2.1 数仓分层设计题的答题主线:ODS、DWD、DWS、ADS 的职责边界
“你们数仓怎么分层的?”是开场高频题。标准答案是四层:ODS 原始数据层、DWD 明细数据层、DWS 汇总数据层、ADS 应用数据层。但光说四层不够,要能讲清每一层做什么、解决什么问题、数据是怎么流转的。
ODS 层保持原样接入,不做清洗,只做压缩和分区。分区通常按天,数据格式可以用 Parquet 或 ORC,压缩用 Snappy 或 ZSTD。DWD 层做清洗、脱敏、维度退化、拉宽,这一层最核心的工作是“降维”:把重复的维度字段直接拉进事实表,避免下游多次关联。DWS 层按主题做轻度汇总,比如按用户、按商品、按日期粒度,输出多维度聚合指标。ADS 层面向具体业务报表,一个报表一张表,字段就是指标所见即所得。
这里有个容易丢分的追问点:为什么 DWD 层要做维度退化?答案不是“方便查询”这么简单,而是为了减少 Join 次数。数仓计算成本的大头在 Shuffle,而 Shuffle 的主要来源就是 Join。维度退化把常用的、变化不频繁的维度字段直接存储在事实表中,代价是存储冗余,换来的是查询时避免大量 Join,用空间换时间。SYNC 到下游的数据服务层时,也能直接支持宽表查询,不需要实时关联维表。
另一个高频追问是:DWS 和 ADS 的区别到底在哪?简洁的区分是 DWS 是主题粒度汇总,ADS 是业务需求专属,前者的粒度是稳定的、可复用的,后者的粒度随业务变化。很多面试者把两者混为一谈,“我们 DWS 就是直接出报表的”——这句话一说出来,面试官就知道你没做过真正的分层数仓,因为严格的分层里,DWS 是服务多业务的,ADS 才是单业务的。
2.2 Hive 数据倾斜题的答题结构:产生原因、定位方法和三类解决方案
数据倾斜是离线计算里最常考的实战问题。PDF 里的答案一般会列十几种情况,但核心主线只有一条:“相同 key 的数据量远超其他 key,导致单个 reducer 处理时间过长”。定位方法也不复杂:Hive 里看 Map 阶段完成后的 Reduce 耗时分布,如果某一个 Reduce 跑了很久,其他 Reduce 已结束,基本就是倾斜了。Spark 里看 Stage 的 Task 耗时,某个 Task 的数据量比其他 Task 大好几个量级,也是典型倾斜。
解决倾斜,按层递进给答案,面试官会更满意:
第一层是语法手段。最常见的:
- 因空值导致倾斜,把空值 key 加随机前缀,打散到多个 Reduce。SQL 写法示例:
-- 原写法:country 字段大量为空,导致空值 key 集中在一个 reduce SELECT country, count(*) AS cnt FROM user_actions GROUP BY country; -- 优化写法:空值/null 加随机前缀打散 SELECT CASE WHEN country IS NULL OR country = '' THEN CONCAT('rand_', FLOOR(RAND() * 100)) ELSE country END AS country, count(*) AS cnt FROM user_actions GROUP BY CASE WHEN country IS NULL OR country = '' THEN CONCAT('rand_', FLOOR(RAND() * 100)) ELSE country END;逻辑说明:第一种写法在 GROUP BY 时,空字符串和 NULL 会被归到同一个组,如果该组数据巨大,就造成单 Reducer 倾斜。第二种写法给空值人为拼接了 0~99 的随机前缀,空值被均匀分布到 100 个组中,最后结果只需要把前缀去掉再聚合一次即可。参数上要注意前缀范围取多少,RAND() 的范围和并发数匹配,不要全表随机导致后续二次聚合也倾斜。
- 因 Join 小表 key 集中导致的倾斜,用 Map Join 广播小表:
-- 开启自动 Map Join 的 Hive 参数设置 SET hive.auto.convert.join=true; SET hive.mapjoin.smalltable.filesize=25000000; SET hive.auto.convert.join.noconditionaltask.size=200000000; -- 具名 Map Join 写法 SELECT /*+ MAPJOIN(dim_user) */ a.user_id, b.user_name FROM fact_orders a JOIN dim_user b ON a.user_id = b.user_id;参数说明:hive.auto.convert.join开启自动转换,hive.mapjoin.smalltable.filesize设定小表阈值,默认 25MB 左右,hive.auto.convert.join.noconditionaltask.size是整体 Map Join 输入大小限制。在生产环境里,维表几十 MB 是常见情况,这两个参数是高频调节点。Map Join 把小表加载进每个 Mapper 的内存,读取大表时直接在内存中查小表数据,避免了 Reduce 阶段的 Join。注意点:小表不能太大,否则 OOM;任务高峰期如果内存紧张,可能还需要配合调大容器内存。
- 因 Join 的关联字段值为空导致倾斜,可以先过滤空值,再做关联计算。逻辑不复杂,但很多没实战经验的人会漏掉这一步,你的答案里主动提这一条,会显得踩过坑。
第二层是参数手段。主要是开启 Hive 的倾斜负载均衡机制:
set hive.groupby.skewindata=true;这个参数开启后,Hive 会自动对 GROUP BY 结果做两轮 MR:第一轮将 key 加随机前缀打散,进行局部聚合;第二轮去掉前缀,按真实 key 做最终聚合。代价是额外的 MapReduce 任务开销,所以只在确认倾斜时才开,默认关闭才是合理的生产状态。这个细节说清楚,比直接说“遇到倾斜就打开这个参数”高级得多——面试官想听的就是你懂“什么时候不该开”。
第三层是业务层面的治本。比如热点 key 本身是某种极端值(如超级大卖家的订单量占了全表 80%),考虑把这些高频 key 拆出来单独处理,用 “热点隔离” 思路:把大 key 过滤出来单独做一轮聚合,再与正常 key 的结果 UNION ALL 合并。伪代码示意:
-- 大 key(比如 seller_id = 'TSLA001')单独聚合 SELECT seller_id, sum(order_amt) AS amt FROM fact_sales WHERE seller_id = 'TSLA001' GROUP BY seller_id UNION ALL -- 其余 key 正常聚合 SELECT seller_id, sum(order_amt) AS amt FROM fact_sales WHERE seller_id != 'TSLA001' GROUP BY seller_id;注意,真实场景里大 key 可能不止一个,建议写成配置文件或参数表,不要让 SQL 代码里硬编码热点值。否则换了一个业务线,一段代码就要改一次。
2.3 大促场景压缩、列存与分区表的面试题应答
文件存储格式是大数据面试里“看起来简单、实际上区分度很高”的一类题。PDF 里通常出现 “你们生产环境为什么用 ORC 而不是 Parquet?” “分区字段选什么类型?” 这类问题。
回答“为什么用 ORC”,不要只背名词“列式存储压缩比高”,要讲数据落盘的原理。ORC 有以下四个特征可以逐条展开:
- 列式存储天然对分析型查询友好,只读取查询涉及的列,IO 量小;
- 支持轻量级索引(Row Group Index 和 Bloom Filter Index),谓词下推时能跳过大量不需要读的数据块。谓词下推是什么?就是 SQL 中 WHERE 条件下推到存储层,在读文件时就过滤掉不符合条件的数据块,而不是全表读出来再过滤。这一步是 ORC 比 Parquet 在 Hive 场景更好用最重要的原因;
- ORC 有自带的压缩方式 ZLIB 和 SNAPPY,压缩比高,尤其对数值型的列效果明显。生产上常用 ZLIB 配合 ORC,压缩率高,解压速度在 Hive 场景可接受;
- 文件结构是 Stripe 级别,每个 Stripe 内部包含索引数据和列数据,读取时可以按 Stripe 粒度切分,利于并发查询。
但要注意,如果面试官问的是 Spark 场景,答案就得反过来:“Spark 里我更常用 Parquet”。原因是 Spark 的 Catalyst 优化器对 Parquet 的原生支持更好,Parquet 的谓词下推和列裁剪在 Spark 里执行效率更高,且 Parquet 与 Spark 的 datasource v2 接口配合更稳定。这里的高分点在于:你能分场景回答,而不是一把尺子量到底。存储格式不是“哪个更好”,而是“在哪个引擎的哪个版本下更合身”。
分区表的面试题必背一条:分区字段不能是数据中粒度太细、取值太多的字段。如果你按 user_id 分区,一个用户一条数据,每个分区就一个文件,元数据量大到 NameNode / Metastore 压力直接爆掉;同时下游扫描分区时频繁列目录,效率极低。合理的分区粒度是按天、按小时,或者按业务线 ID(少量固定值)。另外要记得,分区字段的类型也有讲究:日期分区用 STRING 类型dt='2026-05-20',比用 DATE 类型更好过滤。因为 Hive 对 STRING 类型的等值过滤走了静态分区裁剪,而 DATE 类型在某些版本下无法直接走分区裁剪,全表扫描就发生了。
3. 实时链路面试题:Flink Checkpoint、精确一次与背压的答题结构
实时计算在现在的大数据岗位中占比越来越高,Flink 相关题目是 PDF 里篇幅最大的模块之一。答案框架核心只有三块:状态与 Checkpoint、端到端精确一次(Exactly-Once)、背压的处理与定位。面试官从这三块各抽一道,基本就能判断你的实时功底。
3.1 Flink Checkpoint 机制与两个关键参数:interval 和 timeout
“讲讲 Flink 的 Checkpoint 机制” 是必问题。答题结构建议按“存储、触发、对齐、恢复”四步走,说清楚 Chandy-Lamport 算法的核心。
简单说,Checkpoint 是 Flink 定时把算子状态和源端 offset 一起做快照,保存到外部存储(HDFS 或 OSS)。它的核心是 Barrier 机制:Source 算子周期性插入 barrier 到数据流里,barrier 随数据一起在下游算子之间传递。当某个算子所有输入流都收到了对齐的 barrier 时,就触发本算子状态快照。快照完成后,向 JobManager 确认。当所有算子都确认后,这次 Checkpoint 就是 Completed。
作答时带上实际配置,会让答案更有网感:
execution.checkpointing.interval: 60s execution.checkpointing.timeout: 10min execution.checkpointing.min-pause: 30s execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION state.backend.type: rocksdb参数说明:
interval是 Checkpoint 触发间隔,生产上常见 1~5 分钟。间隔太短会导致状态频繁落盘,影响吞吐;太长会放大故障恢复时间(恢复时要回放到最近一次快照)。我一般先设 60s,观察反压和 Checkpoint 耗时再调;timeout是单次 Checkpoint 超时时间,默认 10 分钟。如果数据量大或状态大,单次 Checkpoint 超过 timeout 就会失败,两个连续失败会导致 Job 重启。线上会遇到问题:Checkpoint 总在 timeout 附近失败,调大 timeout 后有改善,但更治本的方向是减小状态或优化 RocksDB 存储;min-pause是最小间隔,保证上一次 Checkpoint 结束和下一次开始之间至少隔多久,防止频繁快照挤占业务处理时间;externalized-checkpoint-retention设置任务 cancel 后 Checkpoint 是否保留。生产上一定要设置成 RETAIN_ON_CANCELLATION,否则任务重启后状态丢失,需要从 Kafka 最早位点重读,极其危险;state.backend.type是状态后端。大状态场景用 RocksDB,状态超过 JVM Heap 可用内存时 RocksDB 是唯一选择,它利用磁盘存储状态,代价是序列化开销和读写性能打折。
关键追问是:Checkpoint 失败的常见原因有哪些?三个方向:第一,状态太大导致快照超时,典型场景是 RocksDB 的 state 文件导入到 HDFS 过慢;第二,网络或 HDFS 抖动导致存储不可用;第三,任务本身存在反压,barrier 无法及时传递,Checkpoint 一直处于对齐等待状态。第三个原因是面试加分项,因为很多人背了 Chandy-Lamport 但不知道背压会让 barrier 流动变慢,两者是联动的。
3.2 端到端 Exactly-Once 的实现组合:两阶段提交 + Kafka Source/Sink
“Flink 怎么做端到端精确一次?” 这是实时方向最难也最有区分度的问题。标准答案是两阶段提交(Two-Phase Commit,2PC),但要讲出细节才有价值。
Flink 的精确一次语义分三层:Source 端、内部状态、Sink 端。Source 端依赖 Kafka offset 的保存,Flink 把 Kafka consumer 的 offset 存放在 Checkpoint 状态里,恢复时从状态中的 offset 重新消费,这样就保证了 Source 端的精确一次。但注意这里有一个前提:数据在 Kafka 里没有被自动提交 offset 机制干扰。生产环境要把enable.auto.commit设为 false,由 Flink 来管理 offset。
内部状态层:因为状态本身包含在 Checkpoint 里,Flink 内部状态天然精确一次。关键是 Sink 端落外部系统时,如何保证不重复写入,Flink 的 Kafka Producer 使用的是 2PC 协议。
Kafka 2PC 的过程展开是这样的:Flink 的 FlinkKafkaProducer 在 Checkpoint 开始时预提交,数据写入 Kafka 的事务缓冲区;Checkpoint 完成时真正提交事务。如果任务中途失败,Kafka 里未提交的事务数据自动回滚,下游消费者端设置isolation.level=read_committed才能读到已提交的数据,否则可能读到脏数据。实现这个机制需要 Kafka 版本支持事务(0.11 及以上),同时需要唯一的 transactionId 前缀(通过transaction.timeout.ms配置事务超时)。生产配置示例:
CREATE TABLE kafka_sink ( user_id STRING, event_time TIMESTAMP(3) ) WITH ( 'connector' = 'kafka', 'topic' = 'dwd_user_events', 'properties.bootstrap.servers' = 'kafka-1:9092,kafka-2:9092', 'properties.transaction.timeout.ms' = '900000', 'properties.isolation.level' = 'read_committed', 'format' = 'json' );参数说明:transaction.timeout.ms必须小于等于 Broker 端的transaction.max.timeout.ms(默认 15 分钟),否则 Producer 会报错。如果 Flink Checkpoint 间隔是 5 分钟,这个超时设置到 10 分钟左右比较合理。isolation.level=read_committed让 Flink SQL 读取 Kafka 时只消费已经 commit 的事务消息,确保端到端链路一致。
要命的坑是:Kafka 事务的过期时间如果设置太短,事务还没提交就被 Broker 中止,Flink 会一直重试,表现为“找不到 transaction”的异常。排查思路:先看 Flink 日志里的 transactionId,再到 Kafka 服务端看事务状态;常见做法是直接调大transaction.timeout.ms。
3.3 背压问题的定位方式:Web UI 指标和 Kafka 消费 Lag 的联动判断
背压不只是一个问题,更是一个“现象”。面试官会这么问:“收到的告警说有个 Flink 作业消费 Kafka Lag 持续上涨,你怎么排查?”
优秀的回答一定不是直接说“调并行度”,而是给出定位链路。第一步看 Flink Web UI 的 Backpressure 选项卡。Flink 1.13+ 版本背压监控会给出三个状态:OK(空闲)、LOW(低压力)、HIGH(高压力)。重点关注 HIGH 出现的算子。
第二步是区分“Source 端压力”还是“下游处理压力”。看 Source 算子的 Inbound 数据量是否明显下降,或者 Source 算子的 busy 时间占比。所谓 busy 时间占比,就是指该算子线程处理数据和等待数据的比例,如果接近 100%,说明 Source 一直在满负荷拉取数据。同时看 Sink 算子的 Busy 比例低而 Source 的 Busy 比例高,基本可以判定是“全链路反压”——问题不在某一算子,而是整条链路的吞吐不够。
第三步是细查。常见的根因:
- 某个算子存在热点,比如 KeyBy 后某 key 的数据量远大于其他 key,导致该子任务处理慢;
- 外部系统响应变慢,比如 Sink 写到 MySQL/Redis 时连接被打满,Sink 算子成为瓶颈。注意,Sink 变慢并不一定在 Sink 算子显示 HIGH,因为 Flink 的背压传播是自下而上的,最直观的表现反而是 Source 或上游算子变 HIGH;
- RocksDB 状态读写变慢,可能是磁盘 IO 问题,也可能是状态过大导致序列化频繁。
排查命令层面,从 Kafka 端看消费组 Lag:
kafka-consumer-groups.sh --bootstrap-server kafka-1:9092 --describe --group flink_etl_group这个命令输出每个分区当前消费到的 offset、log-end offset 和 LAG。如果某个分区的 LAG 明显大于其他分区,说明该分区的消费子任务存在处理倾斜。比如 12 个分区,其中 3 个分区的 LAG 持续上涨,其他分区 Lag 为 0,这基本可以判定是 KeyBy 分布不均,而不是整体吞吐不足。此时应该审视 key 的选择逻辑,而不是盲目加并行度——加了并行度,分到热 key 的那一个子任务依然还是会 Lag。
4. 大数据系统设计题与 SQL 场景题的“标准答案模板”
PDF 里常见系统设计题和 SQL 场景题。这两类题最怕的就是“不知道从哪里开始答”。系统设计题有模板可以套,SQL 场景题则有固定套路。
4.1 Lambda 架构与 Kappa 架构的选型题:什么时候不得不用 Lambda
“数仓实时和离线两套引擎,你选 Lambda 还是 Kappa?” 这个题看 PDF 是看不会的,要有实际工程经验才能答得让人信服。
Lambda 架构是离线批处理和实时流处理并行,两套代码、两套结果再合并。优点是准确,批处理链路可以修正实时链路的各种问题;缺点是需要维护两套代码,离线结果和实时结果经常对不上,合并逻辑复杂。
Kappa 架构是只维护一套实时计算链路,需要回溯时用 Kafka 重放数据来重新计算结果。节省了开发成本,但问题是:数据量大了以后,从 Kafka 重放全量数据的耗时和成本都不可控,而且 Kafka 的保存时间通常只有几天到几周,如果历史数据超过保留时间,Kappa 就无能为力。
我给面试官的答案通常是:默认用 Kappa,但业务要求“实时+离线双轨验证”的场景,尤其是财务对账、军品价格管控类对准确性要求极高且需要审计留痕的场景,被迫上用 Lambda。你可以表述为:业务对准确率要求到小数点后几位,且实时计算结果需要和 T+1 离线结果做交叉验证,这时候 Kappa 很难让业务方接受。面试时主动说出“Kappa 局限性的核心是数据回溯能力有限”,比你背出两个概念得分高得多。
设计题答题模板里有三个固定环节:数据量估算、写入链路设计、查询链路设计。数据量先估算:日活 1000 万,每条行为日志 1KB,一天原始日志约 10TB,压缩后约 1TB。写入链路设计:日志采集用 Flume 或 Filebeat,写入 Kafka,Topic 按业务线拆,分区数按下游消费并行度定。查询链路设计:实时查询用 Doris/ClickHouse,明细用 HBase,即席分析用 Spark SQL 跑在 Hive 表上。这三段讲完,设计题基本不低于基准分。
4.2 大数据 SQL 场景题:留存率、漏斗转化和会话时长窗口的写法
SQL 场景题是必须手写代码的部分。PDF 里最常出现的三道题:用户留存率、漏斗转化率、会话时长(Session 划分)。每道题都要有明确的 SQL 模板和边界条件说明。
第一题:计算某天的次日留存率。常见写法如下:
WITH user_login AS ( SELECT user_id, dt FROM user_action_log WHERE dt BETWEEN '2026-05-01' AND '2026-05-02' GROUP BY user_id, dt ) SELECT a.dt AS login_date, COUNT(DISTINCT a.user_id) AS dau, COUNT(DISTINCT CASE WHEN b.dt = date_add(a.dt, 1) THEN b.user_id END) AS retained_users, ROUND( COUNT(DISTINCT CASE WHEN b.dt = date_add(a.dt, 1) THEN b.user_id END) / COUNT(DISTINCT a.user_id), 4 ) AS retention_rate FROM user_login a LEFT JOIN user_login b ON a.user_id = b.user_id AND b.dt = date_add(a.dt, 1) GROUP BY a.dt;代码说明:先限定日期范围,把活跃用户去重取出,再按用户 ID 关联次日是否活跃。DATE_ADD(a.dt, 1)计算次日日期。这个 SQL 能跑通,但数据量大时要优化:应该按日期分区裁剪,如果只算一天留存,a表只取2026-05-01,b表只取2026-05-02,加 WHERE 条件缩小扫描范围,避免上面 SQL 扫描两天后做GROUP BY a.dt的全表聚合。生产环境必须这么写,否则几亿行数据会把资源打爆。
第二题:漏斗转化率。典型场景是统计“曝光 → 点击 → 下单 → 支付”的每一层转化。这里有一个常见错误:直接对多张表做 Join 导致数据膨胀和重复计数。正确思路是先把用户行为按事件类型做行转列聚合:
SELECT user_id, MAX(CASE WHEN event_type = 'expose' THEN 1 ELSE 0 END) AS is_expose, MAX(CASE WHEN event_type = 'click' THEN 1 ELSE 0 END) AS is_click, MAX(CASE WHEN event_type = 'order' THEN 1 ELSE 0 END) AS is_order, MAX(CASE WHEN event_type = 'pay' THEN 1 ELSE 0 END) AS is_pay FROM user_event_log WHERE dt = '2026-05-20' GROUP BY user_id;这种写法用MAX(CASE WHEN...)把每个用户是否有某类行为转成一列,避免了多表 Join 和重复计数。得到这个中间结果后,后续漏斗计算就是各列逐级求SUM并依次计算比值:
SELECT COUNT(*) AS total_users, SUM(is_expose) AS step_expose, SUM(is_click) AS step_click, SUM(is_order) AS step_order, ROUND(SUM(is_click) / SUM(is_expose), 4) AS click_rate, ROUND(SUM(is_order) / SUM(is_click), 4) AS order_rate FROM ( -- 上面那个用户级聚合结果作为子查询 ) t;漏斗题的关键考察点不是 SQL 语法,而是你是否能识别出“同一用户在多阶段的行为去重”这个陷阱。基础的错误是 JOIN 事件表时用户有两条下单记录,就计成了两次转化。这里用MAX()去重后,每个用户每个事件只保留 0/1 标记,从根上避免了膨胀。
第三题:会话划分。这是最容易考出区分度的题。SQL 里做 Session 划分的核心是“计算相邻事件时间差,超过阈值则新起一个会话”。这要用到窗口函数LAG()取上一行时间,判断差值后打标,再对打标结果做累加求和:
WITH user_events AS ( SELECT user_id, event_time, LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time) AS prev_time FROM user_action_log WHERE dt = '2026-05-20' ), session_flags AS ( SELECT user_id, event_time, CASE WHEN prev_time IS NULL OR (unix_timestamp(event_time) - unix_timestamp(prev_time)) > 1800 THEN 1 ELSE 0 END AS is_new_session FROM user_events ) SELECT user_id, event_time, SUM(is_new_session) OVER (PARTITION BY user_id ORDER BY event_time) AS session_id, COUNT(*) OVER (PARTITION BY user_id, SUM(is_new_session) OVER (PARTITION BY user_id ORDER BY event_time)) AS session_event_cnt FROM session_flags;代码说明:LAG()函数按用户分区、按时间排序,取到当前事件的上一个事件时间。unix_timestamp()把时间字符串转成秒。如果当前事件与上一事件间隔超过 1800 秒(30 分钟),is_new_session标记为 1,否则为 0。最后用SUM() OVER()对标记逐行累加,得到一个随事件递增的会话编号。同一user_id同一session_id内的所有事件就构成一个会话。
这个 SQL 的低级错误是忘记按PARTITION BY user_id分区,导致不同用户的事件被混在一起算时间差,结果必然错误。追问点在于:如果事件量一天一个人有上万条,这个 SQL 是全量排序,性能很差,生产上应该按用户维度先分桶或先用 Flink 做会话拼接后再落 Hive 表。能在 SQL 答案后主动补充这句,就是有生产经验。
5. 源码级追问与调优参数清单:HashMap 问题、Shuffle 调优和 Lakehouse
到了这一章,PDF 里的题目更多是“为什么这样设计”的底层逻辑和具体的参数配置,属于面试中拉开差距的部分。
5.1 从 8 到 6:HashMap 链表转红黑树的阈值选择逻辑
经典问题:“HashMap 什么时候转红黑树?为什么阈值是 8?” 答案:“链表长度超过 8 且数组长度大于等于 64 时转红黑树,链表长度降到 6 时转回链表。” 追问:“为什么转回链表是 6 而不是 7?”
答案的核心是避免频繁切换。如果阈值都是 7,当链表长度在 7 附近波动时,会反复在链表和红黑树之间切换,而树化和反树化本身是有性能开销的(需要创建树节点、重新组织结构)。用 8 和 6 中间留出 1 的缓冲区间,让切换不要那么频繁。这个设计在 Concurrency 场景里也会被追问,但大数据方向的面试更关心的是你在写 UDF 或内存计算时是否理解 Map 结构的性能边界。比如你在 Spark 里用mapPartitions去做分组聚合,如果单个分区的数据倾斜导致某个 key 的链表过长,HashMap 退化成链表遍历,时间复杂度从 O(1) 变成 O(n)。这也能接回第 2 章的数据倾斜——内存中的数据倾斜也会导致 HashMap 性能退化。
5.2 Shuffle 调优参数对照表:离线批处理与实时流处理的通用手段
大数据面试题的调优题主要围绕 Shuffle。Hive 和 Spark 各有关键参数,但思路一致:减少落盘、减少网络传输、增大并行度。
| 模块 | 参数 | 推荐值 | 作用 |
|---|---|---|---|
| Spark Shuffle | spark.shuffle.file.buffer | 32k~64k | 每个 Shuffle 输出文件的内存缓冲大小,太小会增加磁盘 IO 次数 |
| Spark Shuffle | spark.reducer.maxSizeInFlight | 48m~96m | Reducer 同时拉取的数据量上限,太大会导致网络阻塞 |
| Spark 动态资源 | spark.dynamicAllocation.enabled | true | 开启 Executor 动态伸缩,避免高峰资源不足 |
| Hive Reduce | hive.exec.reducers.bytes.per.reducer | 256MB~1GB | 每个 Reducer 处理的数据字节数,决定 Reduce 数量 |
| Hive Map Join | hive.mapjoin.smalltable.filesize | 10MB~25MB | 小于该值自动转化为 Map Join |
参数列举之后一定要有使用逻辑:先判断瓶颈在 CPU、磁盘 IO 还是网络;再确定是增大并行度、减少 Shuffle 数据量、还是优化 Join 策略。比如 Spark 中大量小文件导致 Task 数过多,可以设置spark.sql.shuffle.partitions从默认 200 降到合理值(视数据量而定),但有 Join 或聚合时不能盲目降,会影响并行度。这里的“合理值”我没有编造具体数字,而是要你去实测观察,面试时按你们的集群规模和任务耗时来分析才是正确答法。
5.3 从数据湖到 Lakehouse 的必答观点:为什么 Spark 3 和 Flink 都在兼容 Iceberg
数据湖和湖仓一体(Lakehouse)是近两年大数据面试里必然出现的方向。PDF 里通常出现“Iceberg 和 Hudi 的区别”“什么是 Lakehouse”这类问题,高频但容易答空。
回答要点是“用数据湖的底层存储能力支撑数仓的分析能力”。数据湖(Data Lake)解决了数据“存得下”的问题,用开放格式(Parquet/ORC)存储所有原始数据,但早期的问题是“查不动”和“没法 ACID 事务”——传统 Hive 表写入时如果程序中途失败,会出现脏目录,且读时无法看到一致性的快照。Lakehouse 的思路是在数据湖的存储之上,引入表格式(Table Format)层,Iceberg 是这类方案的代表。
讲 Iceberg 时重点提三个机制:
- ACID 事务能力:每次写入生成新的 Manifest 元数据文件,通过元数据文件切换实现快照隔离。多个读者可以同时读取不同快照,互不阻塞;
- 时间旅行(Time Travel):可以查询某个历史时间点的表快照,不需要把数据 Copy 出来,这对数据回溯和数据修复极有价值;
- 隐藏分区(Hidden Partitioning):建表时声明分区字段和分区策略,Iceberg 自动处理分区目录,用户不需要关心分区列值溢出或者目录层级过深的问题。很多团队从 Hive 迁到 Iceberg,第一个感受到的收益就是这个——再也不用维护一堆手动脚本去补分区。
面试回答里的落点是:“你在生产环境为什么要考虑 Iceberg?” 回答方向:第一,多个计算引擎(Spark、Flink、Trino)需要同时访问同一张表,并保证读写一致性;第二,需要支持行级更新但不是高频 Update,Iceberg 的 Merge-on-Read 可以平衡查询性能和写入成本;第三,需要增量读取能力——下游数据服务可以只消费新增的快照数据,而不用全表扫描。这三个场景对应的大数据热门趋势是“湖仓一体、批流一体”,把这些关键词嵌进回答里,面试官会认为紧跟行业方向。
最后补一个生产常用的小技巧:在 Iceberg 表上做小文件合并(Compaction)。Spark 作业常见命令如下:
import org.apache.spark.sql.SparkSession val spark = SparkSession.builder() .appName("iceberg-compaction") .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") .config("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog") .config("spark.sql.catalog.local.type", "hadoop") .config("spark.sql.catalog.local.warehouse", "hdfs://namenode:9000/warehouse") .getOrCreate() spark.read.format("iceberg").load("local.db.user_events") .write.format("iceberg") .option("compaction", "true") // 触发 Iceberg 重写小文件 .mode("append") .save("local.db.user_events")说明:通过 Spark 重新读取 Iceberg 表并写回同一张表,Iceberg 会触发文件重写,把碎片小文件合并为大文件。这能显著提升后续查询的扫描效率。参数上主要关注write.target-file-size-bytes(目标文件大小,默认 512MB)和write.merge.is-enabled(是否开启写入时合并)。生产上建议把 Compaction 做成定时任务,在业务低峰期跑,避免影响在线查询。
这套从理论基础到生产参数的回答方式,配合自己真实的带参数排障经历,足以让任何一份“含答案”的大数据面试题 PDF 变成真正能用的面试武器。最后记住一点:面试官永远在找“真正碰过数据的人”,你的每一句参数、每一个坑,都要能给出来源——测过、看过日志、还是读过源码。说不出来源的答案,不如不说。
本文还有配套的精品资源,点击获取