You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求Databricks Lakehouse Medallion Architecture的Python实现示例

Databricks湖屋金牌架构(Medallion Architecture)Python实现示例

金牌架构通过青铜、银、金三层数据流转,实现从原始数据到业务可用数据的渐进式处理,以下是基于PySpark和Delta Lake的实现示例:

青铜层(Bronze Layer):原始数据存储

青铜层用于保存未经处理的原始数据,支持结构化、半结构化等多种格式,通过Delta Lake存储并开启变更数据捕获(CDF),为后续增量处理提供基础。

代码示例:写入青铜层Delta表

from pyspark.sql.functions import current_timestamp

# 读取原始数据源(示例为CSV,可替换为Kafka、云存储等)
raw_data = spark.read.csv("/path/to/raw/source", header=True, inferSchema=True)

# 添加数据摄入时间戳
bronze_data = raw_data.withColumn("ingest_timestamp", current_timestamp())

# 写入Delta表,开启CDF与自动合并Schema
bronze_data.write.format("delta") \
  .option("mergeSchema", "true") \
  .option("enableChangeDataFeed", "true") \
  .mode("append") \
  .save("/delta/bronze/raw_transactions")

# 注册临时视图方便后续操作
spark.sql("CREATE OR REPLACE TEMP VIEW bronze_transactions AS SELECT * FROM delta.`/delta/bronze/raw_transactions`")

银层(Silver Layer):数据清洗与标准化

银层基于青铜层的增量变更数据,执行去重、格式修正、缺失值处理等清洗操作,输出干净、一致的标准化数据集,同样开启CDF支持后续增量聚合。

代码示例:基于CDF增量生成银层数据

from pyspark.sql.functions import to_date

# 读取青铜层的增量变更记录
bronze_changes = spark.read.format("delta") \
  .option("readChangeData", "true") \
  .option("startingVersion", 0) \
  .load("/delta/bronze/raw_transactions")

# 数据清洗逻辑
silver_data = bronze_changes \
  .dropDuplicates(["transaction_id"]) \
  .na.fill({"amount": 0.0, "status": "UNKNOWN"}) \
  .withColumn("transaction_date", to_date("transaction_timestamp")) \
  .drop("_change_type", "_commit_version", "_commit_timestamp")  # 移除CDF元数据字段

# 写入银层Delta表
silver_data.write.format("delta") \
  .option("mergeSchema", "true") \
  .option("enableChangeDataFeed", "true") \
  .mode("append") \
  .save("/delta/silver/clean_transactions")

# 注册临时视图
spark.sql("CREATE OR REPLACE TEMP VIEW silver_transactions AS SELECT * FROM delta.`/delta/silver/clean_transactions`")

金牌层(Gold Layer):业务聚合数据

金牌层面向具体业务场景,基于银层数据进行聚合建模,生成直接供报表、BI工具使用的业务级数据集。

代码示例:生成业务聚合的金牌层数据

# 按日期和地区聚合交易指标
gold_data = spark.sql("""
  SELECT 
    transaction_date,
    region,
    COUNT(transaction_id) AS total_transactions,
    SUM(amount) AS total_amount,
    AVG(amount) AS avg_transaction_amount
  FROM silver_transactions
  GROUP BY transaction_date, region
  ORDER BY transaction_date DESC
""")

# 写入金牌层Delta表(根据业务需求选择overwrite/append/merge模式)
gold_data.write.format("delta") \
  .mode("overwrite") \
  .save("/delta/gold/daily_region_transactions")

# 注册视图供BI工具访问
spark.sql("CREATE OR REPLACE VIEW gold_daily_region_transactions AS SELECT * FROM delta.`/delta/gold/daily_region_transactions`")

关键特性说明

  • 变更数据捕获(CDF):通过enableChangeDataFeed=true开启,支持增量读取变更记录,避免全量扫描,提升处理效率。
  • MergeSchema:启用后自动合并新增字段,适配原始数据结构的变化。
  • 流处理适配:上述示例可修改为readStream/writeStream模式,实现实时数据Pipeline。

内容的提问来源于stack exchange,提问作者Alexander Todorovic

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 22:30:50