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 概念

常见转换(续)

嵌套数据

  • 数组或结构体数据
  • 展开或合并
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 概念

结构化流处理

流式管道

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

Passons à la pratique !

Databricks 概念

Preparing Video For Download...