本页定位 · Data Engineering

AI 项目的命门不是模型,是数据。FDE 模块 08 提到了 ETL/ELT,但没展开。本页补齐数据工程与大数据的完整知识链:从数据采集与编排、数仓维度建模、Spark/Flink 批流处理,到数据湖/湖仓一体,再到数据质量与治理。我会像讲《Java 进阶》那样,把每个知识点拆成"它是什么 → 为什么这样设计 → 怎么用 → 踩过什么坑",让你学完能搭一条"原始日志 → 干净表 → 可被模型/报表消费"的可靠数据链路。基线:Spark 3.x / Flink 1.18+ / Iceberg 1.x。

1

数据采集与 ETL/ELT:把"脏数据"变成"可信资产"

Ingestion · ETL/ELT · Orchestration

数据工程是"后端的数据侧"——核心职责是采集、清洗、建模、供给。模型再聪明,喂的是垃圾数据也只能产出垃圾结论。这一章讲清数据进门的姿势、转换到底发生在哪一步、以及怎么让任务按时按依赖跑。

批处理 vs 流处理:先分清你要哪种节奏

数据进来有两条路。批处理是"攒一堆再算"(小时级、天级),适合报表、训练样本;流处理是"来一条算一条"(毫秒~秒级),适合风控、实时监控。选错节奏,系统会又贵又脆。

维度批处理 Batch流处理 Streaming
触发方式定时/手动触发一批事件触发,持续运行
延迟分钟~天毫秒~秒
典型引擎Spark、MapReduceFlink、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,对源库压力小、延迟低。

-- 增量抽取:只取上次水位之后的变化(用 updated_at 游标) SELECT * FROM orders WHERE updated_at > '2026-09-14 23:00:00' -- 水位线,存在调度表里 AND updated_at <= '2026-09-15 23:00:00'; -- 闭区间,避免重复/漏

清洗与转换(Transform):把脏数据变干净

原始数据几乎一定脏:空值、类型不一致("123" 和 123 混着)、编码乱、重复、时区错。转换层要做的清单:去重、类型标准化、空值填充/标记、单位统一、字段拆分与拼合、非法值过滤、时区归一。

# PySpark 里做一层典型清洗 from pyspark.sql import functions as F clean = (raw .where(F.col("user_id").isNotNull()) # 1. 丢掉主键为空 .withColumn("amount", F.col("amount").cast("double")) # 2. 类型标准化 .withColumn("ts", F.to_utc_timestamp(F.col("ts"), "Asia/Shanghai")) # 3. 时区归一 .dropDuplicates(["order_id"]) # 4. 去重 .fillna({"channel": "unknown"})) # 5. 空值兜底

加载(Load):全量、增量、CDC 怎么落

清洗完要落库。三种策略:全量覆盖(小表、重算安全);增量追加(按日期分区追加,可回刷);Merge/Upsert(按主键更新已有、插入新的,CDC 落湖仓常用)。湖仓表格式(Iceberg)原生支持 MERGE INTO。

-- Iceberg 表的 Upsert:有则更新、无则插入 MERGE INTO dw.orders t USING staging.orders_delta s ON t.order_id = s.order_id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT *;

论为什么加载也要"可重跑"

调度任务今天挂了、明天重跑,必须幂等——跑一次和跑一百次结果一样。增量追加如果不带去重,重跑会重复;全量覆盖最安全但慢;Upsert 兼顾。设计管道时把"重跑幂等"当铁律,半夜告警才不会被自己坑。

调度编排:Airflow/Dagster 把任务串成 DAG

数据任务有强依赖:先抽数 → 再清洗 → 再聚合 → 再导出。编排工具把这套依赖画成 DAG(有向无环图),自动按时调度、失败重试、可观测。Airflow 用 Python 定义 DAG,生态最大;Dagster 强调"资产(asset)"语义,更现代。

# Airflow 风格的 DAG(示意) with DAG("etl_pipeline", schedule="@daily", catchup=False) as dag: extract = PythonOperator(task_id="extract", python_callable=pull) clean = PythonOperator(task_id="clean", python_callable=wash) load = PythonOperator(task_id="load", python_callable=save) extract >> clean >> load # 定义依赖顺序
编排工具核心模型适合谁
AirflowDAG + 任务,生态最大通用批调度、团队已熟悉
Dagster资产(asset)语义,类型感知想让数据资产可追踪的新项目
PrefectPython 原生 flow,混合调度偏 Python 的工程团队
新手最爱踩的调度坑

① catchup 没关:任务停了三天再开,Airflow 默认补跑这三天,瞬间压垮源库。② 任务不幂等:重跑就重数据。③ 把重活塞进一个 task:一个 Python 函数干完所有事,失败全重来、不好并行。拆成有依赖的小 task,可重试、可观测。

任务监控与 SLA:失败重试、告警与数据新鲜度

管道跑起来只是开始,跑得稳、出事有人知才是生产。调度系统要配三件事:重试(瞬时失败自动重试,配指数退避)、告警(任务失败/超时/数据空了立刻通知)、SLA 监控(这张表今天该几点就绪,晚了就报)。

# Airflow:任务级重试与超时(在算子上声明) clean = PythonOperator( task_id="clean", python_callable=wash, retries=3, # 失败自动重试 3 次 retry_delay=timedelta(minutes=5), # 退避间隔 execution_timeout=timedelta(hours=1),# 跑超 1h 判失败 sla=timedelta(hours=2)) # 超过 SLA 触发告警
重试不是万能药

重试只该用于幂等且瞬时的失败(网络抖、源库锁)。若任务本身有 bug(SQL 写错),重试只会重复失败、浪费资源、还可能重复写。先分清失败类型:非幂等写操作重试前要确认"会不会重复落库"。

2

数仓与建模:维度建模让分析变简单

Dimensional Modeling · Star/Snowflake · SCD

数仓不是"把业务库复制一份"。业务库为写入事务优化(范式化、避免冗余),数仓为分析读取优化(反范式、预聚合)。这一章的明星方法是维度建模——让业务方也能看懂的"事实 + 维度"结构。

维度建模思想:面向分析而非事务

事务型建模回答"现在库存多少、订单状态如何";维度建模回答"去年华东区手机品类的销售额按周趋势"。后者需要把"销售额(可度量)"和"地区/品类/时间(怎么切)"分开。核心就两类表:事实表记"发生了什么(数字)",维度表记"这件事的背景(描述)"。

论为什么数仓要反范式

业务库范式化(拆成很多小表)是为了写入不冗余、不乱。但分析时每查一个维度都要 join 一堆表,慢且难懂。数仓把常用描述冗余进宽表/维度表,用空间换查询简单和速度。记住:数仓的设计目标是"人能看懂、查得快",不是"不冗余"。

事实表与维度表:数仓的两类主角

事实表(Fact)存业务过程的可度量事件,如"一笔订单""一次点击",包含外键(指向维度)和度量(金额、数量)。维度表(Dimension)存描述性属性,如"商品""用户""时间""地区",基本不会变、且常被当作分析切片的依据。

特征事实表 Fact维度表 Dimension
存什么事件 + 度量值(金额/次数)描述属性(名称/类别/地区)
规模巨大、持续增长较小、变化慢
主键业务键 + 代理键代理键(surrogate key)
例子销售事实、点击流事实商品维、时间维、用户维

一个常见误区:把事实表和维度表都建得"又宽又全"。事实表每多一个维度外键就多一次 join,维度表塞太多描述列会膨胀。经验是事实表只放可加度量 + 必要外键,描述性的都交给维度表,保持事实表的"瘦"。

论粒度决定事实表能回答什么

同样一笔销售,按"订单"粒度能算订单数,按"订单明细行"粒度才能算件数和品类分布。粒度定得太粗(只到订单),很多下钻分析就做不了;定得太细(到每个 SKU),存储和查询成本又高。先想清楚业务要下钻到哪一层,再钉粒度,这是建模第一步。

星型 vs 雪花模型

星型(Star):一个事实表直接连多个维度表,维度不再拆分,像星星。简单、查询快、最常用。雪花(Snowflake):维度再规范化拆成子维度(如"城市维"拆出"省份维")。省空间但 join 多、查询慢。实际几乎都用星型,雪花只在维度极度冗余时偶尔用。

维度星型雪花
结构事实直连扁平维度维度再拆子维度
查询性能快(少 join)慢(多 join)
冗余有冗余少冗余
易用性高,业务方易懂低

论为什么实际几乎都选星型

现代数仓(尤其是列存 + 缓存维度)下,星型的冗余成本极低,而雪花多出来的 join 在分析师自助查询时极其劝退。除非维度表大到离谱,否则星型是默认答案。雪花那点空间收益,换不来查询体验和性能的损失。

缓慢变化维(SCD):历史怎么留

维度会变:顾客搬了城市、商品改了类目。要不要保留"当时"的样子?这决定指标的可比性。SCD Type 1:直接覆盖(丢了历史);Type 2:加生效/失效时间,保留全量历史版本(最常用);Type 3:加一列存"上一版值"(只留最近一次)。

-- SCD Type 2:同一用户两行,带生效区间 INSERT INTO dim_user (user_id, city, valid_from, valid_to, is_current) VALUES ('u1', '上海', '2024-01-01', '9999-12-31', true); -- 用户搬家后:旧行封口,新行生效 UPDATE dim_user SET valid_to='2026-03-01', is_current=false WHERE user_id='u1' AND is_current=true; INSERT INTO dim_user VALUES ('u1', '北京', '2026-03-02', '9999-12-31', true);
SCD 忘做会导致"指标对不上"

如果只用 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 报表金额凭空翻倍。这是多值维度建模最容易翻车的地方,务必在建模规范里写死。

3

批处理引擎:Spark 的 RDD、DataFrame 与优化

Spark · RDD · DataFrame · Shuffle

Spark 是批处理事实标准。它把计算表达成有向无环图(DAG),在内存里算、容错靠血缘重算。理解它的抽象层级和性能黑洞(shuffle),才能写出既对又快的作业。

RDD:Spark 的底层抽象

RDD(弹性分布式数据集)是不可变、分区的数据集合,是 Spark 最原始的 API。它记录"怎么算出来"的血缘(lineage),某分区丢了能重算而非重跑全部。但 RDD 不感知数据 schema,优化空间小。

# RDD 风格:函数式算子,靠 JVM 对象,无 schema 优化 rdd = sc.textFile("s3://raw/logs/*.gz") counts = (rdd .map(lambda line: line.split("\t")[0]) # 取第一列 .map(lambda u: (u, 1)) .reduceByKey(lambda a, b: a + b)) # 触发 shuffle

论为什么日常几乎不用裸 RDD

RDD 把数据当不透明的 Java 对象,无法做列式优化、无法 Catalyst 优化器重写。DataFrame/Dataset 带 schema,优化器能"看懂"你的意图做谓词下推、列裁剪。除非要极细粒度控制(自定义分区、特殊二进制),否则一律用 DataFrame。

补充一点:RDD 的容错靠 血缘(lineage)——每个 RDD 记得自己是怎么从父 RDD 算来的。某个分区丢失,Spark 沿血缘重算那一条路径,而不是重跑整个作业,这正是它比早期 MapReduce 每次落盘更优雅的地方。理解血缘,就理解了 Spark"算错了能自愈"的底气。

DataFrame/Dataset:声明式、更快

DataFrame 是带 schema 的分布式表,Dataset 是带类型安全的版本(Scala/Java)。你写"我要什么",Catalyst 优化器决定"怎么高效算"。绝大多数作业用它,性能和可维护性都碾压 RDD。

# DataFrame 风格:声明式 SQL 语义,Catalyst 自动优化 result = (spark.read.parquet("s3://lake/orders") .where("amount > 0") # 列裁剪 + 谓词下推 .groupBy("region") .agg(F.sum("amount").alias("gmv"))) result.show()

分区与并行度:任务切多细

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 优化的总原则。

shuffle 分区数默认 200 的坑

Spark 的 spark.sql.shuffle.partitions 默认 200。对几 MB 的小数据集,200 个分区几乎全空跑,调度开销比计算还大;对超大集群又太少,单分区过大撑爆内存。按"单分区 ~128MB"目标手动调,或开启 AQE(自适应查询执行)让它自动合并小分区、动态优化,省去手调。

宽依赖、窄依赖与算子选择

窄依赖:一个父分区只给一个子分区(map、filter),可流水线、某分区丢了自己重算。宽依赖:一个父分区给多个子分区(groupBy、join),必须等上游全就绪、产生 shuffle、阶段(stage)边界。Spark 在宽依赖处切 stage。

# 用 DataFrame 的 join 而不是自己 reduceByKey 拼, # Catalyst 会选择 broadcast join / sort-merge join small = spark.read.table("dim_user") # 小表 big = spark.read.table("fact_order") # 大表 joined = big.join(broadcast(small), "user_id") # 小表广播,避免大表 shuffle

内存与缓存:persist/cache 用对才提速

同一个 DataFrame 被反复用(多次 join、多次聚合),用 cache()/persist() 把它留在内存/磁盘,避免重算血缘。但缓存不是免费的:占内存、要治理 eviction。

# 一次性物化,后续复用不再重算 popular = spark.read.table("fact_click").filter("cnt > 100").persist(StorageLevel.MEMORY_AND_DISK) a = popular.groupBy("page").count() # 第一次触发计算并缓存 b = popular.filter("page = 'home'") # 直接读缓存 popular.unpersist() # 用完释放,别占着
缓存用错的两种浪费

① 缓存一次性的中间结果:只用一个 action 就丢,缓存反而拖慢(多一次物化)。② 忘记 unpersist:缓存堆积挤占执行内存,触发频繁 spill 到磁盘,比不缓存还慢。规则:被 ≥2 个 action 复用的才缓存,用完即 unpersist。

广播变量、累加器与 UDF:Spark 的共享机制

算子里的函数跑在各个 executor 上,怎么共享"只读大对象"(如一份字典)或"全局计数"?广播变量(broadcast)把只读数据发到各节点一次,避免每个 task 都序列化一份;累加器(accumulator)做分布式计数/求和(只在 driver 可读最终结果);UDF把自定义函数注册给 DataFrame 用。

# 广播小字典,避免每个 task 重复传大对象 mapping = spark.sparkContext.broadcast({"A": "苹果", "B": "香蕉"}) result = df.rdd.map(lambda r: (mapping.value[r["code"]], r["amt"])) # 累加器:分布式计数异常行(只在 driver 端可读汇总) bad = spark.sparkContext.accumulator(0) df.foreach(lambda r: bad.add(1) if r.amount < 0 else None)

论广播 vs 直接闭包捕获

小数据直接闭包捕获也行,但大对象(几 MB 的映射表)若每个 task 都序列化一次会爆炸。broadcast 走高效分发、各节点存一份,省网络和内存。经验:大于几 KB 的闭包捕获对象,改用 broadcast。

4

流处理:Flink 与 Kafka 的实时世界

Flink · Kafka · Window · Exactly-once

流处理要回答"此刻正在发生什么"。Kafka 是流的管道与存储,Flink 是真正的逐事件计算引擎。这一章讲清时间语义、窗口、状态与一致性——这些是流和批最不一样的地方。

Kafka:流的"管道"与"存储"

Kafka 用主题(topic)分频道,每个 topic 分多个分区(partition)保证有序与并行,数据按偏移量(offset)持久化、可被消费者重复读。它既是消息队列,也是短期日志存储——这正是 CDC 和流处理的地基。

# 生产:把订单事件写进 topic(key 用 order_id 保证同键同分区有序) producer.send(ProducerRecord("orders", order_id, json)) # 消费:按组(consumer group)分摊分区,用 offset 标记进度 for msg in consumer: process(msg.value) # 处理事件 consumer.commit() # 提交 offset = 标记"已消费"

论为什么 Kafka 能做"重放"

传统 MQ(如早期 RabbitMQ)消费完就删,出错难补救。Kafka 把消息按 offset 保留(可配保留期),消费者想重来就从旧 offset 重读。这让它既是传输管道,又是"可重放的日志",是流处理容错和回溯的关键。

消费者组再均衡会短暂停消费

消费者上线/下线/崩溃会触发 rebalance(再均衡),期间整个组暂停消费、重分配分区,频繁再均衡会让吞吐剧烈抖动,监控上表现为周期性消费停滞。用 CooperativeStickyAssignor 做增量再均衡、调大 session.timeout、避免在处理里长阻塞,能显著减少停顿。

Flink 编程模型:Source → Transform → Sink

Flink 作业是"从哪读(Source)→ 怎么算(Transform)→ 写到哪(Sink)"的 DAG。它把流抽象成 DataStream,算子链(operator chain)把相邻算子拼一起减少线程切换。

// Flink Java:实时统计每用户的点击数 DataStream<Click> clicks = env.addSource(new KafkaSource<>(...)); clicks.keyBy(c -> c.userId) .window(TumblingEventTimeWindows.of(Time.minutes(5))) .aggregate(new CountAgg()) .addSink(new JdbcSink<>(...)); // 写回数据库 env.execute("user-click-count");

时间语义与窗口:事件时间 + 水印

流里"时间"有三种:事件时间(事件发生时刻,最准)、摄入时间(进 Flink 时)、处理时间(算的时候)。真实数据会乱序、会延迟,所以用事件时间 + 水印(Watermark):水印说"事件时间已到 T,早于 T 的我都当到齐了",窗口才关、才出结果。

// 允许 5 秒乱序:水印 = 最大事件时间 - 5s WatermarkStrategy.<Click>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((c, ts) -> c.eventTime); // 窗口类型:滚动(无重叠)/滑动(有重叠)/会话(按静默间隔)
水印设太小 vs 太大

太小(如 0):稍晚到的数据被当"迟到"丢弃,结果不准。太大:窗口迟迟不关,结果延迟产出、状态堆积。设多少取决于数据的真实乱序程度——从日志观察 p99 延迟再定,别拍脑袋。

状态管理与 Checkpoint

流计算要"记住点东西":累加的计数、去重的集合、窗口里攒的数据,这些叫状态(State)。状态存在哪?Checkpoint 是 Flink 的分布式快照,定期把状态持久化到存储。作业挂了,从最近 checkpoint 恢复,状态不丢。

// 开启精准的 checkpoint(屏障对齐,精确一次的基础) env.enableCheckpointing(5000); // 每 5s 一次 env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(2000); env.setStateBackend(new RocksDBStateBackend("s3://ckpt/")); // 大状态用 RocksDB

论状态存在内存还是 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 机制,只是换了个入口。

-- 把点击流按 5 分钟窗口统计,结果写回另一个 topic/表 INSERT INTO user_click_cnt SELECT user_id, TUMBLE_START(ts, INTERVAL '5' MINUTE) AS w, COUNT(*) AS cnt FROM clicks GROUP BY user_id, TUMBLE(ts, INTERVAL '5' MINUTE);

论SQL 与 DataStream 怎么选

常规聚合、join、维表关联用 SQL 足够且快;涉及自定义状态机、复杂事件模式(CEP)或精细控制算子链时,回到 DataStream API。两者可混用:SQL 起大部分活,特殊逻辑用 DataStream 补。

Kafka 事务性 Producer:端到端精确一次的源头

Flink 的精确一次需要"源"也配合。Kafka 事务性 Producer 能把"向多个分区/主题写消息"打包成一个原子事务:要么全部对消费者可见,要么全不可见(靠 isolation.level=read_committed 隔离)。CDC 入湖时常用它,保证"这一批变更要么都进、要么都不进",不会出现"一半写入、一半没写"的中间态。

它和幂等 Producer 常被混为一谈,但层级不同:幂等只解决单分区内重试导致的重复(靠序列号去重);事务进一步保证跨分区原子可见,是解决"消费-转换-生产"链路精确一次的关键。两者通常一起开,幂等是事务的底层基础。

# 开启事务:配 transactional.id,先 initTransactions props.put("enable.idempotence", "true"); props.put("transactional.id", "order-cdc-1"); producer.initTransactions(); producer.beginTransaction(); producer.send(rec1); producer.send(rec2); # 一批原子写入 producer.commitTransaction(); # 提交后才对 read_committed 消费者可见

论为什么 CDC 入湖离不开事务

设想一次 binlog 抓了 1000 条变更,写到第 500 条时作业挂了:非事务下,下游看到"残缺的一半",对账直接错乱。事务保证"要么 1000 条全可见、要么 0 条",配合 Flink 从 checkpoint 重放,才能做到端到端不重不漏。理解这层,才懂精确一次为什么需要整条链路协同。

5

数据湖与湖仓一体:让湖也能做仓的活

Data Lake · Lakehouse · Iceberg/Delta

数仓贵且锁死结构,数据湖便宜但"存了等于没管"。湖仓一体(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 能"反复重算"的底层保障——你不用因为加字段就重建整张表。

# Iceberg:安全加一列(旧文件该列为 null,不破坏历史) spark.sql("""ALTER TABLE lake.orders ADD COLUMN coupon_code STRING AFTER amount""") # 类型放宽也是允许的(int -> long),收窄则要 migrate
变更是否安全推荐做法
加可空列安全直接加,旧数据补 null
类型放宽(int→long)安全兼容方向允许
删列 / 改类型破坏性新表 + 双写过渡
Schema 演进也有不能碰的红线

不是所有变更都免费:把 String 改成 Int 会读不出旧数据;删列虽允许但下游用了就崩;改分区键要重写数据。演进前先查"兼容性矩阵",破坏性变更走新表 + 双写过渡,别硬改生产表。

对象存储上的分区布局与时间旅行

湖上数据按分区目录组织(如 /dt=2026-09-15/region=sh/),查询时只扫相关分区(分区裁剪),省钱省力。时间旅行(Time Travel):表格式记录快照,你能查询"昨天那个版本的数据",用于审计、回滚、复现训练集。

# 查询历史快照(Iceberg 支持按版本号/时间戳) spark.sql("SELECT * FROM lake.orders VERSION AS OF 42") # 第 42 个快照 spark.sql("SELECT * FROM lake.orders TIMESTAMP AS OF '2026-09-14'") # 增量读取:只拿两次快照之间新增的文件(CDC 消费) spark.read.format("iceberg").option("startSnapshotId", 41) .option("endSnapshotId", 42).load("lake.orders")

论时间旅行为什么对 AI 是刚需

模型训练集必须可复现:三个月后想复现当时那版模型,就得拉"当时那一刻"的数据快照,而不是今天被改过的。时间旅行给了你"数据版的 git checkout",是机器学习可复现性的数据底座。

小文件问题与 Compaction

流往湖里持续写,会产生海量小文件(每个 checkpoint 一小批)。小文件多了,查询时要打开成千上万个文件、元数据膨胀,性能雪崩。解法:后台 Compaction(合并)把小文件 rolling 合成大文件,Iceberg 自带 rewrite 工具。

# Iceberg:把小于 128MB 的小文件合并(Spark 行动) CALL catalog.system.rewrite_small_files( table => 'lake.orders', options => map('target-file-size-bytes', '134217728'))
Compaction 遗忘的代价

不合并,三个月后一张表上百万小文件,Trino 列扫描直接超时。Compaction 要定期跑且避开高峰期(它自己也消耗算力)。把它当成和"数据新鲜度"同等重要的运维动作,写进调度。

统一查询引擎:Trino / DuckDB 读同一份湖

湖仓一体的一大红利:多个引擎读同一份 Iceberg 表。重量级用 Trino(Presto 系,分布式交互查询),轻量级/本地用 DuckDB(进程内,秒开小数据)。分析师用 Trino 跑即席 SQL,数据科学用 DuckDB 本地拉样本,无需各自导出。

-- DuckDB 直接读湖上的 Iceberg(本地即席分析) SELECT region, SUM(amount) FROM iceberg_scan('s3://lake/orders') GROUP BY region; -- Trino 则把 Iceberg 当普通 catalog 表查,语法一致

论一份数据,多种算力

传统数仓数据锁死在自家引擎里;湖仓让数据"开放存储 + 开放格式",谁都能读。这意味着你不会被单一引擎绑架,也能按场景选最快最便宜的那个。开放格式(Parquet + Iceberg 元数据)是这自由的根。

隐藏分区与 partition transform:查询不用懂底层路径

传统按 /dt=2026-09-15/ 显式分区,查询要写对路径、且分区字段和存储强耦合——分析师忘了带 dt 条件就全表扫。Iceberg 的 partition transform 让分区成为"逻辑声明":你定义"按 event_time 的月/天/小时分""按 region 分桶",引擎自动把数据落到对应文件,并做隐藏分区裁剪。

关键点:查询写业务条件(WHERE event_time > ...)也能自动命中分区,不用再显式写分区列。分区演进也更自由——改 transform 不影响已有查询写法,老数据按旧规则、新数据按新规则,平滑过渡。

# 建表时声明分区变换(隐藏分区,用户无需感知底层目录) CREATE TABLE lake.events ( id BIGINT, event_time TIMESTAMP, region STRING ) PARTITIONED BY (days(event_time), bucket(16, region)); # 查询无需写分区列,引擎自动裁剪到对应文件 SELECT * FROM lake.events WHERE event_time > TIMESTAMP '2026-09-14';

论为什么"隐藏分区"比"显式分区"省心

显式分区下,分区策略是"潜规则",错了就全表扫、成本高还难发现。隐藏分区把"怎么分"和"怎么查"解耦:按业务时间查,引擎自己知道扫哪些文件,分析师不必背分区规则。这也是湖仓一体在"易用性"上压传统数仓分区的典型体现。

6

数据质量与治理:谁能用、怎么脱敏

Quality · Lineage · Metadata · Compliance

数据规模一大,就必须管:质量(这数据能信吗)、血缘(这列从哪来)、元数据(这是什么)、权限与脱敏(谁能看、怎么藏)。这既是工程问题,也是合规问题,更是 AI 项目能不能用的底线。

数据质量维度与质量门禁

质量不是"有没有数据",而是多维度的:完整性(非空率)、准确性(是否真实)、一致性(跨表口径)、时效性(延迟多久)、唯一性(是否重复)。把它们做成门禁(gate):不达标就阻断下游,避免脏数据扩散。

质量维度怎么查不达标后果
完整性COUNT NULL / 总行数模型特征缺值
准确性与源系统对账指标虚高/虚低
一致性跨表汇总比对报表对不上
时效性最新分区时间差用了过期数据
# Great Expectations 风格:把期望写成可执行的"契约" expect orders.amount to be_between(0, 1_000_000) expect orders.user_id to not_be_null expect table to have_rows # 空表 = 抽数挂了,立刻告警

论为什么数据质量对 AI 特别致命

一条脏数据进训练集,模型行为难解释、难归因,且污染是"无痕"的。必须在入湖/入仓处设质量门禁(空值率、schema 漂移、分布异常告警),否则下游模型与报表全盘失真。这与 FDE 模块 08"数据质量是大前提"一脉相承。

数据血缘:这列到底从哪来

血缘(Lineage)记录数据从源到消费的来龙去脉:这张表由哪些上游加工、这列由哪个字段变来。出事了能快速定位"哪个环节污染了数据",做影响分析时也知道"改这个字段会波及谁"。

论血缘是治理的"地图"

没有血缘,数据团队像在黑屋找线头:一个指标异常,要人肉翻十几个 SQL 才知源头。血缘自动化采集(从 SQL/任务依赖推断)后,点一下就能看到全链路,排障从小时级降到分钟级,也是合规审计的硬要求。

元数据与数据目录

元数据是"关于数据的数据":表叫什么、字段含义、负责人、更新频率、敏感级别。数据目录(Data Catalog)把这些集中管理,让人能搜到、看懂、放心用(如 Unity Catalog、DataHub、Amundsen)。它让"数据资产"真正可被发现。

# 给表/列打标签与说明(Unity Catalog 风格),治理从这里开始 COMMENT ON TABLE dw.orders IS '订单事实表,粒度=订单明细行'; ALTER TABLE dw.orders SET TAGS ('pii:phone', 'domain:sales'); # 列级打标,脱敏策略才能按标签自动套用

敏感数据分级与脱敏

手机号、身份证、人脸、银行卡都是 PII(个人敏感信息)。治理要:① 标记敏感级别;② 按角色控制谁能看明文;③ 动态脱敏(查询时按需遮蔽,如 138****8000);④ 进出湖都留审计。脱敏别只靠"约定",要有技术强制。

# 动态脱敏:非授权角色查 phone 时自动遮蔽 CREATE MASKING POLICY phone_mask AS (val STRING) RETURNS STRING -> CASE WHEN is_role_in('data_engineer') THEN val ELSE regexp_replace(val, '^(\d{3}).*(\d{4})$', '\1****\2') END; ALTER TABLE dw.user MODIFY COLUMN phone SET MASKING POLICY phone_mask;
PII 漏管 = 合规事故

手机号、身份证、人脸等进湖前必须标记与脱敏。等保/PIPL/GDPR 都盯这块,漏一处就是客户现场的重大风险。常见翻车:脱敏只做了展示层、但原始表还能直查;或测试环境用生产真实数据忘了脱敏。分级、脱敏、审计三件套缺一不可。

合规:PIPL / GDPR / 等保的硬约束

合规不是法务的事,是数据工程要落地的:最小必要(别多采)、授权可撤回(用户删除要真能删)、跨境/留存期限(到期清理)、可审计(谁何时查了什么)。这些约束直接决定你的湖仓要支持哪些能力(如行级删除、留存 TTL)。

论把合规当"架构特性"而非"事后补丁"

等出事再补脱敏、补删除,成本和风险都高。从建湖起就把敏感分级、动态脱敏、行级权限、留存清理、访问审计设计进去——它们是和"能查出来"同等重要的能力。治理前置,才不会出现"技术能做、合规不让做"的尴尬。

数据契约与 Schema Registry:把承诺管起来

Schema 是生产者对消费者的承诺。数据契约(Data Contract)把这份承诺显式化:字段、类型、语义、SLA、owner 都写清楚,变更要走评审。Schema Registry 集中存所有 Schema 的版本,生产者的写入和消费者的读取都来校验——破坏性变更在"注册"这步就被卡住。

# 生产者注册/校验 Schema(Confluent Registry 风格) curl -X POST http://registry/v1/subjects/orders-value/versions \ -H 'Content-Type: application/vnd.schemaregistry.v1+json' \ -d '{"schema":"{\"type\":\"record\",\"name\":\"Order\",...}"}' # 兼容性模式设为 BACKWARD:新 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 慢查询再优化。