워크플로를 활용한 프로덕션 파이프라인

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로 데이터 변환하기

노트북 태스크

task1_ingest: 로드 및 정제

  • CSV 로드 → 정제 적용 → Delta 테이블 저장
Databricks에서 Spark SQL로 데이터 변환하기

노트북 태스크

task2_metrics: 카테고리별 매출

  • 정제된 테이블 읽기 → 지표 계산 → 새 테이블에 저장
Databricks에서 Spark SQL로 데이터 변환하기

노트북 태스크

task3_customers: 지출액 기준 순위

  • 정제된 테이블 읽기 → 고객 순위 산정 → 새 테이블에 저장
Databricks에서 Spark SQL로 데이터 변환하기

작업 생성

task1_ingest에서 task2_metrics, task3_customers로 연결되는 의존성 화살표가 있는 세 태스크 완료 DAG를 보여주는 Jobs and Pipelines UI

  • 수동 실행, 작업 예약, 트리거 설정 가능
Databricks에서 Spark SQL로 데이터 변환하기

작업 실행

task1_ingest와 task2_metrics는 녹색으로 성공, task3_customers는 빨간색으로 실패한 작업 실행 DAG

Databricks에서 Spark SQL로 데이터 변환하기

작업 실행

태스크 3의 오류 확인

  • 오류: customer-idCustomer_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로 데이터 변환하기

파이프라인 실행

transactions_bronze 10만 행, transactions_silver 3만 3천 행, category_revenue_gold 6행이 모두 녹색 구체화 뷰로 표시된 Lakeflow 파이프라인 DAG

  • Bronze(10만 행) → Silver(3만 3천 행) → Gold(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...