Databricks 데이터 엔지니어

Databricks 데이터 엔지니어 5강

[Associate] Spark SQL로 ELT (2) — 고급 변환·UDF

이번 강은 ELT의 두 번째로, 중첩 데이터와 고차 함수·UDF 같은 고급 변환을 다룬다.

Databricks 5강 개념도

실제 데이터는 JSON처럼 배열·구조체가 중첩된 경우가 많다. Spark SQL은 이를 다루는 강력한 함수를 제공한다.

이 강의 목표는 다음과 같다.

  • 중첩 구조(struct·array)에 접근·전개한다.
  • 고차 함수로 배열을 변환한다.
  • UDF와 PySpark↔SQL 상호운용을 이해한다.

1. 중첩 구조 다루기

  • 구조체 필드: col.field
  • 배열 원소: col[0]
  • 배열을 행으로 전개: explode()
SELECT user_id, explode(items) AS item
FROM orders;                       -- 배열 items를 행으로 분해

2. JSON 문자열 파싱

문자열로 들어온 JSON은 스키마를 주고 파싱한다.

SELECT from_json(payload, 'id INT, name STRING') AS j FROM raw;

3. 고차 함수 (Higher-order functions)

배열 각 원소에 함수를 적용한다. 전개 없이 배열 채로 처리한다.

SELECT transform(scores, x -> x * 10) AS scaled,      -- 각 원소 변환
       filter(scores, x -> x >= 60) AS passed          -- 조건 필터
FROM tests;

4. UDF와 상호운용

SQL 내장 함수로 부족하면 UDF(사용자 정의 함수)로 확장한다.

from pyspark.sql.functions import udf
@udf("string")
def mask(s): return s[:3] + "***"
spark.udf.register("mask", mask)   # SQL에서도 사용 가능

주의: UDF는 Spark의 최적화(코드 생성·푸시다운) 혜택이 줄어 느릴 수 있다. 가능하면 내장 함수·고차 함수를 우선한다. PySpark DataFrame API와 SQL은 언제든 바꿔 쓸 수 있다.

예제

예제) 주문 테이블의 items 컬럼이 상품 배열이다. 상품별 판매 건수를 세려면?

해설) 먼저 explode(items)로 배열을 상품 단위 행으로 펼친다. 그다음 상품 키로 GROUP BY 해 건수를 센다. 배열을 그대로 두고 조건 처리만 필요하면 explode 대신 filter 같은 고차 함수를 쓰는 것이 효율적이다.

샘플 문제

문) 배열 컬럼을 각 원소가 하나의 행이 되도록 펼치는 함수는?

  • (A) transform()
  • (B) explode()
  • (C) filter()
  • (D) from_json()

정답: (B). explode()가 배열을 행으로 전개한다. transform/filter는 배열을 유지한 채 원소를 변환·필터하는 고차 함수, from_json은 JSON 문자열 파싱이다.

정리

  • 중첩 접근은 col.field·col[i], 배열 전개는 explode().
  • from_json으로 JSON 문자열을 구조화한다.
  • 고차 함수(transform·filter)로 배열을 전개 없이 처리한다.
  • UDF는 확장용이나 최적화가 약하니 내장 함수를 우선한다.

다음 강에서는 스트리밍과 Auto Loader로 증분 처리를 시작한다.

댓글 0