Databricks 中的資料轉換

Databricks 概念

Kevin Barlow

Data Practitioner

資料工程中的 SQL

SQL

  • 資料庫管理員(DBA)熟悉
  • 適合標準操作
  • 可執行預先定義的 UDF
-- Creating a new table in SQL

CREATE TABLE table_name
USING delta
AS (
  SELECT *
  FROM source_table
  WHERE date >= '2023-01-01'
)
Databricks 概念

資料工程的其他語言

Python、R、Scala

  • 軟體工程師熟悉
  • 支援標準與複雜轉換
  • 可使用並自訂函式
#Creating a new table in Pyspark

spark
  .read
  .table('source_table')
  .filter(col('date') >= '2023-01-01')
  .write
  .saveAsTable('table_name')
Databricks 概念

常見轉換

綱要(Schema)操作

  • 新增與移除欄位
  • 重新定義欄位

#Pyspark

df
  .withColumn(col('newCol'), ...)
  .drop(col('oldCol'))

_篩選_

  • 將 DataFrame 縮減為子集
  • 可傳入多個條件

#Pyspark

df
  .filter(col('date') >= target_date)
  .filter(col('id') IS NOT NULL)
Databricks 概念

常見轉換(續)

巢狀資料

  • 陣列或 Struct 資料
  • 展開或合併
df
  .explode(col('arrayCol')) #wide to long
  .flatten(col('items')) #long to wide

_彙總_

  • 依欄位分組
  • 計算摘要統計
df
  .groupBy(col('region'))
  .agg(sum(col('sales')))
Databricks 概念

Auto Loader

Auto Loader 會在新檔案進入資料湖時處理它們。

  • 漸進式處理
  • 高效處理
  • 自動化
spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .load(file_path)

Auto Loader 示意圖

1 https://www.databricks.com/blog/2020/02/24/introducing-databricks-ingest-easy-data-ingestion-into-delta-lake.html
Databricks 概念

Structured Streaming

串流管線

spark.readStream
    .format("kafka")
    .option("subscribe", "<topic>")
    .load()
    .join(table_df, 
      on="<id>", how="left")
    .writeStream
    .format("kafka")
    .option("topic", "<topic>")
    .start()
Databricks 概念

一起來練習吧!

Databricks 概念

Preparing Video For Download...