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

一、当"数据驱动"变成对账游戏
很多公司都在提"数据驱动决策",但实际场景往往是这样的:业务方问"上个月华东区高净值用户的复购率是多少",分析师得从 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 的生产级实现,包含数据质量校验的自动化方案。
落地路线建议:
- 起步阶段:先搭建 ODS + DWD 两层,解决数据汇聚和清洗问题,验证数据质量
- 进阶阶段:引入 DWS 汇总层,按核心业务主题构建宽表,降低下游查询复杂度
- 成熟阶段:建设 ADS 应用层和数据质量监控体系,实现指标口径统一和异常自动告警
数据仓库不是一次性工程,而是持续演进的系统。它的价值不在于架构多完美,而在于数据是否真的变成了"单一事实来源"。