湖流管道中的维度建模

维度建模是一种将黄金层数据组织为事实表和维度表的技术,使分析师和商业智能(BI)工具能够高效地查询数据。 本页解释了如何用Lakeflow管道构建该模型。

Overview

维度建模将数据分为两种表格:

  • 事实表存储了你关心的事件或度量值,例如订单、点击或销售。 每一行代表该事件的一次发生,主要用键和数字度量来描述。
  • 维度表存储这些事件的描述性上下文信息,例如客户、产品或日期。 每一行对应一个业务实体。

星形模式是你把一个事实表放在中间,并通过它们的键连接到多个维度表时得到的形状。 这种布局便于分析师和商业智能工具查询,工程师也易于推理,因为每个表格都有明确的单一职责。

在 Lakeflow 管道中,星型架构很自然地适用于奖牌架构的黄金层。 青铜数据集和白银数据集负责数据摄取和清洗,而黄金数据集则将事实表和维度表实体化,从而让下游用户可以直接查询这些表。 因为流水线会让这些表逐步更新,你就能实现星型模式的查询简便性,无需在BI层单独进行提取、转换、加载(ETL)步骤。

工作原理

你在管道中将维度和事实构建为数据集,并根据各自的变化方式选择相应的数据集类型。 对于大多数金层模型:

  • 维度表 构建为物化视图(或者在需要保留历史记录时,构建为采用缓慢变化维度 (SCD) 2 型的流式表)。 物化视图会在输入发生变化时,基于清理后的银层数据高效地重新计算,并为每个业务实体生成一行数据。
  • 事实表 构建为由银层增量馈送的流式表,以使金层聚合尽可能接近实时。 事实通过键数来指代其维度,而非重复描述属性。

关于这两种数据集类型的更多信息,请参见 “实体化视图 ”和 “流表”。 要在某个维度中跟踪历史记录,请参见 AUTO CDC API:使用管道简化变更数据捕获

密钥与替代密钥

优先使用自然 (源数据中已存在的标识符,如订单号),其中源的自然键稳定且可用,因为它能够很好地聚类和连接。 只有当源重复使用或更改 ID 时,才会使用 替代密钥 (流水线生成的替代标识符)。

当你需要代理密钥时,避免使用像 sha2(natural_key)这样的哈希代理。 哈希是故意随机的,这对液体聚类和Z阶性能不利,因为物理相邻的行最终会分散在文件中。 相反,从稳定的自然密钥确定性地推导出一个保持顺序的代理,使同一业务实体始终映射到同一个代理。 确定性键在对维度进行完全刷新或重建后仍可保留,从而保持现有事实表与维度表之间的联接不受影响。

或者,当上游表仅追加且从未进行全量刷新时,你也可以使用 IDENTITY 列。 由于 IDENTITY 值是在插入行时分配的,因此重建可能会为同一实体重新分配不同的 ID,并悄然破坏仍使用旧值的事实表与维度表之间的联接。

日期维度

构建一个dim_date,作为通过sequence()explode()在某个日期范围内生成的简单物化视图,而不是从数据源摄取它。 它是静态的参考数据,计算成本低,而且简化了基于日期的连接和模型其他部分的窗口设置。

示例

以下示例构建了一个带有客户维度和订单事实表的小星型模式。

维度表

维度表通常是基于清洗后的银层数据构建的物化视图,其中每个业务实体对应一行数据,如以下代码所示:

Python

from pyspark import pipelines as dp

@dp.materialized_view(name="dim_customer", comment="Customer dimension")
def dim_customer():
    return (
        spark.read.table("customers_silver")
        .select("customer_id", "customer_name", "region", "signup_date")
    )

SQL

CREATE OR REFRESH MATERIALIZED VIEW dim_customer
COMMENT "Customer dimension"
AS SELECT customer_id, customer_name, region, signup_date
FROM customers_silver;

事实数据表

事实表保存可测量事件,通过键数引用维度,而非重复描述属性。 保持事实表精简(主要包含键和数值指标),并在查询时使用联接引入描述性详细信息,如以下代码所示:

Python

from pyspark import pipelines as dp

@dp.table(name="fact_orders", comment="One row per order line, keyed to dimensions")
def fact_orders():
    return (
        spark.readStream.table("orders_silver")
        .select(
            "order_id",
            "customer_id",       # foreign key to dim_customer
            "product_id",        # foreign key to dim_product
            "order_date",        # foreign key to dim_date
            "quantity",
            "amount",
        )
    )

SQL

CREATE OR REFRESH STREAMING TABLE fact_orders
COMMENT "One row per order line, keyed to dimensions"
AS SELECT
  order_id,
  customer_id,   -- foreign key to dim_customer
  product_id,    -- foreign key to dim_product
  order_date,    -- foreign key to dim_date
  quantity,
  amount
FROM STREAM(orders_silver);

最佳做法

以下几种做法有助于让星型模式在不断扩展时保持良好状态:

  • 将事实表保留为流式表,将维度表保留为实体化视图,除非你确实需要变更历史,在这种情况下,请使用 AUTO CDCSTORED AS SCD TYPE 2。 请参阅 AUTO CDC API:使用管道简化变更数据捕获
  • 使用下游 BI 工具直接查询黄金层物化视图。 Lakeflow 流水线保持它们的增量刷新,因此你可以在无需单独报告 ETL 步骤的情况下获得近实时结果。
  • 将模型维度和事实作为独立的流程,汇入同一黄金层,因此每个数据集都可以作为一个连贯DAG的一部分进行调度、检查点和刷新。 请参阅 使用 Lakeflow 管道流以增量方式加载和处理数据

其他资源