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 개념

공통 변환(계속)

중첩 데이터

  • 배열 또는 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 개념

구조적 스트리밍

스트리밍 파이프라인

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...