bg11.数据链路
链路这一块
一个用户在电商 App 的首页看到一个商品卡片、点了一下。从这一刻起到"推荐系统下次给他推更准的商品",这条数据要经过:
产生 → 采集 → 传输 → 加工 → 存储 → 交付 → 被消费
七段。每段都有自己的技术选型、自己的失败模式、以及自己的负责人。基础篇讲的是每一段内部怎么运转;本系列讲的是这七段之间如何咬合。
flowchart LR
subgraph A["① 产生 · 客户端/服务端"]
APP["App 埋点 SDK"]
SVC["业务服务日志"]
DB[("业务数据库<br/>MySQL")]
end
subgraph B["② 采集 + ③ 传输"]
GW["上报网关"]
CDC["CDC 采集<br/>读数据库变更日志"]
MQ["Kafka<br/>消息管道"]
end
subgraph C["④ 加工"]
FLINK["Flink<br/>实时链"]
SPARK["Spark<br/>离线链"]
end
subgraph D["⑤ 存储"]
LAKE[("数据湖<br/>Iceberg on 对象存储")]
WH[("数据仓库 / OLAP<br/>ByteHouse 等")]
end
subgraph E["⑥ 交付 + ⑦ 消费"]
BI["报表 / BI"]
API["数据 API"]
FEAT["特征 → 推荐/搜索/广告模型"]
CROWD["人群包 → 营销"]
RISK["实时风控"]
end
APP --> GW --> MQ
SVC --> MQ
DB --> CDC --> MQ
MQ --> FLINK
MQ --> LAKE
LAKE --> SPARK --> LAKE
FLINK --> LAKE
FLINK --> WH
LAKE --> WH
WH --> BI
WH --> API
LAKE --> FEAT
FLINK --> FEAT
WH --> CROWD
FLINK --> RISK
抽象的图看不出重点。用刚才那次点击,看它在每一层具体长什么样。
形态一:客户端产生的原始事件(JSON)
{
"event": "product_click", // 事件名,标识"发生了什么"
"event_id": "9f3c1a7e-4b2d-11f0-a1", // 本次事件的唯一 ID,用于下游去重
"client_time": 1753846231120, // 客户端本地时间戳(毫秒),不可信,原因见 12 篇
"device_id": "d-8823ef01", // 设备标识,登录前唯一能用的身份
"user_id": 100238841, // 用户 ID,未登录时为 null
"properties": { // 事件的业务属性,不同事件字段不同
"item_id": 60012345,
"position": 3, // 在列表中的第几个坑位
"page": "home_feed",
"trace_id": "t-77aa12" // 关联同一次曝光→点击的追踪 ID
}
}
形态二:落到 ODS 层(原始层)——几乎不加工,只加采集元信息
-- ODS 层建表:一行 = 一个原始事件,原文保留
CREATE TABLE ods_app_event (
event STRING,
event_id STRING,
client_time BIGINT,
device_id STRING,
user_id BIGINT,
properties STRING, -- 整个 JSON 原样存字符串,不解析
server_time BIGINT, -- 网关收到的时间,由采集端补充
ingest_time BIGINT -- 写入数仓的时间,由加工端补充
) PARTITIONED BY (dt STRING, hour STRING); -- 按 server_time 派生的日期+小时分区
形态三:加工到 DWD 层(明细层)——解析、清洗、打宽
-- DWD 层:字段结构化 + 关联维度信息(打宽),一行仍是一个事件
CREATE TABLE dwd_item_click_di ( -- di = daily incremental,每日增量
event_time TIMESTAMP, -- 已做时间校正的业务时间
user_id BIGINT,
device_id STRING,
item_id BIGINT,
cate_id BIGINT, -- 从商品维度表关联进来
shop_id BIGINT, -- 从商品维度表关联进来
position INT,
page STRING,
trace_id STRING,
is_valid BOOLEAN -- 清洗结论:是否为有效点击(去重/反作弊后)
) PARTITIONED BY (dt STRING);
形态四:汇总到 DWS 层(汇总层)——按分析视角预先聚合
-- DWS 层:一行 = 一个"用户 × 类目 × 日期"的汇总,不再是单个事件
CREATE TABLE dws_user_cate_click_1d (
dt STRING,
user_id BIGINT,
cate_id BIGINT,
click_cnt BIGINT, -- 当日点击次数
click_item_cnt BIGINT -- 当日点击的不同商品数
) PARTITIONED BY (dt);
形态五:交付给消费方
同一份 DWS 数据,被三个下游以三种完全不同的方式取走:
-- 用法 A:分析师看报表——发散查询,今天想看类目,明天想看时段
SELECT cate_id, SUM(click_cnt) AS clicks
FROM dws_user_cate_click_1d
WHERE dt = '2026-07-29'
GROUP BY cate_id ORDER BY clicks DESC LIMIT 10;
-- 用法 B:推荐模型要特征——固定查询,但要求毫秒级、高并发
-- 结果不会留在数仓,会被推到线上 KV 存储供模型实时读取
SELECT user_id,
COLLECT_LIST(CONCAT(CAST(cate_id AS STRING), ':', CAST(click_cnt AS STRING))) AS cate_pref
FROM dws_user_cate_click_1d
WHERE dt = '2026-07-29'
GROUP BY user_id;
-- 用法 C:营销要圈人——只要一个用户 ID 集合,不要明细
SELECT user_id
FROM dws_user_cate_click_1d
WHERE dt = '2026-07-29' AND cate_id = 50008 AND click_cnt >= 3;
模拟输入与运行结果(以上述那一条 product_click 事件为例,假设当天该用户对类目 50008 共点击 3 次):
| 层 | 该用户在这一层的数据量 | 具体内容 |
|---|---|---|
| ODS | 3 行 | 3 条原始 JSON 事件,含重复上报的可能 |
| DWD | 3 行 | 3 条结构化点击明细,is_valid 已标注 |
| DWS | 1 行 | ('2026-07-29', 100238841, 50008, 3, 2) |
| 用法 A 输出 | 参与到 cate_id=50008 的聚合中,贡献 3 次点击 |
|
| 用法 C 输出 | 该 user_id 被选中(click_cnt >= 3 成立) |
为什么会是这个结果:每往上一层,行数下降、可复用性上升、灵活性下降。ODS 保留一切但没法直接用;DWS 查得飞快但只能回答"用户 × 类目 × 日"这一个视角的问题——想看"用户 × 时段"就得再建一张表。这个"行数—灵活性"的反向关系,是后面所有分层设计取舍的根源。
事实表(fact table):记录"发生了什么"的表。一行 = 一次业务事件或一个业务过程的状态。特征是行数极大、随时间持续增长、包含可累加的数值列(点击次数、金额、时长)。
维度表(dimension table):记录"参与者是谁、是什么样的"的表。一行 = 一个实体(一个用户、一个商品、一个店铺)。特征是行数相对小、变化慢、列多且以描述性文本为主。
上面的 dwd_item_click_di 是事实表;商品信息表(item_id, item_name, cate_id, brand_id, shop_id, ...)是维度表。
指标(metric / measure):一个可度量的数值结论,例如"支付 GMV"、"点击率"、"人均使用时长"。
维度(dimension,作为分析概念时):观察指标的切分视角,例如按"日期""类目""城市""渠道"来看。
⚠️ 撞名澄清:
维度这个词在本系列里有两个含义。① 上一段的"维度表"里,它指实体;② 这里指分析的切分视角。两者高度相关(维度表提供了可切分的视角),但不是一回事。本系列在可能混淆处会写成"维度表"或"分析维度"以示区分。
口径(definition / semantics):一个指标精确的计算规则——数据来源、过滤条件、聚合方式、时间归属。
口径是本系列最重要的一个词。举例,"今天的 GMV 是多少"这句话本身是不可计算的,因为至少存在这些互不相同的口径:
| 口径名 | 计算规则 | 典型使用方 |
|---|---|---|
| 下单 GMV | 所有已创建订单的金额之和,不管付没付钱 | 运营看流量转化 |
| 支付 GMV | 只算已支付成功的订单 | 财务、经营分析 |
| 支付后退款净额 | 支付 GMV 减去当期退款 | 财务结算 |
| 剔除刷单 GMV | 支付 GMV 再减去风控标记的异常订单 | 汇报口径 |
四个数可以差出 10% 以上。"指标名相同、口径不同"是数据团队最高频的事故类型
元数据(metadata):描述数据的数据。包括表有哪些字段、字段类型和含义、表的分区方式、由哪个任务产出、多久更新一次、归谁负责、数据量多大。
数据字典(data dictionary):面向人的元数据集合——把"这个字段是什么意思、取值范围是什么、什么情况下为空"写清楚,供使用者查阅。
血缘(data lineage):数据的生产依赖关系图——哪张表由哪些表经过哪个任务产出,一个字段的值来源于上游的哪些字段。
数据资产(data asset):被纳入统一管理、有明确归属人、有质量保障、可被检索和申请使用的数据表 / 指标 / 数据服务。
| 层 | 全称 | 一行代表什么 | 加工程度 |
|---|---|---|---|
| ODS | Operational Data Store,原始/贴源层 | 一条原始记录,与源端一一对应 | 几乎不加工,只补采集元信息 |
| DWD | Data Warehouse Detail,明细层 | 一个清洗后的业务事件 | 解析、清洗、去重、关联维度打宽 |
| DWS | Data Warehouse Summary,汇总层 | 一个"实体 × 分析维度 × 时间粒度"的汇总 | 按常用视角预聚合 |
| ADS | Application Data Store,应用层 | 一个直接对应某张报表/某个接口的结果 | 面向具体应用定制 |
实时链与离线链:
分叉之前:最早只有离线链——每天凌晨跑一次批任务,把昨天的数据算出来。这在报表时代够用。
痛点出现:业务开始要"现在"的数据。大促实时大屏要看当前 GMV,风控要在下单的一瞬间判断是否拦截,推荐要用用户 5 分钟前的行为。T+1 的延迟对这些场景等于无效。
于是加了一条实时链:Kafka → Flink → 实时结果。
代价立刻出现:同一个指标现在有两套代码算——离线一套 Spark SQL,实时一套 Flink SQL。两套代码、两套口径、两倍的维护成本,而且两边算出来的数几乎必然对不上。
上面这种"实时链 + 离线链并存"的形态有个名字:Lambda 架构。它由三部分组成:
- 批处理层(batch layer):全量数据,慢但准确,是最终事实来源
- 速度层(speed layer):只处理最新的增量,快但可能不准
- 服务层(serving layer):把两者的结果合并后对外提供
┌────────────────────────────┐
┌───────▶│ 批处理层 Spark(T+1,准确) │──┐
数据源 ─┤ └────────────────────────────┘ ├─▶ 服务层(合并) ─▶ 查询
│ ┌────────────────────────────┐ │
└───────▶│ 速度层 Flink(秒级,近似) │──┘
└────────────────────────────┘
Lambda 解决了什么:同时拿到低延迟和最终准确性。实时链先给一个近似值,离线链在次日用全量数据把它覆盖修正。
代价(这是重点):
- 同一份业务逻辑写两遍。批和流的算子语义不同(比如"去重"在批里是
DISTINCT,在流里是有状态的 keyed state 加超时清理),不是复制粘贴能解决的。 - 口径漂移不可避免。两套代码只要有一处理解偏差,数就对不上,而排查"实时和离线差 0.3% 是谁的问题"极其耗时。
- 两套资源、两套运维。
- 修正逻辑本身很复杂。服务层怎么判断某个时间段该用批的结果还是流的结果,是一个真实的工程难题。
Kappa 架构的主张很直接:砍掉批处理层,只保留一条流链路。需要重算历史时,就把消息队列里的数据从头重放一遍。
它成立的前提(这三条决定了 Kappa 能不能用):
- 消息队列要能保存足够长时间的历史数据——重放 3 个月,Kafka 就得存 3 个月,成本高昂
- 流引擎重放的吞吐要足够高,否则追平历史要跑很久
- 流处理的结果要能覆盖式写入,支持重算后替换
流批一体是另一条路线,它不砍掉批,而是让批和流用同一套代码、同一份表。Flink 的 SQL 层在这个方向上做了工作(同一段 SQL 既能以批模式跑也能以流模式跑),Iceberg 这类表格式则提供了"同一张表既能被流写入也能被批读取"的基础。
【推断】 流批一体在代码层面(同一段 SQL)已经比较可用,但在口径层面并不自动一致——因为批模式看到的是完整数据,流模式看到的是带 watermark 的近似完整数据,迟到数据的处理策略不同就会产生差异。所以"用了流批一体就不会对不上数"是不成立的。
【推断】 根据公开资料所反映的通行做法,多数团队的实际形态既不是纯 Lambda 也不是纯 Kappa,而是按指标重要性分级:
- 大促实时大屏、风控这类:走实时链,接受一定近似
- 财务结算、对外披露数据:只认离线链的结果
- 大多数普通报表:只有离线链,不做实时
也就是说,不是所有指标都值得付两条链的成本。哪些指标上实时,本质是一道成本-价值题
| 角色 | 主要交付物 | 关心什么 |
|---|---|---|
| 客户端 / 服务端工程师 | 埋点代码、业务库表结构 | 功能上线、包体积、接口性能 |
| 数据平台工程师 | Kafka / 计算引擎 / 调度系统 / 湖仓的可用性 | 集群稳定性、资源利用率、成本 |
| 数仓工程师(也叫数据开发 / ETL 工程师) | ODS→ADS 的模型与任务、指标口径 | 数据准确性、任务准时产出、模型可复用 |
| 数据治理 / 数据产品 | 指标体系、数据字典、质量规则、资产目录 | 口径统一、可发现性、合规 |
| 数据分析师 / 算法工程师 | 报表、分析结论、特征与模型 | 数据可用性、口径可解释、特征时效 |
两个高发事故接缝:
接缝一:客户端 ↔ 数据团队(埋点)。 前面讲过 KPI 错配。典型事故:客户端为了复用代码,把 A 页面的曝光事件也用在了 B 页面,没通知数据团队;下游按"这个事件只来自 A 页面"写的过滤逻辑瞬间失效,且不报错,只是数变大了——这类故障最难发现。
接缝二:数据团队 ↔ 算法团队(特征)。 训练模型时用的是离线数仓算的特征,线上推理时用的是实时链算的特征,两边口径不一致。这个问题有专门的名字叫 training-serving skew(训练-服务偏差),后文专门处理。它的危害是模型线上效果无声下降,同样不报错。