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 の基本概念

一般的な変換

スキーマ操作

  • 列の追加・削除
  • 列の再定義

#Pyspark

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

フィルタリング

  • DataFrame を部分集合に絞る
  • 複数条件を指定

#Pyspark

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

一般的な変換(続き)

ネストデータ

  • Array または 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 の基本概念

Let's practice!

Databricks の基本概念

Preparing Video For Download...