使用 workflows 构建生产流水线

在 Databricks 中使用 Spark SQL 进行数据转换

Disha Mukherjee

Lead Data Engineer

为何选择 Delta Lake?

安全发光的数字金库,内含整齐的数据表和防护盾牌,现代扁平风格

 

$$

  • ACID 事务 → 回滚失败写入
  • 模式约束 → 拦截类型不匹配
  • 版本控制 → 查询任意历史状态
在 Databricks 中使用 Spark SQL 进行数据转换

写入 Delta

df_valid.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("transactions_clean")

print(f"Rows written: {df_valid.count():,}")
Rows written: 33,223

$$

$$

  • Unity Catalog 中出现一个新的 Delta 表
在 Databricks 中使用 Spark SQL 进行数据转换

Notebook 任务

task1_ingest:加载并清洗

  • 加载 CSV → 清洗 → 写入 Delta 表
在 Databricks 中使用 Spark SQL 进行数据转换

Notebook 任务

task2_metrics:按类别统计收入

  • 读取清洗表 → 计算指标 → 写入新表
在 Databricks 中使用 Spark SQL 进行数据转换

Notebook 任务

task3_customers:按消费额排名

  • 读取清洗表 → 排名客户 → 保存到新表
在 Databricks 中使用 Spark SQL 进行数据转换

创建作业

Jobs 与 Pipelines 界面,显示三任务作业 DAG,依赖从 task1_ingest 到 task2_metrics 再到 task3_customers

  • 可手动运行、定时调度并设置触发器
在 Databricks 中使用 Spark SQL 进行数据转换

运行作业

作业运行 DAG,task1_ingest 与 task2_metrics 绿色成功,task3_customers 红色失败

在 Databricks 中使用 Spark SQL 进行数据转换

运行作业

检查任务三中的错误

  • 错误:customer-id 应为 Customer_ID
在 Databricks 中使用 Spark SQL 进行数据转换

运行作业

图形视图(成功)

时间线视图

在 Databricks 中使用 Spark SQL 进行数据转换

什么是 Lakeflow?

$$

$$

对比:命令式(Jobs)| 声明式(Lakeflow)

 

  • Jobs → 我们逐步管理
  • Lakeflow → 仅声明表应包含的内容
  • Databricks 负责顺序、重试与算力
在 Databricks 中使用 Spark SQL 进行数据转换

@dlt.table 模式

@dlt.table(name="transactions_bronze")
def transactions_bronze():
    return spark.read.format("csv").schema(schema).load(FILE_PATH)

@dlt.table(name="transactions_silver") def transactions_silver(): return dlt.read("transactions_bronze").na.drop(...).filter(...)
@dlt.table(name="category_revenue_gold") def category_revenue_gold(): return dlt.read("transactions_silver").groupBy("Category").agg(...)
在 Databricks 中使用 Spark SQL 进行数据转换

流水线运行

Lakeflow 流水线 DAG,显示 transactions_bronze 10 万行,transactions_silver 3.3 万行,category_revenue_gold 6 行,全部为绿色物化视图

  • 铜层(100K 行)→ 银层(33K 行)→ 金层(6 行)
在 Databricks 中使用 Spark SQL 进行数据转换

Notebooks、Jobs,还是 Lakeflow?

$$

分层:Notebooks、Databricks Jobs、Lakeflow Pipelines

 

$$

  • Notebooks → 探索、原型
  • Databricks Jobs → 多步骤、定时流水线
  • Lakeflow → 全托管、声明式
在 Databricks 中使用 Spark SQL 进行数据转换

让我们一起练习吧!

在 Databricks 中使用 Spark SQL 进行数据转换

Preparing Video For Download...