AI 项目的命门不是模型,是数据。FDE 模块 08 提到了 ETL/ELT,但没展开。本页补齐数据工程与大数据的完整知识链:从数据采集与编排、数仓维度建模、Spark/Flink 批流处理,到数据湖/湖仓一体,再到数据质量与治理。我会像讲《Java 进阶》那样,把每个知识点拆成"它是什么 → 为什么这样设计 → 怎么用 → 踩过什么坑",让你学完能搭一条"原始日志 → 干净表 → 可被模型/报表消费"的可靠数据链路。基线:Spark 3.x / Flink 1.18+ / Iceberg 1.x。
数据采集与 ETL/ELT:把"脏数据"变成"可信资产"
数据工程是"后端的数据侧"——核心职责是采集、清洗、建模、供给。模型再聪明,喂的是垃圾数据也只能产出垃圾结论。这一章讲清数据进门的姿势、转换到底发生在哪一步、以及怎么让任务按时按依赖跑。
批处理 vs 流处理:先分清你要哪种节奏
数据进来有两条路。批处理是"攒一堆再算"(小时级、天级),适合报表、训练样本;流处理是"来一条算一条"(毫秒~秒级),适合风控、实时监控。选错节奏,系统会又贵又脆。
| 维度 | 批处理 Batch | 流处理 Streaming |
|---|---|---|
| 触发方式 | 定时/手动触发一批 | 事件触发,持续运行 |
| 延迟 | 分钟~天 | 毫秒~秒 |
| 典型引擎 | Spark、MapReduce | Flink、Kafka Streams |
| 典型场景 | 日报表、训练集、对账 | 反欺诈、监控告警、实时大屏 |
日报、月报本就是 T+1 批处理,强行上实时又贵又脆弱。按业务 SLA 选:风控/监控要实时(Flink),报表/画像可批处理(Spark)。先问"迟一晚上行不行",行就别上流。
论Lambda 与 Kappa:批流怎么共存
Lambda:同时跑批层(准确)和流层(快速),查询时合并两者——健壮但维护两套代码、两套口径。Kappa:只用流,用重放历史消息替代批层——架构简单,但要求流系统能重放全量。现代湖仓(Iceberg + Flink)让 Kappa 越来越可行,新项目可优先考虑单一流架构,少一套心累。
ETL vs ELT:转换到底发生在哪一步
这俩字母顺序就是答案。ETL:抽取(Extract)→ 在中间层转换(Transform)干净 → 再加载(Load)进仓库。传统数仓这么做,因为仓库算力弱,得先洗好再进。ELT:抽取 → 原样加载进湖/仓 → 在仓库里用引擎算力再转换。现代存储便宜、引擎(Spark/DuckDB)快,ELT 更灵活,是当前主流。
论为什么 ELT 成了主流
ELT 把"原始数据"也留在湖里,意味着同一份原始数据可以反复换口径重算,不用每次都重新抽。做错指标?改个 SQL 重跑即可,原始事实不动。这对"口径经常变"的业务分析是巨大优势。ETL 的清洗结果一旦落地,想换逻辑就得重抽,成本高。
抽取(Extract):数据从哪来、怎么拿
数据源五花八门:业务库(MySQL/PG)、日志(埋点/服务器)、第三方 API、消息队列。抽取方式分三类:全量(整表拉)、增量(按时间戳/自增 ID 拉新)、CDC(监听 binlog,近实时捕获增删改)。生产推荐 CDC,对源库压力小、延迟低。
清洗与转换(Transform):把脏数据变干净
原始数据几乎一定脏:空值、类型不一致("123" 和 123 混着)、编码乱、重复、时区错。转换层要做的清单:去重、类型标准化、空值填充/标记、单位统一、字段拆分与拼合、非法值过滤、时区归一。
加载(Load):全量、增量、CDC 怎么落
清洗完要落库。三种策略:全量覆盖(小表、重算安全);增量追加(按日期分区追加,可回刷);Merge/Upsert(按主键更新已有、插入新的,CDC 落湖仓常用)。湖仓表格式(Iceberg)原生支持 MERGE INTO。
论为什么加载也要"可重跑"
调度任务今天挂了、明天重跑,必须幂等——跑一次和跑一百次结果一样。增量追加如果不带去重,重跑会重复;全量覆盖最安全但慢;Upsert 兼顾。设计管道时把"重跑幂等"当铁律,半夜告警才不会被自己坑。
调度编排:Airflow/Dagster 把任务串成 DAG
数据任务有强依赖:先抽数 → 再清洗 → 再聚合 → 再导出。编排工具把这套依赖画成 DAG(有向无环图),自动按时调度、失败重试、可观测。Airflow 用 Python 定义 DAG,生态最大;Dagster 强调"资产(asset)"语义,更现代。
| 编排工具 | 核心模型 | 适合谁 |
|---|---|---|
| Airflow | DAG + 任务,生态最大 | 通用批调度、团队已熟悉 |
| Dagster | 资产(asset)语义,类型感知 | 想让数据资产可追踪的新项目 |
| Prefect | Python 原生 flow,混合调度 | 偏 Python 的工程团队 |
① catchup 没关:任务停了三天再开,Airflow 默认补跑这三天,瞬间压垮源库。② 任务不幂等:重跑就重数据。③ 把重活塞进一个 task:一个 Python 函数干完所有事,失败全重来、不好并行。拆成有依赖的小 task,可重试、可观测。
任务监控与 SLA:失败重试、告警与数据新鲜度
管道跑起来只是开始,跑得稳、出事有人知才是生产。调度系统要配三件事:重试(瞬时失败自动重试,配指数退避)、告警(任务失败/超时/数据空了立刻通知)、SLA 监控(这张表今天该几点就绪,晚了就报)。
重试只该用于幂等且瞬时的失败(网络抖、源库锁)。若任务本身有 bug(SQL 写错),重试只会重复失败、浪费资源、还可能重复写。先分清失败类型:非幂等写操作重试前要确认"会不会重复落库"。
数仓与建模:维度建模让分析变简单
数仓不是"把业务库复制一份"。业务库为写入事务优化(范式化、避免冗余),数仓为分析读取优化(反范式、预聚合)。这一章的明星方法是维度建模——让业务方也能看懂的"事实 + 维度"结构。
维度建模思想:面向分析而非事务
事务型建模回答"现在库存多少、订单状态如何";维度建模回答"去年华东区手机品类的销售额按周趋势"。后者需要把"销售额(可度量)"和"地区/品类/时间(怎么切)"分开。核心就两类表:事实表记"发生了什么(数字)",维度表记"这件事的背景(描述)"。
论为什么数仓要反范式
业务库范式化(拆成很多小表)是为了写入不冗余、不乱。但分析时每查一个维度都要 join 一堆表,慢且难懂。数仓把常用描述冗余进宽表/维度表,用空间换查询简单和速度。记住:数仓的设计目标是"人能看懂、查得快",不是"不冗余"。
事实表与维度表:数仓的两类主角
事实表(Fact)存业务过程的可度量事件,如"一笔订单""一次点击",包含外键(指向维度)和度量(金额、数量)。维度表(Dimension)存描述性属性,如"商品""用户""时间""地区",基本不会变、且常被当作分析切片的依据。
| 特征 | 事实表 Fact | 维度表 Dimension |
|---|---|---|
| 存什么 | 事件 + 度量值(金额/次数) | 描述属性(名称/类别/地区) |
| 规模 | 巨大、持续增长 | 较小、变化慢 |
| 主键 | 业务键 + 代理键 | 代理键(surrogate key) |
| 例子 | 销售事实、点击流事实 | 商品维、时间维、用户维 |
一个常见误区:把事实表和维度表都建得"又宽又全"。事实表每多一个维度外键就多一次 join,维度表塞太多描述列会膨胀。经验是事实表只放可加度量 + 必要外键,描述性的都交给维度表,保持事实表的"瘦"。
论粒度决定事实表能回答什么
同样一笔销售,按"订单"粒度能算订单数,按"订单明细行"粒度才能算件数和品类分布。粒度定得太粗(只到订单),很多下钻分析就做不了;定得太细(到每个 SKU),存储和查询成本又高。先想清楚业务要下钻到哪一层,再钉粒度,这是建模第一步。
星型 vs 雪花模型
星型(Star):一个事实表直接连多个维度表,维度不再拆分,像星星。简单、查询快、最常用。雪花(Snowflake):维度再规范化拆成子维度(如"城市维"拆出"省份维")。省空间但 join 多、查询慢。实际几乎都用星型,雪花只在维度极度冗余时偶尔用。
| 维度 | 星型 | 雪花 |
|---|---|---|
| 结构 | 事实直连扁平维度 | 维度再拆子维度 |
| 查询性能 | 快(少 join) | 慢(多 join) |
| 冗余 | 有冗余 | 少冗余 |
| 易用性 | 高,业务方易懂 | 低 |
论为什么实际几乎都选星型
现代数仓(尤其是列存 + 缓存维度)下,星型的冗余成本极低,而雪花多出来的 join 在分析师自助查询时极其劝退。除非维度表大到离谱,否则星型是默认答案。雪花那点空间收益,换不来查询体验和性能的损失。
缓慢变化维(SCD):历史怎么留
维度会变:顾客搬了城市、商品改了类目。要不要保留"当时"的样子?这决定指标的可比性。SCD Type 1:直接覆盖(丢了历史);Type 2:加生效/失效时间,保留全量历史版本(最常用);Type 3:加一列存"上一版值"(只留最近一次)。
如果只用 Type 1 覆盖,去年同期"上海用户"的销售额会被算到今天的"北京用户"头上,同比彻底失真。凡是分析要按历史口径切片的维度(地区、品类、组织架构),几乎都要上 Type 2。
粒度与一致性维度:建模最容易翻车的两点
粒度(Grain)是事实表"一行代表什么"。一行是"一笔订单"还是"一个订单里的一个商品行"?粒度不清,同一指标会算出不同数。声明粒度是第一件事。一致性维度(Conformed Dimension):多个事实表共用同一套维度(如时间维),否则"销售额"和"流量"无法在同一报表里对齐。
论"一个事实表一个粒度"
绝不要把"订单级"和"订单明细级"混在一张事实表。一行既记整单金额又记商品金额,sum 金额时会把整单重复算 N 次。先钉死粒度(通常是最小可加单元),再决定加哪些度量。这是新手建模翻车第一名。
总线架构与数据集市:分层给不同角色定制
企业级数仓常按 Kimball 总线架构 组织:底层是一致性维度和事实表(共享),上层按业务线切成多个数据集市(Data Mart)——给销售、运营、财务各一套贴合他们视角的汇总表。这样既统一口径,又不让所有人挤一张大表。
| 分层 | 内容 | 服务于谁 |
|---|---|---|
| ODS | 原始贴源层 | 排查/回灌 |
| DWD | 明细层(清洗后事实/维度) | 分析师 |
| DWS | 汇总层(主题宽表) | 报表/应用 |
| ADS | 应用数据集市 | 业务线/看板 |
论分层不是越多越好
小团队三层(ODS/DWD/ADS)足够,硬上五六层只会增加链路长度和出错环节。分层是为了"职责清晰 + 复用",不是炫技。每加一层都要回答"谁复用它",答不上就别加。
桥接表与多值维度:一个订单多个标签怎么建模
经典维度建模假设"一个事实行对应一个维度成员",但真实业务常多值:一个订单贴了 ["促销","会员","新品"] 三个标签,一个患者有多种诊断。直接把多值塞成逗号字符串会毁掉可查询性——SQL 只能 LIKE,无法精确聚合、索引也失效。
标准解法是 桥接表(Bridge Table):建一张"事实↔标签"的多对多关系表,配合权重,让查询能正确展开与聚合。它把"多值"从维度字段里解放出来,变成可 join 的关系,建模才可持续。
举个真实场景:大促时一个订单同时命中多个营销标签,运营要按"标签"维度统计 GMV。用桥接表就是标准 join + group by,既快又准;用逗号字符串就只能全表 LIKE,慢且容易数错。
| 方案 | 做法 | 问题 |
|---|---|---|
| 逗号字符串 | tags="促销,会员" | 无法精确聚合、索引失效 |
| 多列 tag1/2/3 | 预开 N 列 | 超 N 个就丢,稀疏浪费 |
| 桥接表 | fact↔bridge↔dim 多对多 | 查询稍复杂,但正确可扩展 |
论桥接表要小心"双重计数"
多值维度 join 后,一个订单会展开成多行,如果直接 SUM(amount) 会把金额按标签数放大。正确做法是用桥接表的权重(weight)或"先按订单去重聚合再关联标签",否则 BI 报表金额凭空翻倍。这是多值维度建模最容易翻车的地方,务必在建模规范里写死。
批处理引擎:Spark 的 RDD、DataFrame 与优化
Spark 是批处理事实标准。它把计算表达成有向无环图(DAG),在内存里算、容错靠血缘重算。理解它的抽象层级和性能黑洞(shuffle),才能写出既对又快的作业。
RDD:Spark 的底层抽象
RDD(弹性分布式数据集)是不可变、分区的数据集合,是 Spark 最原始的 API。它记录"怎么算出来"的血缘(lineage),某分区丢了能重算而非重跑全部。但 RDD 不感知数据 schema,优化空间小。
论为什么日常几乎不用裸 RDD
RDD 把数据当不透明的 Java 对象,无法做列式优化、无法 Catalyst 优化器重写。DataFrame/Dataset 带 schema,优化器能"看懂"你的意图做谓词下推、列裁剪。除非要极细粒度控制(自定义分区、特殊二进制),否则一律用 DataFrame。
补充一点:RDD 的容错靠 血缘(lineage)——每个 RDD 记得自己是怎么从父 RDD 算来的。某个分区丢失,Spark 沿血缘重算那一条路径,而不是重跑整个作业,这正是它比早期 MapReduce 每次落盘更优雅的地方。理解血缘,就理解了 Spark"算错了能自愈"的底气。
DataFrame/Dataset:声明式、更快
DataFrame 是带 schema 的分布式表,Dataset 是带类型安全的版本(Scala/Java)。你写"我要什么",Catalyst 优化器决定"怎么高效算"。绝大多数作业用它,性能和可维护性都碾压 RDD。
分区与并行度:任务切多细
Spark 把数据切成分区(partition),每个分区一个 task 并行跑。分区太少 → 并行度不足、单 task OOM;太多 → 调度开销爆炸。读数据时按文件/大小自动定分区,重算前用 repartition/coalesce 调。
| 操作 | 是否 shuffle | 用途 |
|---|---|---|
| repartition(n) | 是(重洗牌) | 增/减分区,均匀 |
| coalesce(n) | 否(合并相邻) | 只减分区,写前合并小文件 |
| 默认读分区 | 否 | 按 block 大小自动定 |
太小:比如 2 个分区处理 100G,单 task 爆内存、跑得慢。太大:几万个 1KB 小文件,每个 task 调度开销比计算还大("小文件地狱")。经验值:单分区目标 128MB~256MB,类似 HDFS block。写结果前用 coalesce 控制输出文件数。
Shuffle:批处理性能的头号杀手
groupByKey、reduceByKey、join、repartition 都要 shuffle:把数据按 key 跨节点重新分发。它要落盘 + 网络传输,是 Spark 作业最慢的环节。优化心法:减少 shuffle 数据量、选对算子、避免大 key。
论reduceByKey 为什么比 groupByKey 好
groupByKey 先把所有 value 原样 shuffle 过去再聚合,网络搬的是全量。reduceByKey 先在本地做一轮 combiner 预聚合,再 shuffle 减少后的结果,网络量小得多。能本地先合的就本地先合,这是 shuffle 优化的总原则。
Spark 的 spark.sql.shuffle.partitions 默认 200。对几 MB 的小数据集,200 个分区几乎全空跑,调度开销比计算还大;对超大集群又太少,单分区过大撑爆内存。按"单分区 ~128MB"目标手动调,或开启 AQE(自适应查询执行)让它自动合并小分区、动态优化,省去手调。
宽依赖、窄依赖与算子选择
窄依赖:一个父分区只给一个子分区(map、filter),可流水线、某分区丢了自己重算。宽依赖:一个父分区给多个子分区(groupBy、join),必须等上游全就绪、产生 shuffle、阶段(stage)边界。Spark 在宽依赖处切 stage。
内存与缓存:persist/cache 用对才提速
同一个 DataFrame 被反复用(多次 join、多次聚合),用 cache()/persist() 把它留在内存/磁盘,避免重算血缘。但缓存不是免费的:占内存、要治理 eviction。
① 缓存一次性的中间结果:只用一个 action 就丢,缓存反而拖慢(多一次物化)。② 忘记 unpersist:缓存堆积挤占执行内存,触发频繁 spill 到磁盘,比不缓存还慢。规则:被 ≥2 个 action 复用的才缓存,用完即 unpersist。
广播变量、累加器与 UDF:Spark 的共享机制
算子里的函数跑在各个 executor 上,怎么共享"只读大对象"(如一份字典)或"全局计数"?广播变量(broadcast)把只读数据发到各节点一次,避免每个 task 都序列化一份;累加器(accumulator)做分布式计数/求和(只在 driver 可读最终结果);UDF把自定义函数注册给 DataFrame 用。
论广播 vs 直接闭包捕获
小数据直接闭包捕获也行,但大对象(几 MB 的映射表)若每个 task 都序列化一次会爆炸。broadcast 走高效分发、各节点存一份,省网络和内存。经验:大于几 KB 的闭包捕获对象,改用 broadcast。
流处理:Flink 与 Kafka 的实时世界
流处理要回答"此刻正在发生什么"。Kafka 是流的管道与存储,Flink 是真正的逐事件计算引擎。这一章讲清时间语义、窗口、状态与一致性——这些是流和批最不一样的地方。
Kafka:流的"管道"与"存储"
Kafka 用主题(topic)分频道,每个 topic 分多个分区(partition)保证有序与并行,数据按偏移量(offset)持久化、可被消费者重复读。它既是消息队列,也是短期日志存储——这正是 CDC 和流处理的地基。
论为什么 Kafka 能做"重放"
传统 MQ(如早期 RabbitMQ)消费完就删,出错难补救。Kafka 把消息按 offset 保留(可配保留期),消费者想重来就从旧 offset 重读。这让它既是传输管道,又是"可重放的日志",是流处理容错和回溯的关键。
消费者上线/下线/崩溃会触发 rebalance(再均衡),期间整个组暂停消费、重分配分区,频繁再均衡会让吞吐剧烈抖动,监控上表现为周期性消费停滞。用 CooperativeStickyAssignor 做增量再均衡、调大 session.timeout、避免在处理里长阻塞,能显著减少停顿。
Flink 编程模型:Source → Transform → Sink
Flink 作业是"从哪读(Source)→ 怎么算(Transform)→ 写到哪(Sink)"的 DAG。它把流抽象成 DataStream,算子链(operator chain)把相邻算子拼一起减少线程切换。
时间语义与窗口:事件时间 + 水印
流里"时间"有三种:事件时间(事件发生时刻,最准)、摄入时间(进 Flink 时)、处理时间(算的时候)。真实数据会乱序、会延迟,所以用事件时间 + 水印(Watermark):水印说"事件时间已到 T,早于 T 的我都当到齐了",窗口才关、才出结果。
太小(如 0):稍晚到的数据被当"迟到"丢弃,结果不准。太大:窗口迟迟不关,结果延迟产出、状态堆积。设多少取决于数据的真实乱序程度——从日志观察 p99 延迟再定,别拍脑袋。
状态管理与 Checkpoint
流计算要"记住点东西":累加的计数、去重的集合、窗口里攒的数据,这些叫状态(State)。状态存在哪?Checkpoint 是 Flink 的分布式快照,定期把状态持久化到存储。作业挂了,从最近 checkpoint 恢复,状态不丢。
论状态存在内存还是 RocksDB
状态小(几 MB)用内存后端,最快。状态大(GB 级,如长窗口去重)内存扛不住,用 RocksDB(落本地盘 + 可溢写),代价是序列化开销。选后端本质是"状态规模 vs 性能"的权衡,别无脑上 RocksDB。
Exactly-once:端到端一致性
精确一次(Exactly-once)不是"每条只算一次"(不可能,至少一次才是物理现实),而是"最终结果等于精确一次":靠 checkpoint 恢复 + 两阶段提交(sink 与 checkpoint 对齐)实现。Flink 内部精确一次易得,端到端还要 sink 支持事务(如 Kafka 事务性 producer、支持 upsert 的数据库)。
论别神话 Exactly-once
它对"重复计算导致金额翻倍"这类场景是刚需;但对"多算一次点击统计"往往 at-least-once + 下游去重就够了,且更轻。先问业务能否容忍重复,再决定是否付出两阶段提交的复杂度与延迟。不是所有流都要精确一次。
Flink SQL:用 SQL 写流处理
不是所有人都想写 Java/Scala DataStream。Flink SQL 让分析师和工程师用熟悉的标准 SQL 声明流:把 Kafka topic 当"动态表",INSERT INTO ... SELECT 就是持续的计算。它底层仍是同一套状态/checkpoint 机制,只是换了个入口。
论SQL 与 DataStream 怎么选
常规聚合、join、维表关联用 SQL 足够且快;涉及自定义状态机、复杂事件模式(CEP)或精细控制算子链时,回到 DataStream API。两者可混用:SQL 起大部分活,特殊逻辑用 DataStream 补。
Kafka 事务性 Producer:端到端精确一次的源头
Flink 的精确一次需要"源"也配合。Kafka 事务性 Producer 能把"向多个分区/主题写消息"打包成一个原子事务:要么全部对消费者可见,要么全不可见(靠 isolation.level=read_committed 隔离)。CDC 入湖时常用它,保证"这一批变更要么都进、要么都不进",不会出现"一半写入、一半没写"的中间态。
它和幂等 Producer 常被混为一谈,但层级不同:幂等只解决单分区内重试导致的重复(靠序列号去重);事务进一步保证跨分区原子可见,是解决"消费-转换-生产"链路精确一次的关键。两者通常一起开,幂等是事务的底层基础。
论为什么 CDC 入湖离不开事务
设想一次 binlog 抓了 1000 条变更,写到第 500 条时作业挂了:非事务下,下游看到"残缺的一半",对账直接错乱。事务保证"要么 1000 条全可见、要么 0 条",配合 Flink 从 checkpoint 重放,才能做到端到端不重不漏。理解这层,才懂精确一次为什么需要整条链路协同。
数据湖与湖仓一体:让湖也能做仓的活
数仓贵且锁死结构,数据湖便宜但"存了等于没管"。湖仓一体(Lakehouse)用开放表格式把 ACID、Schema 演进、SQL 带到对象存储上,逐步统一两者。这一章讲清它凭什么成立。
数据湖:廉价存一切,但"暗物质"警告
数据湖把原始数据(日志、图片、任意格式)堆在对象存储(S3/OSS),按字节算钱、无限扩容。但只存不管会变成"数据沼泽":谁也说不清里面有什么、能不能信。所以光有湖不够,得有表格式和治理。
论为什么"便宜存储"会反噬
因为太便宜,团队无节制地堆数据,三年后存储费不低、且大量无人知晓的"暗物质"带来合规与质量风险。湖的价值要靠表格式(可查询、可演进)+ 元数据目录 + 质量门禁兑现,否则就是昂贵的垃圾场。
值得一提的是,数据湖并不等于"随便堆文件"。要发挥价值,它得配合 分区策略、列存文件格式(Parquet/ORC)、以及元数据管理;否则就退化成开篇说的"数据沼泽"。湖仓一体正是把这三件套工程化、产品化的结果,让"便宜存储"真正变成"可查询、可演进的资产"。
湖仓一体与开放表格式:Iceberg / Delta / Hudi
它们把一份 Parquet 文件"包装成"带 ACID 的表:用元数据文件记录哪些数据文件属于当前快照、谁增谁删。Iceberg 引擎中立(Spark/Flink/Trino 都能读)、隐藏分区友好;Delta 与 Spark 生态紧耦合;Hudi 擅长流式 upsert。三者让"对象存储上的表"获得仓的能力。
| 表格式 | 强项 | 典型生态 |
|---|---|---|
| Iceberg | 引擎中立、隐藏分区、海量分区 | Spark/Flink/Trino/DuckDB |
| Delta Lake | 与 Spark 一体、易上手 | Spark/Databricks |
| Apache Hudi | 流式 upsert、近实时入库 | Spark/Flink |
Schema 演进:加列/改类型不破坏旧数据
表格式支持Schema 演进:加可选列、改名(带 alias)、改类型(兼容方向)。旧数据没这列就补默认值,新查询照常跑。这是 ELT 能"反复重算"的底层保障——你不用因为加字段就重建整张表。
| 变更 | 是否安全 | 推荐做法 |
|---|---|---|
| 加可空列 | 安全 | 直接加,旧数据补 null |
| 类型放宽(int→long) | 安全 | 兼容方向允许 |
| 删列 / 改类型 | 破坏性 | 新表 + 双写过渡 |
不是所有变更都免费:把 String 改成 Int 会读不出旧数据;删列虽允许但下游用了就崩;改分区键要重写数据。演进前先查"兼容性矩阵",破坏性变更走新表 + 双写过渡,别硬改生产表。
对象存储上的分区布局与时间旅行
湖上数据按分区目录组织(如 /dt=2026-09-15/region=sh/),查询时只扫相关分区(分区裁剪),省钱省力。时间旅行(Time Travel):表格式记录快照,你能查询"昨天那个版本的数据",用于审计、回滚、复现训练集。
论时间旅行为什么对 AI 是刚需
模型训练集必须可复现:三个月后想复现当时那版模型,就得拉"当时那一刻"的数据快照,而不是今天被改过的。时间旅行给了你"数据版的 git checkout",是机器学习可复现性的数据底座。
小文件问题与 Compaction
流往湖里持续写,会产生海量小文件(每个 checkpoint 一小批)。小文件多了,查询时要打开成千上万个文件、元数据膨胀,性能雪崩。解法:后台 Compaction(合并)把小文件 rolling 合成大文件,Iceberg 自带 rewrite 工具。
不合并,三个月后一张表上百万小文件,Trino 列扫描直接超时。Compaction 要定期跑且避开高峰期(它自己也消耗算力)。把它当成和"数据新鲜度"同等重要的运维动作,写进调度。
统一查询引擎:Trino / DuckDB 读同一份湖
湖仓一体的一大红利:多个引擎读同一份 Iceberg 表。重量级用 Trino(Presto 系,分布式交互查询),轻量级/本地用 DuckDB(进程内,秒开小数据)。分析师用 Trino 跑即席 SQL,数据科学用 DuckDB 本地拉样本,无需各自导出。
论一份数据,多种算力
传统数仓数据锁死在自家引擎里;湖仓让数据"开放存储 + 开放格式",谁都能读。这意味着你不会被单一引擎绑架,也能按场景选最快最便宜的那个。开放格式(Parquet + Iceberg 元数据)是这自由的根。
隐藏分区与 partition transform:查询不用懂底层路径
传统按 /dt=2026-09-15/ 显式分区,查询要写对路径、且分区字段和存储强耦合——分析师忘了带 dt 条件就全表扫。Iceberg 的 partition transform 让分区成为"逻辑声明":你定义"按 event_time 的月/天/小时分""按 region 分桶",引擎自动把数据落到对应文件,并做隐藏分区裁剪。
关键点:查询写业务条件(WHERE event_time > ...)也能自动命中分区,不用再显式写分区列。分区演进也更自由——改 transform 不影响已有查询写法,老数据按旧规则、新数据按新规则,平滑过渡。
论为什么"隐藏分区"比"显式分区"省心
显式分区下,分区策略是"潜规则",错了就全表扫、成本高还难发现。隐藏分区把"怎么分"和"怎么查"解耦:按业务时间查,引擎自己知道扫哪些文件,分析师不必背分区规则。这也是湖仓一体在"易用性"上压传统数仓分区的典型体现。
数据质量与治理:谁能用、怎么脱敏
数据规模一大,就必须管:质量(这数据能信吗)、血缘(这列从哪来)、元数据(这是什么)、权限与脱敏(谁能看、怎么藏)。这既是工程问题,也是合规问题,更是 AI 项目能不能用的底线。
数据质量维度与质量门禁
质量不是"有没有数据",而是多维度的:完整性(非空率)、准确性(是否真实)、一致性(跨表口径)、时效性(延迟多久)、唯一性(是否重复)。把它们做成门禁(gate):不达标就阻断下游,避免脏数据扩散。
| 质量维度 | 怎么查 | 不达标后果 |
|---|---|---|
| 完整性 | COUNT NULL / 总行数 | 模型特征缺值 |
| 准确性 | 与源系统对账 | 指标虚高/虚低 |
| 一致性 | 跨表汇总比对 | 报表对不上 |
| 时效性 | 最新分区时间差 | 用了过期数据 |
论为什么数据质量对 AI 特别致命
一条脏数据进训练集,模型行为难解释、难归因,且污染是"无痕"的。必须在入湖/入仓处设质量门禁(空值率、schema 漂移、分布异常告警),否则下游模型与报表全盘失真。这与 FDE 模块 08"数据质量是大前提"一脉相承。
数据血缘:这列到底从哪来
血缘(Lineage)记录数据从源到消费的来龙去脉:这张表由哪些上游加工、这列由哪个字段变来。出事了能快速定位"哪个环节污染了数据",做影响分析时也知道"改这个字段会波及谁"。
论血缘是治理的"地图"
没有血缘,数据团队像在黑屋找线头:一个指标异常,要人肉翻十几个 SQL 才知源头。血缘自动化采集(从 SQL/任务依赖推断)后,点一下就能看到全链路,排障从小时级降到分钟级,也是合规审计的硬要求。
元数据与数据目录
元数据是"关于数据的数据":表叫什么、字段含义、负责人、更新频率、敏感级别。数据目录(Data Catalog)把这些集中管理,让人能搜到、看懂、放心用(如 Unity Catalog、DataHub、Amundsen)。它让"数据资产"真正可被发现。
敏感数据分级与脱敏
手机号、身份证、人脸、银行卡都是 PII(个人敏感信息)。治理要:① 标记敏感级别;② 按角色控制谁能看明文;③ 动态脱敏(查询时按需遮蔽,如 138****8000);④ 进出湖都留审计。脱敏别只靠"约定",要有技术强制。
手机号、身份证、人脸等进湖前必须标记与脱敏。等保/PIPL/GDPR 都盯这块,漏一处就是客户现场的重大风险。常见翻车:脱敏只做了展示层、但原始表还能直查;或测试环境用生产真实数据忘了脱敏。分级、脱敏、审计三件套缺一不可。
合规:PIPL / GDPR / 等保的硬约束
合规不是法务的事,是数据工程要落地的:最小必要(别多采)、授权可撤回(用户删除要真能删)、跨境/留存期限(到期清理)、可审计(谁何时查了什么)。这些约束直接决定你的湖仓要支持哪些能力(如行级删除、留存 TTL)。
论把合规当"架构特性"而非"事后补丁"
等出事再补脱敏、补删除,成本和风险都高。从建湖起就把敏感分级、动态脱敏、行级权限、留存清理、访问审计设计进去——它们是和"能查出来"同等重要的能力。治理前置,才不会出现"技术能做、合规不让做"的尴尬。
数据契约与 Schema Registry:把承诺管起来
Schema 是生产者对消费者的承诺。数据契约(Data Contract)把这份承诺显式化:字段、类型、语义、SLA、owner 都写清楚,变更要走评审。Schema Registry 集中存所有 Schema 的版本,生产者的写入和消费者的读取都来校验——破坏性变更在"注册"这步就被卡住。
论契约前置,别等线上炸
没有 Registry,生产者偷偷改字段类型,消费者半夜 500。有了契约 + 兼容性校验,破坏性变更在 CI/注册阶段就被拒绝,逼着你走"加字段 + 弃用 + 删除"的正规流程。它是数据治理从"人盯人"升级到"系统卡点"的关键一环。
数据工程把原始脏数据变成可信资产:先分清批/流节奏,用 ETL(中间转换)或 ELT(入湖再转,主流)管线,抽取用全量/增量/CDC,加载讲幂等;数仓用维度建模(事实表记度量、维度表记描述,星型优于雪花,SCD Type 2 留历史,粒度与一致性维度是翻车重灾区);批处理靠 Spark 的 DataFrame + Catalyst,优化重心在分区与 shuffle(reduceByKey 优于 groupByKey、broadcast 小表、按需缓存);流处理用 Kafka 做可重放管道 + Flink 做逐事件计算,核心是事件时间/水印/状态/Checkpoint 与端到端 Exactly-once;湖仓一体靠 Iceberg/Delta 开放表格式带来 ACID、Schema 演进与时间旅行;最后以质量门禁、血缘、元数据目录、分级脱敏与合规为底线,把干净可复现的数据供给 AI 模型。
1.ETL 和 ELT 的核心区别?为什么现在 ELT 更主流?
查看答案
ETL 在加载前于中间层转换,ELT 先原样入湖/仓再在仓内转换。ELT 主流因为:存储便宜、现代引擎算力强、原始数据留在湖里可反复换口径重算,不用每次重新抽取。适合数据种类多、转换逻辑常变的场景。
2.星型模型和雪花模型怎么选?为什么实际几乎都选星型?
查看答案
星型事实直连扁平维度,少 join、查询快、业务方易懂;雪花把维度再规范化拆子维度,省空间但多 join、慢、难懂。列存下星型冗余成本极低,换不来雪花的查询体验损失,所以星型是默认答案。
3.Spark 里 reduceByKey 为什么比 groupByKey 更适合大数据聚合?
查看答案
groupByKey 把全量 value 原样 shuffle 后聚合,网络搬运量大;reduceByKey 先在每个分区本地做 combiner 预聚合,再 shuffle 减少后的结果,网络量小得多。这是"本地先合"的 shuffle 优化总原则。
4.Flink 的 Exactly-once 到底保证什么?它真的是"每条只算一次"吗?
查看答案
不是物理上的每条只算一次(那不可能),而是"最终结果等效于精确一次":靠 checkpoint 恢复 + 两阶段提交(sink 与 checkpoint 对齐)实现。端到端还要 sink 支持事务。轻量场景用 at-least-once + 下游去重往往更划算。
5.数据湖的"时间旅行"对 AI 项目为什么是刚需?
查看答案
模型训练集必须可复现:三个月后想复现当时那版模型,得拉"当时那一刻"的数据快照而非今天被改过的。时间旅行给了数据版的 git checkout,是机器学习可复现性的底座。
6.为什么 PII 漏管是合规事故而非小事?
查看答案
手机号、身份证、人脸等进湖后若未分级、未脱敏、未审计,违反 PIPL/GDPR/等保,且常见翻车是脱敏只做展示层但原始表可直查、或测试环境用生产真实数据。分级、脱敏、审计三件套缺一不可,应作为架构特性前置设计。
下一步往哪走
路学完数据工程之后
① 串 AI:回 ai.html / tech-python-ml,你现在是它们的"数据供给方"——特征仓库让线上线下特征一致。
② 串 FDE:模块 08 数据链路与可观测,本页是它的技术底座;模块 07 脱敏,本页第 6 章给了落地。
③ 动手:用 Docker 起 Kafka + Spark/Flink,跑一条"日志 → 清洗 → 聚合 → Iceberg"的最小管道,故意触发一次 shuffle 慢查询再优化。