从单机到集群
上一本《数据治理与清洗手册》里,小哲把数据洗得干干净净——可数据量也爆炸了:日增 5000 万条日志、千万用户、账小灵每秒都在产生对话流。MySQL 查询越来越慢,报表要跑 3 小时。周师傅说:是时候进入「分布式」的世界了。这一本,从单机到集群,把大数据工程师的知识栈一层层点亮。
数据量爆了
单机再强,也扛不住「多」师傅!数据治理刚做完,数据量又爆了——日增 5000 万条行为日志、1000 万用户,账小灵每秒几十万条消息。MySQL 慢查询一堆,昨天的报表今天还没跑出来……
恭喜你,正式踏入「大数据」的门槛。大数据的本质就一个字——多:数据多到一台机器装不下、算不动。解法也只有一条路:分布式——把数据切开,让几百台机器一起干活。
这一本,我带你点亮大数据工程师的知识栈,十站路线——
记住一句话:大数据不是「一种技术」,是「用一群机器解决单机解决不了的问题」的工程哲学。走,第一站,先立世界观。
分布式世界观:CAP 与扩展
水平扩展 · CAP 理论 · 分而治之 · 数据本地性为什么单机不行?为什么不是买台更贵的机器?答案在「扩展」两个字里——垂直扩展有天花板,水平扩展没有。
先分清两种「变大」:垂直扩展——换更猛的 CPU、加更大的内存,一根筋但贵,还有物理上限;水平扩展——多加几台普通机器,把它们组织起来一起干活,便宜、无限、还能扛故障。大数据选的就是水平扩展。
但分布式不是免费午餐——它绕不开一个哲学问题:CAP 理论。分布式系统在一致性(C)、可用性(A)、分区容错(P)三者中只能保证两个——网络一定会断(分区必然存在),所以实际是「CP 还是 AP」的选择题。想通了它,你就理解为什么有的系统「最终一致」、有的「宁可拒绝也不给旧数据」。
【分而治之:大数据的总纲】
100 亿条数据 → 切成 1 万块 → 100 台机器各算 1 万条
→ 汇总 → 结果。这就是 MapReduce 和一切分布式计算的雏形
【CAP 理论(必考)】
C 一致性:所有节点同时看到同一份数据
A 可用性:任何时刻都能读写(不报错)
P 分区容错:网络断了系统还能工作
→ 网络断是必然(P 必须选),所以只能 CP 或 AP 二选一
→ MySQL 集群选 CP,缓存 Redis 选 AP(最终一致)
【数据本地性(Data Locality)】
与其把数据搬到计算处,不如把计算搬到数据处
→ 谁存着数据,谁就来算——少搬数据,快十倍
【故障是常态】
100 台机器里,每天坏一台是「正常」的
分布式系统必须「设计成能容忍故障」:副本、重试、自愈
【面试必问金句】
「分布式系统的本质:一群会坏的单机,协同完成单机做不了的事」
报表查询从 3 小时变成 30 小时?升级单机 MySQL 顶配要 50 万,且半年又不够。改分布式:3 台普通服务器组 Hadoop 集群,日增 5000 万日志照单全收,报表 20 分钟跑完——成本只有五分之一,且以后加机器就行(水平扩展)。
世界观术语
- 垂直 vs 水平扩展:换更强的单机 vs 加更多机器——大数据永远选后者。
- CAP 理论:一致性 / 可用性 / 分区容错——分布式系统的「三选二」哲学。
- 分而治之(Divide & Conquer):大数据计算的底层思想。
- 数据本地性:计算跟着数据走,少搬数据。
- 故障容错:机器会坏是常态——副本、重试、自动恢复是标配。
- 最终一致:AP 系统的常见选择——「过一会儿一定一致」。
本站收获:分布式世界观 = 水平扩展(加机器)+ 分而治之(切数据)+ CAP(做取舍)+ 容错(防坏机器)。想通这四条,大数据的大门就打开了。
存储:HDFS 与对象存储
HDFS · NameNode/DataNode · 副本 · 对象存储数据住哪?单机磁盘装不下、还怕坏。分布式文件系统 HDFS,把几百台机器的磁盘「拼」成一个巨大的虚拟硬盘。
HDFS(Hadoop 分布式文件系统)是大数据存储的老大哥,核心设计:
- NameNode(老大):只记「元数据」——哪个文件、分成哪些块、块在哪台机器。不做数据搬运,所以很轻。
- DataNode(小弟们):真正存数据块(默认 128MB/块)——几百台机器各存各的块。
- 副本机制:每块默认存 3 份,放在不同机器(甚至不同机架)——坏两台机器数据都不丢。
现代趋势:对象存储(S3/OSS/MinIO)——把「文件」抽象成「对象 + 键」,无限扩容、便宜,云上大数据的新底座。
┌─────────────┐
│ NameNode │ ← 老大:只管元数据(目录/块位置/副本)
└──────┬──────┘
┌──────┼──────┐
┌────┴──┐ ┌┴─────┐ ┌┴─────┐
│DN-01 │ │DN-02 │ │DN-03 │ ← DataNode 小弟们:各存数据块
└───────┘ └──────┘ └──────┘
写入 300MB 文件:
切成 3 块(128+128+44MB)
块1 → DN-01, DN-02, DN-03(3 副本,跨机器!)
块2 → DN-02, DN-03, DN-01
→ 任意坏 2 台,数据不丢
【对象存储(现代新宠)】
S3 / 阿里云 OSS / MinIO(自建)
优势:无限扩容、按量付费、和 Spark/Flink 无缝对接
大数据新项目:对象存储 + 计算分离 已成主流
【考点:NameNode 单点】
老大挂了整个集群不可用 → 高可用方案(Active/Standby + JournalNode)
全量行为日志(日增 5000 万条)写入 HDFS/对象存储:按天分区存放,3 副本保证不丢;要重算历史数据时,Spark 直接扫存储层——「存得下、坏不了、随时算」。
存储术语
- HDFS:分布式文件系统——几百台磁盘拼成一个大硬盘。
- NameNode / DataNode:元数据老大 / 数据块小弟——职责分离。
- 块(Block)与副本:默认 128MB 一块、3 副本——容错的基础。
- 机架感知:副本尽量放不同机架——防「一个机柜断电全没」。
- 对象存储(S3/OSS):云上无限存储——湖仓时代的存储底座。
- 存储计算分离:存储和计算独立扩容——现代数据平台架构。
本站收获:存储 = HDFS(老大管目录、小弟管数据、3 副本保命)或对象存储(云上无限扩展)。数据先「住得下」,才能「算得动」。
批处理:MapReduce 与 Spark
MapReduce · Shuffle · Spark RDD · Spark SQL数据存好了,怎么算?「分而治之」落地成 MapReduce;而 Spark 把它的「中间落盘」换成「内存计算」——快一百倍。批处理,是大数据的看家本领。
MapReduce 是 Google 提出的编程模型,就两步:Map(把数据切分、各自处理,产出「键值对」)→ Shuffle(按 key 把相同键的送到同一台机器)→ Reduce(对每组键值对汇总)。词频统计就是它的经典教学案例。
但 MapReduce 每一步都落盘(写磁盘),太慢。于是 Spark 出场:把中间结果放内存,用抽象 RDD/DataFrame 表达计算,还带血缘(算到一半挂了,按血缘重算,不用从头来)。再往上一层 Spark SQL:写 SQL 就能跑分布式计算——大数据工程师的日常就是写 SQL。
-- 100 亿条订单,算各分类营收(Spark 自动分布式执行)
SELECT category, SUM(amount) AS revenue
FROM orders
WHERE dt BETWEEN '2026-09-01' AND '2026-09-30'
GROUP BY category
ORDER BY revenue DESC;
-- 背后发生了什么:
-- 读 HDFS/对象存储 → 自动分区并行扫描(Map)
-- 按 category 打散重排(Shuffle)
-- 各节点局部聚合 → 汇总(Reduce)
-- 100 台机器并行,20 分钟出结果
【Spark 核心概念】
RDD / DataFrame:分布在不同机器上的数据集合
宽依赖 vs 窄依赖:窄依赖不用 Shuffle(快),宽依赖要(慢)
血缘(Lineage):计算步骤记录,故障可重算
惰性求值:只记「计算计划」,真正要结果才执行
Spark SQL / DataFrame API:写起来像 SQL,跑起来分布式
每天凌晨:Spark 作业扫描全天订单 → 计算营收/留存/分类分析 → 结果写回数仓 → 早上 8 点老板看板自动更新。从 MySQL 单机 3 小时,到 Spark 集群 20 分钟——批处理是大数据平台的「夜班工人」。
批处理术语
- MapReduce:Map(各自处理)→ Shuffle(按键分组)→ Reduce(汇总)——分而治之的落地。
- Spark:内存计算引擎——比 MapReduce 快几十上百倍。
- RDD / DataFrame:分布式数据集抽象——DataFrame 更高级、更 SQL。
- 宽依赖 / 窄依赖:要不要 Shuffle——性能优化的关键判断。
- 血缘(Lineage):计算可重放——容错的聪明办法。
- Shuffle(洗牌):数据跨节点重排——分布式计算最贵的环节。
- 批处理 vs 流处理:一次算完一堆(离线)vs 来一条算一条(实时)。
本站收获:批处理 = MapReduce 思想 + Spark 内存计算 + Spark SQL 写 SQL。离线报表、日结统计、全量分析——都靠这台「夜班发动机」。
实时计算:Kafka 与 Flink
流处理 · 窗口 · Checkpoint · Exactly-Once报表可以等 20 分钟,但风控不能——盗刷发生 3 秒后就该拦截。实时计算(流处理),让数据「来一条,算一条」。
实时计算双雄:
- Kafka:消息队列(数据的高速公路)——生产者把消息发进来,消费者按需取走;Topic/Partition 切分、Offset 记录进度、消费者组并行消费。它是整个实时体系的「血管」。
- Flink:真正的流计算引擎——数据来一条处理一条;有状态(记住之前的累计值);窗口(按时间/数量把流切成小批);Checkpoint(状态定期快照,挂了能恢复);Exactly-Once(不重不丢,账一分不差)。
-- Flink SQL:流上直接写 SQL
INSERT INTO revenue_per_minute
SELECT TUMBLE_START(ts, INTERVAL '1' MINUTE) AS win_start,
SUM(amount) AS revenue
FROM orders_stream -- Kafka 里的实时订单流
GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE); -- 1 分钟窗口
【Flink 核心概念】
有状态计算:记住历史累计(每分钟营收要基于之前的状态)
窗口:滚动窗口/滑动窗口/会话窗口
Checkpoint:状态快照 → 故障从快照恢复(不从头算)
Exactly-Once:精确一次语义——账不能多记也不能少记
背压(Backpressure):下游慢了,上游自动减速(别挤爆)
【Kafka 核心概念】
Topic(主题)/ Partition(分区)/ Offset(偏移量)
消费者组:一组消费者并行消费同一 Topic
Broker:Kafka 节点;副本机制保消息不丢
支付事件进 Kafka → Flink 实时计算:同一账号 1 分钟内异地大额支付 → 秒级告警拦截(盗刷止损);同时账小灵根据用户实时行为流更新推荐——「来一条、算一条、应一条」,全程毫秒级。
实时计算术语
- Kafka:分布式消息队列——数据流的中枢神经系统。
- Flink:流计算引擎——有状态、窗口、精确一次。
- 窗口(Window):把无限流切成有限批——滚动/滑动/会话。
- Checkpoint / 状态:定期快照 + 状态恢复——实时计算的容错。
- Exactly-Once:不重不丢——钱的场景必须它。
- 背压:上下游速度不匹配时的自动调节——防雪崩。
- 离线 vs 实时:批处理(T+1 报表)vs 流处理(秒级响应)——两者互补,不是替代。
本站收获:实时计算 = Kafka 当血管(消息流转)+ Flink 当心脏(流式计算)。窗口切流、Checkpoint 保命、Exactly-Once 保账——实时体系让数据「热乎着用」。
采集通道:日志与 CDC
埋点日志 · Filebeat/Flume · Canal CDC · 数仓入口数据平台再强,数据进不来也是白搭。采集通道,就是数据的「进水管」——日志、埋点、数据库变更,各走各的管道进 Kafka。
三大采集场景,三种管道:
- 日志/埋点:App 和网页的行为埋点(点击、浏览、下单)→ 采集器(Filebeat/Flume/Logstash)读取日志文件 → 清洗 → 发进 Kafka → 落数仓。埋点规范(事件名、属性、版本)是质量的前提——衔接第八本的数据治理。
- 数据库变更(CDC):业务库(MySQL)里的订单变了,怎么同步到数仓?Canal / Debezium 监听 binlog(数据库的变更日志),把「增删改」翻译成消息流进 Kafka——实时同步的黄金通道。
- 文件/API:第三方 CSV、接口数据 → 定时拉取或 webhook 接入。
【埋点日志流】
App 埋点 → 日志文件 → Filebeat(采集) → Kafka(transport)
→ Flink 清洗/解析 → 数仓 ODS 层 → 质量规则校验(第8本)
【数据库 CDC 流】
MySQL 订单表 ──binlog──→ Canal → Kafka(CDC topic)
→ Flink → 实时数仓 / 宽表更新
→ 好处:业务库不用改一行代码,变更自动同步
【采集设计要点】
埋点规范:事件名 snake_case、必填属性、版本号(防乱埋)
幂等采集:重复投递不产生重复数据(配合去重)
分区策略:按时间/用户 hash 分区,保证有序与并行
背压与重试:下游慢了自动缓冲,别丢数据
【Kafka 在采集中的角色】
削峰填谷:业务高峰瞬间百万条,Kafka 先存着慢慢消费
解耦:生产者和消费者互不等待——「高速公路服务区」
大促瞬间每秒 50 万条订单和埋点涌入——直接写库必崩。方案:全部先进 Kafka(高速公路服务区先停着),Flink 按能力慢慢消费、实时计算实时大屏;数据库用 Canal CDC 同步到数仓。高峰扛住了,数据一条没丢。
采集术语
- 埋点:App/网页里埋的行为上报点——数据采集的源头。
- Filebeat / Flume / Logstash:日志采集器——读文件、轻量转发。
- CDC(变更数据捕获):监听数据库 binlog 同步变更——Canal / Debezium。
- 削峰填谷:Kafka 缓冲突发流量——实时体系的「蓄水池」。
- 解耦:生产消费分离——上游挂了不影响下游。
- 数据接入规范:埋点规范、幂等、分区策略——采集质量即数据质量。
本站收获:采集 = 日志走 Filebeat、库变更走 Canal CDC、大流量靠 Kafka 削峰。数据进得来、进得稳、进得规范——平台才有米下锅。
OLAP 分析:Hive 与 ClickHouse
OLTP vs OLAP · Hive · 列式存储 · ClickHouse老板问「这个月的营收按渠道分布看看」,SQL 要跑 3 小时?OLAP 分析引擎就是为「复杂查询秒回」而生的。
先分清两种数据库:OLTP(在线事务处理)——MySQL 这类,专门干「增删改查一笔」的活,单条快;OLAP(在线分析处理)——专门干「扫描几亿行算聚合」的活,分析快。术业有专攻。
分析层两大代表:
- Hive:SQL on Hadoop——把 SQL 翻译成 MapReduce/Spark 作业跑在集群上,能处理超大表,但慢(分钟级),适合离线报表。
- ClickHouse / Doris / StarRocks:现代 OLAP 引擎——列式存储(只读需要的列)+ 向量化计算,万亿行也能秒级返回,适合交互式看板。
查询:SELECT category, SUM(amount) FROM orders GROUP BY category
【MySQL(行式存储)】
一行的所有列挨着存 → 就算只要 category+amount,
也得把整行的其他列全读出来 → 10 亿行 × 全列 = 慢
【ClickHouse(列式存储)】
每列单独存 → 只读 category、amount 两列!
读的数据量少 90% + 向量化批量计算
→ 10 亿行聚合,秒级返回
【Hive vs ClickHouse 分工】
Hive:超大表离线全量算(小时级)→ 结果灌给 ClickHouse
ClickHouse:对外提供秒级查询(老板看板直接查它)
→ 「重活给 Hive,快活给 OLAP」
【数仓技术栈配合】
明细/汇总(DWS)→ ClickHouse/Doris 加速层 → BI 看板
宽表建模、物化视图、预聚合——OLAP 提速三板斧
原来老板点开看板,后台跑 Hive SQL 要等几分钟;改造后:离线任务先把汇总结果灌进 ClickHouse,看板直接查 ClickHouse——「按渠道看营收」「按地区下钻」全部 3 秒内返回。老板满意了,数据分析师也解放了。
OLAP 术语
- OLTP vs OLAP:事务处理(单条快)vs 分析处理(批量快)——两种数据库哲学。
- Hive:SQL on Hadoop——能算超大表,适合离线。
- 列式存储:按列存、按需读——OLAP 快的头号功臣。
- ClickHouse / Doris / StarRocks:现代 OLAP 引擎——万亿行秒级查询。
- 向量化计算:CPU 批量算一列——比逐行快一个数量级。
- 物化视图 / 预聚合:把常用结果提前算好——看板提速三板斧。
本站收获:OLAP = 列式存储 + 向量化 + 预聚合。Hive 干重活、ClickHouse 干快活——「老板要看秒回,重活留到半夜」。
调度资源:YARN 与编排
YARN · 任务编排 · 队列 · 依赖告警集群里同时跑几百个任务,谁先谁后、谁抢谁的资源、挂了下游怎么办?调度与资源管理,是大数据平台的「交警 + 后勤部长」。
YARN(Yet Another Resource Negotiator)是 Hadoop 的资源管理器:集群的内存/CPU 都归它管,谁申请、给多少、用完还。现代平台也在演进到 K8s 管资源(云原生趋势)。
光有资源还不够——几百个任务要有「编排」:
- 任务编排(Airflow / DolphinScheduler):定义任务依赖——「清洗完才能算汇总,汇总完才能出报表」。
- 调度策略:定时(每天 02:00)、失败重试、超时告警、上下游自动触发。
- 资源隔离:按团队/优先级分队列——报表队和实时队互不抢饭。
02:00 定时触发
│
┌────────┴────────┐
采集任务A 采集任务B(并行!)
└────────┬────────┘
│
清洗任务(等 A、B 都完成)
│
┌─────────┼─────────┐
汇总任务C 汇总任务D 质量校验
└─────────┼─────────┘
│
报表生成任务(等全部完成)
│
08:00 老板看到日报 ✓
【编排要点】
DAG:有向无环图——任务依赖的「流程图」
失败重试:网络抖了?重试 3 次再说
超时告警:跑超 4 小时 → 钉钉提醒
补数据:昨天失败了?支持「回刷」历史日期
【资源管理】
YARN 队列:dev/prod 隔离,生产任务优先
K8s 化:容器调度 Spark/Flink——云原生新趋势
成本:错峰跑批(晚上便宜)、弹性扩缩容
每天凌晨,DolphinScheduler 准时拉起 300 个任务:采集→清洗→汇总→质量校验→出报表,依赖关系清清楚楚;某天采集 B 失败,自动重试 3 次仍失败 → 只告警「采集 B」,下游自动等待不空跑;修复后一键「补数据」,全链路重跑,老板照常 8 点看到日报。
调度术语
- YARN:集群资源管理器——CPU/内存的「物业公司」。
- DAG 编排:任务依赖图——Airflow / DolphinScheduler 的核心。
- 队列与优先级:资源按团队/任务等级隔离分配。
- 失败重试 / 告警:自动化运维三件套——重试、告警、补数据。
- K8s 化:云原生调度——Spark/Flink 跑在容器里,弹性伸缩。
- 错峰与成本:大任务排晚上跑、用竞价实例——省钱也是工程师的 KPI。
本站收获:调度 = YARN/K8s 管资源 + DAG 管依赖 + 重试告警管稳定性。几百个任务不乱套,全靠这「交警 + 后勤」。
生态选型:Hadoop 全家桶
Hadoop 生态 · HBase · Zookeeper · 云上大数据Hadoop 不是一个软件,是一个「家族」——几十个组件各管一摊。这一站,把全家桶认全,再讲讲怎么选型。
Hadoop 生态的「分工表」,记牢:
- HDFS:存储(第 2 站讲过)
- YARN:资源管理(第 7 站讲过)
- MapReduce / Spark:计算(第 3 站)
- Hive:SQL 分析(第 6 站)
- HBase:分布式 NoSQL——海量数据的「随机读写」数据库(实时查询场景,比如订单明细秒查)。
- Zookeeper:分布式协调——「谁当领导」的裁判(选主、配置中心)。
- Kafka / Flink:消息与流计算(第 4 站)
选型大趋势:自建 vs 云上——自建(CDH/HDP)灵活但运维重;云上(EMR/Databricks/阿里云 MaxCompute)开箱即用、弹性付费。中小企业上云是大势。
┌─────────────────────────────────────────┐
│ 应用层:BI看板 / 数据服务API / 账小灵 │
├─────────────────────────────────────────┤
│ 分析层:Hive | Spark SQL | ClickHouse │
├─────────────────────────────────────────┤
│ 计算层:Spark(批) | Flink(流) | MR(老) │
├─────────────────────────────────────────┤
│ 存储层:HDFS | 对象存储 | HBase | 数仓 │
├─────────────────────────────────────────┤
│ 资源层:YARN | K8s │
├─────────────────────────────────────────┤
│ 协调层:Zookeeper(选主/配置) │
├─────────────────────────────────────────┤
│ 采集层:Filebeat | Canal CDC → Kafka │
└─────────────────────────────────────────┘
【HBase 何时用】
海量 + 随机读写 + 准实时查询
例:亿级订单明细「按用户查最近 100 条」
行键设计(RowKey)是性能关键——热门考点
【自建 vs 云上(选型决策)】
自建:数据敏感/量超大/已有机房 → CDH/Doris 自建集群
云上:起步快/弹性/省运维 → EMR / MaxCompute / Databricks
混合:核心自建 + 弹性上云(峰时扩容)
【核心组件部署形态】
传统:Hadoop 全家桶一套集群
现代:对象存储 + Spark/Flink + OLAP,轻量云原生
小哲公司评估后选了「云上组合」:对象存储存全量、EMR 跑 Spark 离线、Kafka+Flink 实时、ClickHouse 出看板——3 个工程师两周搭完,按量付费,从「要养一个运维团队」变成「按需租用」。等规模大了再考虑自建。
生态术语
- Hadoop 生态:存储(HDFS)+ 计算(MR/Spark)+ 资源(YARN)+ SQL(Hive)+ 协调(ZK)的家族。
- HBase:分布式列式 NoSQL——海量随机读写。
- Zookeeper:分布式协调服务——选主、分布式锁、配置中心。
- 云上大数据:EMR、Databricks、MaxCompute——开箱即用的托管平台。
- 数据服务化:数仓结果 → API 服务 → 业务/AI 消费(衔接账小灵)。
- 技术选型:没有银弹——按数据量、团队、预算选。
本站收获:生态 = 全家桶各司其职(存/算/管/查/协调)+ 现代趋势(对象存储 + 轻量组件 + 上云)。会选型,比会敲命令更值钱。
湖仓一体:现代数据栈
数据湖 · Iceberg/Hudi/Delta · ACID · 时间旅行数据湖便宜灵活但「不可信」(更新麻烦、ACID 弱);数仓可靠但贵。湖仓一体 = 湖的存储 + 仓的能力——一个平台,全量历史 + 可靠更新。
回顾一下(衔接第八本《数据治理》第 9 站):数据湖什么格式都收、便宜,但做不了「更新和事务」;数仓干净可靠,但贵、灵活度低。湖仓一体用「表格式」技术(Iceberg / Hudi / Delta Lake)给湖里的数据加上数仓能力:
- ACID 事务:写入要么全成要么全不成——不再「算到一半数据缺一半」。
- 时间旅行(Time Travel):按时间戳查「过去的版本」——数据错了能回到昨天。
- Schema 演进:加字段不用重写全表——业务变更不再伤筋动骨。
- 增量更新:只改变化的部分,不重算全量。
【架构演进】
第一代:HDFS + Hive(仓)与 数据湖(分开两套)→ 数据双份、同步难
第二代:湖仓一体 —— 对象存储 + Iceberg/Delta + Spark/Flink
↓
一份存储(湖) + 一套表格式(ACID/时间旅行)+ 多种引擎(批/流/BI)
【Iceberg 核心概念】
Table Format(表格式):给数据湖加「事务层」
Snapshot(快照):每次写入产生一个快照 → 时间旅行的基础
Metadata(元数据):表结构、分区、快照的管理
【应用方式】
离线:Spark 读写 Iceberg 表(全量 + 增量)
实时:Flink 流式写入 Iceberg(分钟级可见)
查询:Spark SQL / 引擎直接查
【和数仓的关系(不是替代!)】
湖仓一体 = 底座;数仓分层思想照用(ODS/DWD/DWS/ADS)
只是「住在湖里」+「有了 ACID」+「能时间旅行」
某天清洗任务写错了逻辑,覆盖了昨天的订单明细。传统数仓:从备份恢复,折腾 3 小时;Iceberg 湖仓:一条 SQL「回到昨天的快照」,5 分钟找回,业务零感知——时间旅行就是数据工程师的后悔药。
湖仓术语
- 湖仓一体:湖的便宜灵活 + 仓的 ACID 可靠——现代数据平台主流。
- 表格式(Iceberg/Hudi/Delta):给数据湖加事务能力的「中间层」。
- ACID 事务:原子性/一致性/隔离性/持久性——数据可信的保证。
- 时间旅行:查任意历史快照——数据恢复的后悔药。
- Schema 演进:表结构平滑变更——业务进化不推倒重来。
- 批流一体:同一套表,离线和实时都能写——Spark 批 + Flink 流。
本站收获:湖仓一体 = 一份存储 + 事务能力 + 时间旅行 + 批流共用。它是第八本「数据治理」的平台级落地——数据湖从此「能信、能改、能后悔」。
性能调优:工程师成长
数据倾斜 · 小文件 · Spark/Flink 调优 · 职业路线同一个 SQL,菜鸟跑 2 小时,专家跑 10 分钟——差距就在调优。调优是大数据工程师的「内功」,也是面试的重头戏。
大数据性能问题,九成出在四个地方:
- 数据倾斜(Skew):某个 key 数据特别多,一台机器累死、其他机器闲着——「木桶的短板」。
- 小文件问题:几百万个小文件,读文件的开销比算数据还大——「一口袋米粒」。
- Shuffle 太重:跨节点搬数据太多次——「快递费比货还贵」。
- 资源与内存:并行度不够、内存溢出——「没吃饱饭还硬干活」。
【数据倾斜】
症状:99% 任务 5 分钟跑完,最后 1% 跑了 3 小时
原因:某个 key(如「北京」)数据量巨大
解法:加盐打散(随机前缀拆分再合并)、两阶段聚合、
广播小表(Join 时小表广播到每台机器)
【小文件问题】
症状:几百万个小文件,任务启动就花了 10 分钟
解法:合并小文件(Coalesce/Repartition 控制分区数)、
写入前设置目标文件大小、分区裁剪
【Shuffle 优化】
症状:任务 70% 时间在「搬数据」
解法:减少宽依赖、合理分区键、使用广播变量、
聚合类算子优先(reduceByKey 优于 groupByKey)
【资源与内存】
症状:OOM 内存溢出 / 并行度不够
解法:调整 Executor 内存/并行度、Flink 状态 TTL、
背压排查(下游消费不过来)
【Flink 特有】
状态过大 → 调大 RocksDB 状态后端 / 清理过期状态
窗口乱序 → Watermark 水位线设计
【面试加分金句】
「先看数据特征(倾斜/大小/分布),再谈调优——
调优的本质是『让每一台机器都吃饱,且不多干活』」
成长路线
- 数据倾斜:加盐、两阶段聚合、广播——面试必考题。
- 小文件:合并、控制分区——写入即治理。
- Shuffle / 内存:减少搬运、合理资源——性能大头。
- 监控与定位:Spark UI / Flink UI 看阶段耗时——先定位再优化。
- 职业路线:SQL 工程师 → 数据开发 → 大数据架构师 → 数据平台负责人;SQL 是门槛,调优是内功,架构是天花板。
- 持续学习:Spark/Flink 源码、云原生(K8s)、数据产品思维——技术 + 业务两条腿。
某日报任务跑 3 小时,Spark UI 一看:90% 时间卡在最后一个 Stage——「北京」这个 key 倾斜。加盐两阶段聚合后,任务 12 分钟跑完。这次调优的收益:每天省 2.8 小时算力 + 报表提前 3 小时——工程师的价值就体现在这种「十分钟的洞察」。
本站收获:调优 = 定位(看 UI 找瓶颈)→ 分析(倾斜/小文件/Shuffle)→ 对症(加盐/合并/广播)。「让每台机器都吃饱,且不多干活」——这就是大数据工程师的内功。
技术栈全景 + 选型速查
面试前最后一页技术栈一句话速查
| 环节 | 主流技术 | 一句话人话 |
|---|---|---|
| 采集 | Filebeat / Flume / Canal / Debezium | 日志和数据库变更的「进水管」 |
| 传输 | Kafka | 数据高速公路服务区(削峰解耦) |
| 存储 | HDFS / 对象存储(S3/OSS) / HBase | PB 级数据住的地方(3 副本保命) |
| 批计算 | Spark / MapReduce | 夜班工人:离线批量算 |
| 流计算 | Flink | 白班收银员:来一条算一条 |
| SQL 分析 | Hive / Spark SQL | 分布式 SQL——大表离线查 |
| OLAP | ClickHouse / Doris / StarRocks | 秒级看板引擎——老板的最爱 |
| 调度 | Airflow / DolphinScheduler / YARN / K8s | 流水线交警 + 物业公司 |
| 湖仓 | Iceberg / Hudi / Delta Lake | 给数据湖加 ACID 和后悔药 |
| 云上 | EMR / Databricks / MaxCompute | 开箱即用的大数据全家桶 |
场景 → 选型速查
| 业务场景 | 推荐方案 |
|---|---|
| 离线日报 / 月报 | 对象存储 + Spark + Hive → 结果进 OLAP |
| 实时大屏 / 风控 | Kafka + Flink + 实时数仓 |
| 老板交互式看板 | ClickHouse / Doris(秒级) |
| 亿级订单明细秒查 | HBase(RowKey 设计好) |
| 全量历史 + 增量更新 + 后悔药 | 湖仓一体(Iceberg/Delta) |
| 团队小、预算少 | 云上托管(EMR/云数仓)——别自建 |
是「用一群机器解决单机解决不了的问题」的工程哲学。
批处理养报表,流处理守实时,湖仓一体保后悔药——
大数据工程师的价值,就是把「不可能算完」变成「按时算完」。