데이터 집계와 조인 효율화

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

Disha Mukherjee

Lead Data Engineer

세 가지 질문, 하나의 깔끔한 데이터셋

recraft: half: 여러 화면에서 다채로운 차트와 비즈니스 보고서를 검토하는 데이터 분석가, 투명 배경

 

$$

  • 어떤 카테고리가 매출을 가장 많이 이끌까요?
  • 상위 고객은 누구인가요?
  • 부서 정보로 결과를 어떻게 보강할 수 있을까요?
Databricks에서 Spark SQL로 데이터 변환하기

groupBy()와 agg()의 작동 방식

 

$$

  • groupBy() - 키 컬럼을 기준으로 행을 분할
  • agg() - 각 그룹에 함수를 병렬로 적용
  • 둘 다 지연(lazy) 방식 - 액션이 실행을 트리거
  • agg() 한 번 호출 = 데이터 스캔 한 번

recraft: half: 데이터 행이 세 개의 레이블된 그룹으로 분할되고 각 그룹 아래에 합계가 표시된 모습, 투명 배경

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

카테고리별 매출 - 코드

category_revenue = (
    df_valid
    .groupBy("Category")

.agg( F.round(F.sum("Transaction_Amount"), 2).alias("total_revenue"), F.count("Transaction_Amount").alias("transaction_count"), F.round(F.avg("Transaction_Amount"), 2).alias("avg_transaction"), )
.orderBy(F.col("total_revenue").desc())
)
Databricks에서 Spark SQL로 데이터 변환하기

카테고리별 매출 - 출력 결과

category_revenue.show(truncate=False)
+-----------+-----------------+-----------------+---------------+
|Category   |total_revenue    |transaction_count|avg_transaction|
+-----------+-----------------+-----------------+---------------+
|Clothing   |338266605.24     |6742             |50173.04       |
|Dining     |337631718.64     |6719             |50250.29       |
|Electronics|331102843.72     |6614             |50060.91       |
|Savings    |329346609.07     |6596             |49931.26       |
|Groceries  |329003451.57     |6537             |50329.43       |
|Unknown    |880591.27        |15               |58706.08       |
+-----------+-----------------+-----------------+---------------+
Databricks에서 Spark SQL로 데이터 변환하기

차원 테이블로 데이터 보강하기

recraft: half: 두 데이터베이스 테이블이 빛나는 연결 다리로 이어진 모습, 투명 배경

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
Databricks에서 Spark SQL로 데이터 변환하기

표준 왼쪽 조인

df_joined = df_valid.join(df_dim, on="Category", how="left")

df_joined.select( "Customer_ID", "Category", "Department", "Transaction_Amount" ).show(5, truncate=False)
+-----------+-----------+----------+------------------+
|Customer_ID|Category   |Department|Transaction_Amount|
+-----------+-----------+----------+------------------+
|CUST003    |Electronics|Tech      |5752.36           |
|CUST009    |Clothing   |Retail    |28959.12          |
|CUST010    |Savings    |Finance   |72098.18          |
|CUST011    |Savings    |Finance   |49771.05          |
|CUST020    |Groceries  |Food      |69825.14          |
+-----------+-----------+----------+------------------+
Databricks에서 Spark SQL로 데이터 변환하기

숨겨진 비용 - 셔플

recraft: half: 연결된 서버와 노드 네트워크를 통해 데이터 패킷이 화살표 방향으로 이동하는 모습, 투명 배경

 

  • 일치하는 키는 반드시 같은 머신에 위치해야 합니다
  • Spark는 행을 정렬하기 위해 네트워크를 통해 데이터를 전송합니다
  • 셔플 = 대규모 처리 시 병목
Databricks에서 Spark SQL로 데이터 변환하기

Spark UI

recraft: full: 작업 단계별 수평 막대 차트, 진행률 표시줄, 데이터 처리량 지표가 표시된 어두운 분석 모니터링 대시보드, 투명 배경

 

  • JobsStages - 실행 내용과 소요 시간 확인
  • Shuffle Read/Write - 이동한 데이터 양 확인
  • 파이프라인 병목 진단에 유용
Databricks에서 Spark SQL로 데이터 변환하기

.explain()으로 쿼리 계획 읽기

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain() - 실행 계획 확인
  • 노드 15 & 16 - 라우팅 인수 없음
Databricks에서 Spark SQL로 데이터 변환하기

브로드캐스트 조인 - 해결책

# Wrap the small table in F.broadcast()
df_broadcast = df_valid.join(
    F.broadcast(df_dim),
    on="Category",
    how="left"
)

print(f"Broadcast joined rows: {df_broadcast.count():,}")
Broadcast joined rows: 33,223
Databricks에서 Spark SQL로 데이터 변환하기

변경 전후 - 실행 계획 확인

df_broadcast.explain(mode="formatted")
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...            
(15) PhotonShuffleExchangeSink
Arguments: SinglePartition
(16) PhotonShuffleMapStage
Arguments: EXECUTOR_BROADCAST, [id=#11839]
  • 차원 테이블에 라우팅 지침이 적용됨
  • 대형 테이블 df_valid이동하지 않음
  • 단 한 번의 변경 = 대규모 처리 시 분 단위에서 초 단위로
Databricks에서 Spark SQL로 데이터 변환하기

연습해 봅시다!

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

Preparing Video For Download...