求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
相关产品推荐
相关产品推荐

