全文共 2,641 字 预计阅读 8 分钟
bg

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 解决了什么:同时拿到低延迟和最终准确性。实时链先给一个近似值,离线链在次日用全量数据把它覆盖修正。

代价(这是重点)

  1. 同一份业务逻辑写两遍。批和流的算子语义不同(比如"去重"在批里是 DISTINCT,在流里是有状态的 keyed state 加超时清理),不是复制粘贴能解决的。
  2. 口径漂移不可避免。两套代码只要有一处理解偏差,数就对不上,而排查"实时和离线差 0.3% 是谁的问题"极其耗时。
  3. 两套资源、两套运维
  4. 修正逻辑本身很复杂。服务层怎么判断某个时间段该用批的结果还是流的结果,是一个真实的工程难题。

Kappa 架构的主张很直接:砍掉批处理层,只保留一条流链路。需要重算历史时,就把消息队列里的数据从头重放一遍。

它成立的前提(这三条决定了 Kappa 能不能用):

  1. 消息队列要能保存足够长时间的历史数据——重放 3 个月,Kafka 就得存 3 个月,成本高昂
  2. 流引擎重放的吞吐要足够高,否则追平历史要跑很久
  3. 流处理的结果要能覆盖式写入,支持重算后替换

流批一体是另一条路线,它不砍掉批,而是让批和流用同一套代码、同一份表。Flink 的 SQL 层在这个方向上做了工作(同一段 SQL 既能以批模式跑也能以流模式跑),Iceberg 这类表格式则提供了"同一张表既能被流写入也能被批读取"的基础。

【推断】 流批一体在代码层面(同一段 SQL)已经比较可用,但在口径层面并不自动一致——因为批模式看到的是完整数据,流模式看到的是带 watermark 的近似完整数据,迟到数据的处理策略不同就会产生差异。所以"用了流批一体就不会对不上数"是不成立的。

【推断】 根据公开资料所反映的通行做法,多数团队的实际形态既不是纯 Lambda 也不是纯 Kappa,而是按指标重要性分级

  • 大促实时大屏、风控这类:走实时链,接受一定近似
  • 财务结算、对外披露数据:只认离线链的结果
  • 大多数普通报表:只有离线链,不做实时

也就是说,不是所有指标都值得付两条链的成本。哪些指标上实时,本质是一道成本-价值题

角色 主要交付物 关心什么
客户端 / 服务端工程师 埋点代码、业务库表结构 功能上线、包体积、接口性能
数据平台工程师 Kafka / 计算引擎 / 调度系统 / 湖仓的可用性 集群稳定性、资源利用率、成本
数仓工程师(也叫数据开发 / ETL 工程师) ODS→ADS 的模型与任务、指标口径 数据准确性、任务准时产出、模型可复用
数据治理 / 数据产品 指标体系、数据字典、质量规则、资产目录 口径统一、可发现性、合规
数据分析师 / 算法工程师 报表、分析结论、特征与模型 数据可用性、口径可解释、特征时效

两个高发事故接缝

接缝一:客户端 ↔ 数据团队(埋点)。 前面讲过 KPI 错配。典型事故:客户端为了复用代码,把 A 页面的曝光事件也用在了 B 页面,没通知数据团队;下游按"这个事件只来自 A 页面"写的过滤逻辑瞬间失效,且不报错,只是数变大了——这类故障最难发现。

接缝二:数据团队 ↔ 算法团队(特征)。 训练模型时用的是离线数仓算的特征,线上推理时用的是实时链算的特征,两边口径不一致。这个问题有专门的名字叫 training-serving skew(训练-服务偏差),后文专门处理。它的危害是模型线上效果无声下降,同样不报错。

Back to Blog