이번 강은 ELT의 두 번째로, 중첩 데이터와 고차 함수·UDF 같은 고급 변환을 다룬다.
실제 데이터는 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
댓글은 운영자만 작성할 수 있어요.