Towards Data Science

The Medallion Data Architecture: An Introduction

8.5内容质量

TL;DR · AI 摘要

Medallion数据架构通过Bronze、Silver、Gold三层分层处理数据,提升数据质量和可维护性。

核心要点

  • Bronze层存储原始数据,需记录数据源、加载时间等元数据
  • Silver层通过数据清洗和规则校验提升数据质量
  • Gold层提供业务规则驱动的聚合视图,用于报表和分析

结构提纲

按章节快速跳转。

  1. 数据ETL流水线随着规模扩大变得难以信任和维护。

  2. Databricks提出青铜-白银-黄金三层架构术语并推广。

  3. 青铜层存储原始数据,白银层进行清洗,黄金层生成聚合视图。

  4. 数据不可变且仅追加,需记录完整元数据。

  5. 通过规则引擎清洗无效数据并建立数据一致性。

  6. 基于业务规则生成物化视图,支持报表和分析系统。

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • Medallion架构
    • 青铜层
      • 原始数据存储
      • 元数据记录
    • 白银层
      • 数据清洗
      • 规则校验
    • 黄金层
      • 聚合视图
      • 业务报表

金句 / Highlights

值得收藏与分享的关键句。

#数据工程#数据架构#Databricks#Delta Lake
打开原文

勋章数据架构:简介 | Towards Data Science

数据工程

勋章数据架构:简介

关于青铜层、白银层和黄金层的实用指南,附带可运行的 Python 和 DuckDB 示例

Thomas Reid

2026 年 8 月 4 日

14 分钟阅读

分享

AI 生成图片

随着数据 ETL 流水线范围的扩大,其可信度往往会降低,而且在无错误运行、文档记录和调试方面也会变得更加困难。

一个系统会传来 CSV 文件,另一个系统会传来 JSON 数据,还有其他地方会传来 Parquet 文件。数周甚至数月过去后,人们突然发现没人能确定哪些数据版本是可信的,哪些规则被应用过,或者为什么昨天的仪表板会出现错误。

勋章架构是对这一问题的实用回应。它将数据平台划分为三层,通常称为青铜层、白银层和黄金层。在每一层的边界处,应有对该层所包含数据的清晰、文档化的描述。这一点在青铜层尤为关键,因为这是数据首次摄入的地方,因此你需要尽可能详细地记录数据来源、加载该数据的人员或系统、加载时间、加载频率等信息。

在理想情况下,每一层的数据应通过 SQL、Python、dbt 等工具进行处理。

勋章架构的起源

青铜层、白银层和黄金层的术语最早由 Databricks 提出。Databricks 是一家数据和人工智能公司,其云平台帮助组织使用 Apache SparkDelta Lake 等技术处理、管理和分析大规模数据集。

Databricks 将勋章架构描述为一种多层模式,数据在通过三层时质量逐步提升。

通常,青铜层用于存储从源系统接收到的原始、未经筛选的数据。记录通常是不可变的,并且只追加。

白银层包含经过清理的青铜层数据。例如,空记录、无效日期、缺失字段等会在存储到此层前被修复或删除。

黄金层通常包含根据业务规则定义的专用聚合数据集,这些数据集以 SQL(物化)视图的形式从白银层派生。例如,数据仪表板和管理报告通常由黄金层的数据构建,因为这些数据是准确的,体积更小,并能带来更高的准确性和更低的处理时间。

当然,像这样的系统在数据出现时就已存在。大多数数据库工程师在听到“勋章”一词之前,早就使用过“暂存区”来将数据引入系统,然后再将其分发到需要的地方。这是一种简单的双层勋章系统。Databricks 只是添加了另一层,为其赋予了花哨的名称并推广了它。

每一层应包含什么?

让我们更详细地看看每一层理想情况下应包含的内容。请注意,在实际系统中,黄金层、白银层和青铜层通常对应现代数据库或数据仓库中的不同模式。

青铜层

青铜层记录了从源系统接收到的内容。有用的青铜层数据还可能包括与源字段(如)一起的摄入元数据:

  • 数据源系统及源文件或事件标识符
  • 摄入时间戳和/或业务生效日期
  • 摄入的记录数量
  • 批次或加载运行标识符

在这一层处理错误和其他类型的数据问题的方式非常重要。

对于错误或缺失的数据值,应保持原样,并根据需要在银层进行隔离。如果数据加载因网络故障等原因中途失败,应标记为失败或被替代,并作为新批次重新加载。

如果收到额外或延迟的数据,请将其作为另一个批次追加,并记录其来源、摄入时间和业务生效日期。如果同一份交付被提交两次,使用文件哈希、批次标识符或来源键来防止意外重复。

无论采取何种方法,摄入过程都应具备幂等性。多次处理同一来源交付不应产生重复记录或以其他方式改变最终状态。

银层

银层会对青铜层数据集应用规则,使记录足够可靠和准确以供使用。典型的转换工作包括:

  • 解析并强制数据类型
  • 标准化日期、货币、国家代码和单位
  • 去重记录
  • 隔离劣质数据
  • 连接参考数据

银层通常应保留业务级别的详细信息。它是产品和下游系统可以依赖的干净、基础数据。

在此层级上出错可能会严重破坏下游系统和流程。例如,银层的order_total列应定义货币和数字类型。order_id应有明确的唯一性规则。如果某行未通过这些规则,管道需要明确的处理结果(例如插入隔离表),而不是静默忽略。

金层

金层围绕特定的业务用例和流程进行组织。金层通常包括:

  • 按天、月或地区汇总和聚合的数据集(例如总销售额、活跃用户数)
  • 为快速查询构建的星型模式或数据集市(连接较少)
  • 为特定团队(如财务、市场或运营)定制的独立数据集

将所有内容串联在一起的是一张图,展示了典型且非常简单的勋章系统可能的结构。

实施勋章模式需要哪些工具?

没有统一的实现方式,但作为起点,我通常会使用某种数据库、数据仓库或基于云的对象存储来实现勋章架构,其中金层、银层和青铜层通常在数据库中是不同的模式,或在对象存储中是不同的文件夹。这可以在从本地笔记本电脑上的SQLite到基于云的大型集群上的AWS Redshift数据湖,或AWS S3/Azure Blob/Google Cloud Storage等任何环境中运行。

对于基于云的对象存储,还需要考虑要使用的开放表格式。最常用的三种是Hudi、Apache Iceberg和Delta表。

就软件工具而言,我认为勋章模式只是数据工程(DE)工作的一部分。因此,数据工程师日常使用的工具与设置和维护勋章系统所用的工具相同。SQL将是你的主要工具,同时要记住像dbt这样的其他工具也依赖于SQL。除了SQL外,Python、Spark和其他编程语言通常也会被使用。

对于基于云的架构,你还可以使用特定于该平台的工具。我主要使用 AWS,因此可能会使用 AWS Athena 进行数据查询,使用 AWS Glue 进行管道开发工作,并使用 Step 进行编排。

注意,除了是本文提到的各种系统和产品(例如 DuckDB)的用户外,我与它们没有任何关联或商业关系。

实际案例:使用 Python 和 DuckDB 处理零售订单

在这个示例中,我使用的是小型在线零售商的夜间 CSV 导出文件。该文件在用于报告前需要进行一些处理。订单可能会重复,某些日期无法解析,负金额需要被拒绝。管道在夜间运行,以便运营团队能在次日 07:00 前获得按地区和货币划分的已支付和退款销售总额。

该管道包含五个阶段:

  • 将每个 CSV 导入内容原样存储到仅追加的 Bronze 表中。
  • 在 Silver 表中将字段转换为正确类型,验证值并移除重复订单。
  • 将被拒绝的行移动到隔离表中进行调查。
  • 将已接受的订单聚合到 Gold 表中的每日区域销售数据。
  • 防止相同源文件被重复摄入。

使用 DuckDB 作为数据库使示例保持简洁,但如果需要扩展到更大的湖仓架构,各层之间的契约可以直接映射。

我们的项目结构将类似于以下布局。

code
retail-medallion/
├── data/
│   └── incoming/
│       └── orders_2026-07-19.csv    <= 由你手动创建
├── pipeline.py                      <= 由你手动创建
└── warehouse.duckdb                 <= 该数据库文件由管道创建

创建虚拟环境并安装 DuckDB

code
D:\projects\retail-medallion> python3 -m venv .venv
# Windows PowerShell: .\.venv\Scripts\Activate.ps1
# macOS/Linux: source .venv/bin/activate
D:\projects\retail-medallion> python3 -m pip install duckdb pytz tabulate

创建输入文件

这只是一个简单的 CSV 文件,打开你最喜欢的文本编辑器并输入以下数据。将其保存为 data/incoming 文件夹下的 orders_2026-07-19.csv 文件。

code
order_id,ordered_at,customer_id,region,amount,currency,status
1001,2026-07-19T09:10:00Z,C001,North,125.50,GBP,paid
1002,2026-07-19T10:05:00Z,C002,South,89.99,GBP,paid
1002,2026-07-19T10:05:00Z,C002,South,89.99,GBP,paid
1003,not-a-date,C003,North,45.00,GBP,paid
1004,2026-07-19T11:42:00Z,C004,West,-10.00,GBP,paid
1005,2026-07-19T12:20:00Z,C005,North,210.00,GBP,refunded

重复和无效行是故意设置的,这对验证管道在数据质量差时的表现是一个很好的测试。

我们的管道代码

将以下代码保存到项目根目录下的 pipeline.py 文件中。

code
from __future__ import annotations

import hashlib
import sys
from pathlib import Path

import duckdb

DATABASE = Path("warehouse.duckdb")

def file_hash(path: Path) -> str:
    digest = hashlib.sha256()
    with path.open("rb") as source:
        for block in iter(lambda: source.read(1024 * 1024), b""):
            digest.update(block)
    return digest.hexdigest()
python
def initialise(connection: duckdb.DuckDBPyConnection) -> None:
    connection.execute("CREATE SCHEMA IF NOT EXISTS bronze")
    connection.execute("CREATE SCHEMA IF NOT EXISTS silver")
    connection.execute("CREATE SCHEMA IF NOT EXISTS gold")
    connection.execute("""
        CREATE TABLE IF NOT EXISTS bronze.ingestion_batches (
            source_hash VARCHAR PRIMARY KEY,
            source_file VARCHAR NOT NULL,
            ingested_at TIMESTAMPTZ NOT NULL DEFAULT current_timestamp
        )
    """)
    connection.execute("""
        CREATE TABLE IF NOT EXISTS bronze.orders_raw (
            order_id VARCHAR,
            ordered_at VARCHAR,
            customer_id VARCHAR,
            region VARCHAR,
            amount VARCHAR,
            currency VARCHAR,
            status VARCHAR,
            source_file VARCHAR NOT NULL,
            source_hash VARCHAR NOT NULL,
            ingested_at TIMESTAMPTZ NOT NULL
        )
    """)

def ingest_bronze(connection: duckdb.DuckDBPyConnection, source: Path) -> bool:
    source = source.resolve()
    digest = file_hash(source)
    already_loaded = connection.execute(
        "SELECT 1 FROM bronze.ingestion_batches WHERE source_hash = ?", [digest]
    ).fetchone()
    if already_loaded:
        print(f"Skipping {source.name}: this exact file has already been loaded")
        return False

    connection.begin()
    try:
        connection.execute(
            """
            INSERT INTO bronze.orders_raw
            SELECT
                order_id, ordered_at, customer_id, region, amount,
                currency, status, ?, ?, current_timestamp
            FROM read_csv(?, header = true, all_varchar = true)
            """,
            [source.name, digest, str(source)],
        )
        connection.execute(
            """INSERT INTO bronze.ingestion_batches (source_hash, source_file)
            VALUES (?, ?)""",
            [digest, source.name],
        )
        connection.commit()
    except Exception:
        connection.rollback()
        raise
    print(f"Loaded {source.name} into bronze")
    return True
python
def build_silver(connection: duckdb.DuckDBPyConnection) -> None:
    connection.execute("""
        CREATE OR REPLACE TEMP VIEW typed_orders AS
        SELECT
            trim(order_id) AS order_id,
            try_cast(ordered_at AS TIMESTAMPTZ) AS ordered_at,
            trim(customer_id) AS customer_id,
            upper(trim(region)) AS region,
            try_cast(amount AS DECIMAL(18, 2)) AS amount,
            upper(trim(currency)) AS currency,
            lower(trim(status)) AS status,
            source_file,
            source_hash,
            ingested_at,
            row_number() OVER (
                PARTITION BY trim(order_id)
                ORDER BY ingested_at DESC, source_file DESC
            ) AS duplicate_rank
        FROM bronze.orders_raw
    """)
    valid = """
        order_id IS NOT NULL AND order_id <> ''
        AND ordered_at IS NOT NULL
        AND customer_id IS NOT NULL AND customer_id <> ''
        AND amount IS NOT NULL AND amount >= 0
        AND currency IN ('GBP', 'EUR', 'USD')
        AND status IN ('paid', 'refunded', 'cancelled')
        AND duplicate_rank = 1
    """
    connection.execute(f"""
        CREATE OR REPLACE TABLE silver.orders AS
        SELECT * EXCLUDE (duplicate_rank)
        FROM typed_orders
        WHERE {valid}
    """)
    connection.execute(f"""
        CREATE OR REPLACE TABLE silver.orders_quarantine AS
        SELECT
            * EXCLUDE (duplicate_rank),
            CASE
                WHEN duplicate_rank > 1 THEN 'duplicate order_id'
                WHEN ordered_at IS NULL THEN 'invalid ordered_at'
                WHEN amount IS NULL THEN 'invalid amount'
                WHEN amount < 0 THEN 'negative amount'
                WHEN currency NOT IN ('GBP', 'EUR', 'USD') THEN 'unsupported currency'
                WHEN status NOT IN ('paid', 'refunded', 'cancelled') THEN 'invalid status'
                ELSE 'missing required value'
            END AS rejection_reason
        FROM typed_orders
        WHERE NOT ({valid})
    """)

def build_gold(connection: duckdb.DuckDBPyConnection) -> None:
    connection.execute("""
        CREATE OR REPLACE TABLE gold.daily_sales_by_region AS
        SELECT
            cast(ordered_at AS DATE) AS order_date,
            region,
            currency,
            count(*) FILTER (WHERE status = 'paid') AS paid_orders,
            sum(amount) FILTER (WHERE status = 'paid') AS gross_sales,
            count(*) FILTER (WHERE status = 'refunded') AS refunded_orders,
            sum(amount) FILTER (WHERE status = 'refunded') AS refunded_value
        FROM silver.orders
        GROUP BY order_date, region, currency
        ORDER BY order_date, region, currency
    """)

def check_quality(connection: duckdb.DuckDBPyConnection) -> None:
    duplicate_count = connection.execute(
        "SELECT count(*) - count(DISTINCT order_id) FROM silver.orders"
    ).fetchone()[0]
    null_key_count = connection.execute(
        "SELECT count(*) FROM silver.orders WHERE order_id IS NULL"
    ).fetchone()[0]
    if duplicate_count or null_key_count:
        raise RuntimeError("Silver quality contract failed")

def print_query(connection: duckdb.DuckDBPyConnection, query: str) -> None:
    result = connection.execute(query)
    print(" | ".join(column[0] for column in result.description))
    for row in result.fetchall():
        print(" | ".join("NULL" if value is None else str(value) for value in row))
python
def main(source: Path) -> None:
    with duckdb.connect(str(DATABASE)) as connection:
        initialise(connection)
        ingest_bronze(connection, source)
        build_silver(connection)
        check_quality(connection)
        build_gold(connection)
        print("\nGold output")
        print_query(connection, "SELECT * FROM gold.daily_sales_by_region")
        print("\nQuarantined records")
        print_query(
            connection,
            """SELECT order_id, ordered_at, amount, rejection_reason
            FROM silver.orders_quarantine""",
        )

if __name__ == "__main__":
    if len(sys.argv) != 2:
        raise SystemExit("Usage: python pipeline.py path/to/orders.csv")
    main(Path(sys.argv[1]))
code
使用以下命令运行

python3 pipeline.py data/incoming/orders_2026-07-19.csv

输出结果?

code
已将 orders_2026-07-19.csv 加载到青铜层

Gold output
order_date | region | currency | paid_orders | gross_sales | refunded_orders | refunded_value
2026-07-19 | NORTH  | GBP      | 1           | 125.50      | 1               | 210.00
2026-07-19 | SOUTH  | GBP      | 1           | 89.99       | 0               | NULL

Quarantined records
order_id | ordered_at                | amount | rejection_reason
1004     | 2026-07-19 12:42:00+01:00 | -10.00 | negative amount
1003     | NULL                      | 45.00  | invalid ordered_at
1002     | 2026-07-19 11:05:00+01:00 | 89.99  | duplicate order_id

运行完成后,黄金层将为每个日期、地区和货币组合生成一行数据,分别显示已支付订单和退款订单的统计数据。包含错误日期、负金额或重复订单ID的记录将无法进入黄金层。这些记录会被保留在 silver.orders_quarantine 表中以便后续检查。

在我的示例中,我选择通过文件哈希机制来简化流程,禁止重复加载相同输入文件到青铜层。因此,如果你第二次运行该命令,会发现青铜层的数据加载部分会被完全跳过,因为文件哈希已经存在。在生产系统中,你同样需要考虑如何处理青铜层的数据重载问题。对于白银层和黄金层来说,数据重载通常不是大问题,因为这些层级的数据应该始终能够从青铜层数据中重新生成。只要确保这一点,其他环节就会顺利进行。

你可以使用以下代码直接查看奖章层级数据:

code
import duckdb
from tabulate import tabulate

def show_table(
    connection: duckdb.DuckDBPyConnection,

    query: str,
) -> None:
    result = connection.execute(query)
    headers = [column[0] for column in result.description]

    print(f"\n{title}")
    print(tabulate(result.fetchall(), headers=headers, tablefmt="psql"))

with duckdb.connect("warehouse.duckdb") as connection:
    show_table(
        connection,
        "BRONZE - Raw orders",
        """
        SELECT
            order_id,
            ordered_at,
            customer_id,
            region,
            amount,
            currency,
            status,
            source_file
        FROM bronze.orders_raw
        ORDER BY order_id
        """,
    )

    show_table(
        connection,
        "SILVER - Validated orders",
        """
        SELECT
            order_id,
            ordered_at,
            customer_id,
            region,
            amount,
            currency,
            status
        FROM silver.orders
        ORDER BY order_id
        """,
    )

show_table( connection, "SILVER - Quarantined orders", """ SELECT order_id, ordered_at, amount, rejection_reason FROM silver.orders_quarantine ORDER BY order_id """, )

show_table( connection, "GOLD - Daily sales by region", """ SELECT * FROM gold.daily_sales_by_region ORDER BY order_date, region """, )

code

Which results in the following output.

BRONZE - Raw orders +------------+----------------------+---------------+----------+----------+------------+----------+-----------------------+ | order_id | ordered_at | customer_id | region | amount | currency | status | source_file | |------------+----------------------+---------------+----------+----------+------------+----------+-----------------------| | 1001 | 2026-07-19T09:10:00Z | C001 | North | 125.5 | GBP | paid | orders_2026-07-19.csv | | 1002 | 2026-07-19T10:05:00Z | C002 | South | 89.99 | GBP | paid | orders_2026-07-19.csv | | 1002 | 2026-07-19T10:05:00Z | C002 | South | 89.99 | GBP | paid | orders_2026-07-19.csv | | 1003 | not-a-date | C003 | North | 45 | GBP | paid | orders_2026-07-19.csv | | 1004 | 2026-07-19T11:42:00Z | C004 | West | -10 | GBP | paid | orders_2026-07-19.csv | | 1005 | 2026-07-19T12:20:00Z | C005 | North | 210 | GBP | refunded | orders_2026-07-19.csv | +------------+----------------------+---------------+----------+----------+------------+----------+-----------------------+

SILVER - Validated orders +------------+---------------------------+---------------+----------+----------+------------+----------+ | order_id | ordered_at | customer_id | region | amount | currency | status | |------------+---------------------------+---------------+----------+----------+------------+----------| | 1001 | 2026-07-19 10:10:00+01:00 | C001 | NORTH | 125.5 | GBP | paid | | 1002 | 2026-07-19 11:05:00+01:00 | C002 | SOUTH | 89.99 | GBP | paid | | 1005 | 2026-07-19 13:20:00+01:00 | C005 | NORTH | 210 | GBP | refunded | +------------+---------------------------+---------------+----------+----------+------------+----------+

SILVER - Quarantined orders +------------+---------------------------+----------+--------------------+ | order_id | ordered_at | amount | rejection_reason | |------------+---------------------------+----------+--------------------| | 1002 | 2026-07-19 11:05:00+01:00 | 89.99 | duplicate order_id | | 1003 | | 45 | invalid ordered_at | | 1004 | 2026-07-19 12:42:00+01:00 | -10 | negative amount | +------------+---------------------------+----------+--------------------+

/think

GOLD - 各地区每日销售情况 +--------------+----------+------------+---------------+---------------+-------------------+------------------+ | order_date | region | currency | paid_orders | gross_sales | refunded_orders | refunded_value | |--------------+----------+------------+---------------+---------------+-------------------+------------------| | 2026-07-19 | NORTH | GBP | 1 | 125.5 | 1 | 210 | | 2026-07-19 | SOUTH | GBP | 1 | 89.99 | 0 | | +--------------+----------+------------+---------------+---------------+-------------------+------------------+

code

## 总结

作为数据库和数据工程师,我们经常在ETL作业中听到关于勋章模式(Medallion pattern)的讨论,老实说,你可能已经多次实现了它的简化版本。在本文中,我试图从基本原理出发,为你展示如何实现一个实用的勋章架构。

不要误解我。我给你展示的例子确实是一个玩具示例,它使用了有限的输入数据和本地数据库,但构建更大规模生产系统的必要原则已经具备。

在生产环境中,你必须决定是使用企业级关系型数据库管理系统(如Postgres或Oracle),还是使用基于云的对象存储(如AWS S3)。如果选择后者,你需要考虑使用哪种事务表存储格式(如Hudi、Delta表或Iceberg)。你还需要考虑是否需要使用管道编排工具(如Airflow或Dagster)。

我甚至还没有谈到分层边界所需的自动化检查类型。例如包括:

- 青铜层行数统计和数据源完整性
- 银层键唯一性、接受值检查和引用完整性
- 金层与银层总量的对账
- 新鲜度和数据量阈值
- 隔离率和模式漂移的预警

但这些只是锦上添花。关键是要理解勋章模式的基本原理,并认识到它如何以及在哪些场景下可以融入你现有的或新的ETL管道中。

勋章架构之所以有效,是因为它使数据中的差异变得可见。接收到的数据与验证后的数据不同,而验证后的数据并不自动适用于特定的业务决策。

作者:Thomas Reid

查看Thomas Reid的所有文章

数据建模

,

ETL管道

勋章架构

编程

Python

分享本文

- 在Facebook上分享
- 在LinkedIn上分享
- 在X上分享

Towards Data Science 是一份社区出版物。提交你的见解以触达全球受众,并通过TDS作者支付计划获得收益。

将href更新为你的实际投稿链接

为TDS撰写文章

✦ 结束CTA ✦