数据仓库搭建:从数据孤岛到单一事实来源

____simple_html_dom__voku__html_wrapper____>

数据仓库搭建:从数据孤岛到单一事实来源

cover

一、当"数据驱动"变成对账游戏

很多公司都在提"数据驱动决策",但实际场景往往是这样的:业务方问"上个月华东区高净值用户的复购率是多少",分析师得从 MySQL 拉订单、Kafka 抓行为日志、Oracle 取财务数据,再跑到飞书多维表格里找运营数据。最后拼出来的数字,跟财务那边算出来的永远对不上。

这就是典型的数据孤岛。数据仓库要解决的,就是把散在各处的数据收拢、清洗、建模,最终形成一个"单一事实来源"(Single Source of Truth)。下面我结合实战经验,拆解一个中等规模数仓的搭建过程。

二、四层架构:ODS-DWD-DWS-ADS 的流转逻辑

现代数仓普遍采用四层分层架构,数据从原始层到应用层,逐步从"脏乱散"走向"净整聚"。

graph TB
    subgraph 数据源
        S1[MySQL 业务库]
        S2[Kafka 行为日志]
        S3[API 外部数据]
        S4[CSV 手工填报]
    end
    subgraph 数仓分层
        ODS[ODS 贴源层<br/>原始数据 1:1 映射]
        DWD[DWD 明细层<br/>清洗 + 标准化]
        DWS[DWS 汇总层<br/>主题聚合]
        ADS[ADS 应用层<br/>指标输出]
    end
    S1 -->|全量/增量同步| ODS
    S2 -->|实时消费| ODS
    S3 -->|定时拉取| ODS
    S4 -->|手动上传| ODS
    ODS -->|数据清洗| DWD
    DWD -->|维度建模| DWS
    DWS -->|指标计算| ADS
    ADS --> R1[BI 看板]
    ADS --> R2[数据产品]
    ADS --> R3[算法特征]

ODS 层(贴源层):跟源系统结构完全一致,不做业务逻辑转换,只做技术层面的格式适配(比如时区统一、编码转换)。这层相当于"数据保险箱",保留最原始的数据快照,下游有问题可以回溯到 ODS 层重新计算。

DWD 层(明细层):在 ODS 基础上做清洗和标准化。清洗包括去重(同一订单被同步两次)、缺失值处理(用户 ID 为空的记录)、异常值过滤(金额为负数的订单)。标准化包括统一日期格式、统一枚举值编码、统一字段命名。

DWS 层(汇总层):按业务主题(如用户、商品、订单)进行轻度汇总,生成宽表。比如把用户基本信息、订单统计、行为统计合并成一张"用户主题宽表",避免下游每次都做多表关联。

ADS 层(应用层):面向具体业务场景的指标输出,比如"各渠道月度转化率"、"用户流失预警名单"。这层直接对接 BI 工具和数据产品。

三、生产级实现:基于 dbt + PostgreSQL 的数仓搭建

下面用 dbt(Data Build Tool)作为数仓建模引擎,PostgreSQL 作为存储后端,实现完整的四层数仓架构。

-- ============================================================
-- ODS 层:贴源层,与业务库结构 1:1 映射
-- ============================================================
-- ODS 订单表:保留源系统全部字段,仅做类型标准化
CREATE TABLE IF NOT EXISTS ods_orders (
    id              BIGINT          NOT NULL,
    user_id         BIGINT          NOT NULL,
    product_id      BIGINT,
    channel         VARCHAR(50),
    amount          NUMERIC(12, 2),
    order_status    VARCHAR(20),
    created_at      TIMESTAMP,
    updated_at      TIMESTAMP,
    -- 技术字段:数据同步元信息
    _sync_time      TIMESTAMP       DEFAULT CURRENT_TIMESTAMP,
    _source_system  VARCHAR(20)     DEFAULT 'mysql_order',
    PRIMARY KEY (id, _sync_time)
);
-- ODS 用户行为表:从 Kafka 消费写入,保留原始事件结构
CREATE TABLE IF NOT EXISTS ods_user_events (
    event_id        VARCHAR(64)     NOT NULL,
    user_id         BIGINT,
    event_type      VARCHAR(50),
    event_time      TIMESTAMP,
    page_url        VARCHAR(500),
    device_type     VARCHAR(20),
    -- 技术字段
    _sync_time      TIMESTAMP       DEFAULT CURRENT_TIMESTAMP,
    _source_system  VARCHAR(20)     DEFAULT 'kafka_event',
    PRIMARY KEY (event_id)
);
-- ============================================================
-- DWD 层:明细层,清洗 + 标准化
-- ============================================================
-- DWD 订单明细表:过滤无效数据,统一枚举值
CREATE TABLE IF NOT EXISTS dwd_order_detail AS
SELECT
    o.id                                AS order_id,
    o.user_id,
    o.product_id,
    -- 渠道枚举标准化:将源系统的多种写法统一
    CASE
        WHEN LOWER(TRIM(o.channel)) IN ('app', 'mobile_app', '手机端') THEN 'APP'
        WHEN LOWER(TRIM(o.channel)) IN ('web', 'pc', '电脑端') THEN 'WEB'
        WHEN LOWER(TRIM(o.channel)) IN ('mini', 'miniprogram', '小程序') THEN 'MINI'
        ELSE 'OTHER'
    END                                 AS channel,
    -- 金额清洗:负数和极端值标记为异常
    CASE
        WHEN o.amount < 0 THEN NULL
        WHEN o.amount > 100000 THEN NULL  -- 单笔超 10 万视为异常,需人工复核
        ELSE o.amount
    END                                 AS amount,
    -- 订单状态标准化
    UPPER(TRIM(o.order_status))         AS order_status,
    o.created_at                        AS order_time,
    o.updated_at                        AS update_time,
    o._sync_time                        AS etl_time
FROM ods_orders o
WHERE
    -- 过滤条件:排除测试数据和无效记录
    o.user_id IS NOT NULL
    AND o.created_at >= '2024-01-01'  -- 只保留有效时间范围的数据
    AND o.id NOT IN (SELECT id FROM ods_orders GROUP BY id HAVING COUNT(*) > 1)
    -- 去重:同一订单取最新同步时间
    AND o._sync_time = (
        SELECT MAX(_sync_time) FROM ods_orders WHERE id = o.id
    );
-- ============================================================
-- DWS 层:汇总层,按用户主题聚合
-- ============================================================
-- DWS 用户主题宽表:将订单、行为、属性合并为一张宽表
CREATE TABLE IF NOT EXISTS dws_user_profile AS
SELECT
    u.user_id,
    -- 基本属性
    u.city,
    u.vip_level,
    -- 订单统计
    COALESCE(orders.total_orders, 0)        AS total_orders,
    COALESCE(orders.total_amount, 0)        AS total_amount,
    COALESCE(orders.avg_amount, 0)          AS avg_order_amount,
    COALESCE(orders.last_order_time, NULL)  AS last_order_time,
    -- 行为统计
    COALESCE(events.total_events, 0)        AS total_events,
    COALESCE(events.view_count, 0)          AS view_count,
    COALESCE(events.click_count, 0)         AS click_count,
    -- 衍生指标:最近一次购买距今天数
    CURRENT_DATE - orders.last_order_time::DATE AS recency_days
FROM dim_users u
-- 左连接订单统计,保留无订单用户
LEFT JOIN (
    SELECT
        user_id,
        COUNT(*) AS total_orders,
        SUM(amount) AS total_amount,
        AVG(amount) AS avg_amount,
        MAX(order_time) AS last_order_time
    FROM dwd_order_detail
    WHERE amount IS NOT NULL
    GROUP BY user_id
) orders ON u.user_id = orders.user_id
-- 左连接行为统计
LEFT JOIN (
    SELECT
        user_id,
        COUNT(*) AS total_events,
        COUNT(*) FILTER (WHERE event_type = 'PAGE_VIEW') AS view_count,
        COUNT(*) FILTER (WHERE event_type = 'CLICK') AS click_count
    FROM dwd_user_event
    GROUP BY user_id
) events ON u.user_id = events.user_id;
-- ============================================================
-- ADS 层:应用层,面向业务场景的指标输出
-- ============================================================
-- ADS 渠道转化率报表:按月按渠道统计
CREATE TABLE IF NOT EXISTS ads_channel_conversion AS
SELECT
    DATE_TRUNC('month', order_time)::DATE  AS month_date,
    channel,
    COUNT(DISTINCT user_id)                AS order_users,
    SUM(amount)                            AS total_amount,
    AVG(amount)                            AS avg_amount,
    -- 转化率需要与行为数据关联计算
    ROUND(
        COUNT(DISTINCT user_id)::NUMERIC / NULLIF(
            (SELECT COUNT(DISTINCT user_id)
             FROM dwd_user_event e
             WHERE DATE_TRUNC('month', e.event_time)::DATE = DATE_TRUNC('month', order_time)::DATE
               AND e.channel = dwd_order_detail.channel),
        ) * 100, 2
    ) AS conversion_rate_pct
FROM dwd_order_detail
WHERE amount IS NOT NULL
GROUP BY 1, 2
ORDER BY 1, 2;

配套的数据质量校验脚本,用于在每层 ETL 完成后自动检测数据异常:

import logging
from dataclasses import dataclass, field
from datetime import datetime
import pandas as pd
from sqlalchemy import create_engine, text
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("data_quality")
@dataclass
class QualityRule:
    """数据质量规则定义:每条规则对应一个检查项"""
    table_name: str
    rule_name: str
    rule_sql: str
    expected_result: str  # "empty" 表示查询结果应为空,"non_empty" 表示应有结果
    severity: str = "warning"  # warning / critical
class DataQualityChecker:
    """数据质量校验器:在 ETL 完成后自动执行规则检查"""
    def __init__(self, db_url: str):
        self.engine = create_engine(db_url, pool_pre_ping=True)
        self.rules: list[QualityRule] = []
        self.results: list[dict] = []
    def add_rule(self, rule: QualityRule) -> None:
        self.rules.append(rule)
    def add_default_rules(self) -> None:
        """添加通用数据质量规则,覆盖空值、重复、一致性三类检查"""
        # 规则1:DWD 层订单金额不应有空值(空值说明清洗逻辑有遗漏)
        self.add_rule(QualityRule(
            table_name="dwd_order_detail",
            rule_name="金额空值检查",
            rule_sql="SELECT COUNT(*) FROM dwd_order_detail WHERE amount IS NULL",
            expected_result="zero",  # 期望结果为 0
            severity="critical",
        ))
        # 规则2:DWD 层订单不应有重复 order_id
        self.add_rule(QualityRule(
            table_name="dwd_order_detail",
            rule_name="订单去重检查",
            rule_sql=(
                "SELECT order_id, COUNT(*) as cnt FROM dwd_order_detail "
                "GROUP BY order_id HAVING COUNT(*) > 1 LIMIT 10"
            ),
            expected_result="empty",
            severity="critical",
        ))
        # 规则3:DWS 层用户数应与维度表一致(不应多也不应少)
        self.add_rule(QualityRule(
            table_name="dws_user_profile",
            rule_name="用户数一致性检查",
            rule_sql=(
                "SELECT ABS("
                "(SELECT COUNT(*) FROM dws_user_profile) - "
                "(SELECT COUNT(*) FROM dim_users)"
                ") AS diff"
            ),
            expected_result="zero",
            severity="warning",
        ))
        # 规则4:ADS 层转化率应在合理范围内(0-100%)
        self.add_rule(QualityRule(
            table_name="ads_channel_conversion",
            rule_name="转化率范围检查",
            rule_sql=(
                "SELECT * FROM ads_channel_conversion "
                "WHERE conversion_rate_pct < 0 OR conversion_rate_pct > 100"
            ),
            expected_result="empty",
            severity="warning",
        ))
    def run_checks(self) -> list[dict]:
        """执行所有质量规则,返回检查结果"""
        for rule in self.rules:
            try:
                with self.engine.connect() as conn:
                    result = conn.execute(text(rule.rule_sql)).fetchall()
                    # 判断是否通过
                    if rule.expected_result == "empty":
                        passed = len(result) == 0
                    elif rule.expected_result == "zero":
                        passed = result[0][0] == 0
                    else:
                        passed = True
                    check_result = {
                        "table": rule.table_name,
                        "rule": rule.rule_name,
                        "passed": passed,
                        "severity": rule.severity,
                        "checked_at": datetime.now().isoformat(),
                        "detail": str(result[:5]) if not passed else None,
                    }
                    self.results.append(check_result)
                    status = "PASS" if passed else "FAIL"
                    logger.info(f"[{status}] {rule.table_name}.{rule.rule_name}")
            except Exception as e:
                self.results.append({
                    "table": rule.table_name,
                    "rule": rule.rule_name,
                    "passed": False,
                    "severity": "critical",
                    "error": str(e),
                    "checked_at": datetime.now().isoformat(),
                })
                logger.error(f"[ERROR] {rule.table_name}.{rule.rule_name}: {e}")
        return self.results
    def has_critical_failure(self) -> bool:
        """是否存在 critical 级别的检查失败"""
        return any(
            not r["passed"] and r["severity"] == "critical"
            for r in self.results
        )
# 使用示例
if __name__ == "__main__":
    checker = DataQualityChecker("postgresql://user:pass@localhost:5432/warehouse")
    checker.add_default_rules()
    results = checker.run_checks()
    if checker.has_critical_failure():
        logger.error("存在 critical 级别质量检查失败,阻断下游任务")
        # 在 Airflow 中此处应 raise 触发任务失败
    else:
        logger.info("所有 critical 检查通过,可继续下游任务")

几个关键设计点:ODS 层保留 _sync_time_source_system 技术字段,用于数据溯源和增量同步。DWD 层的渠道枚举标准化用 CASE WHEN 把源系统的多种写法统一,这是数据治理里最常见的"脏数据"场景。DataQualityChecker 在 ETL 完成后自动执行规则检查,critical 级别失败时阻断下游任务,避免脏数据污染整个数仓。

四、数仓建设的隐性成本:不只是技术问题

搭建数据仓库的技术方案并不复杂,真正的挑战在于组织层面的隐性成本:

第一,数据治理的持续性投入。 数据质量不会因为建了数仓就自动变好。源系统的变更(比如新增字段、修改枚举值)会持续影响数仓的稳定性。一个中等规模数仓,每月大概有 5-10 次源系统变更需要跟进。如果没有专职的数据治理角色,数仓会在 3-6 个月内退化到"数据又对不上了"的状态。

第二,维度建模的认知负担。 Kimball 维度建模理论看似简单,但在实际业务中,维度和事实的边界往往模糊。比如"渠道"在订单事实表中是维度,但在渠道分析场景里又变成了事实的主体。维度建模的过度规范化会导致表数量爆炸(一个中等业务可能产生 50+ 张 DWS 宽表),增加维护成本。

第三,实时与离线的架构割裂。 本文展示的是离线数仓架构(T+1 数据),但业务方越来越需要实时数据。引入实时层(比如 Flink + Kafka)后,离线和实时两套架构的指标口径如何保持一致,是个没有银弹的工程难题。当前主流方案是"Lambda 架构"或"Kappa 架构",但两者都有各自的复杂度代价。

适用边界总结:

  • 适合:数据源 >= 3 个、跨系统分析需求频繁、有专职数据团队的场景
  • 不适合:数据源单一、分析需求简单、团队规模 < 2 人的场景(直接用 BI 工具直连业务库更高效)

五、总结

本文从数据孤岛问题出发,拆解了 ODS-DWD-DWS-ADS 四层数仓架构的设计逻辑,并给出了基于 PostgreSQL + dbt 的生产级实现,包含数据质量校验的自动化方案。

落地路线建议:

  1. 起步阶段:先搭建 ODS + DWD 两层,解决数据汇聚和清洗问题,验证数据质量
  2. 进阶阶段:引入 DWS 汇总层,按核心业务主题构建宽表,降低下游查询复杂度
  3. 成熟阶段:建设 ADS 应用层和数据质量监控体系,实现指标口径统一和异常自动告警

数据仓库不是一次性工程,而是持续演进的系统。它的价值不在于架构多完美,而在于数据是否真的变成了"单一事实来源"。

© 版权声明

相关文章