조직도처럼 상하위 관계를 가진 표에 Nested Set 모델의 좌·우 번호(l_no · r_no)를 붙이는 로직을 PySpark 로 옮길 때의 방법과 함정을 정리한다. 원본은 pandas 로 전역 카운터를 쓰며 재귀 DFS 를 도는 코드였다.
Nested Set 번호 부여는 본질적으로 순차 처리다. 전역 카운터를 증가시키며 깊이 우선으로 순회해야 하므로, 순수 분산 연산만으로 같은 결과를 보장하기 어렵다. 계층 데이터를 드라이버로 모아 순회한 뒤 결과를 다시 DataFrame 으로 만드는 방식이 Spark 에서 계층 구조를 다루는 일반적인 패턴이다. 계층 데이터가 드라이버 메모리에 들어갈 규모일 때만 쓴다.
변환 작업 전반의 규칙은 레거시 코드를 PySpark 로 변환할 때의 규칙 에 있다.
가장 먼저 깨지는 지점이다. 원본 SQL 에 ORDER BY comp_cd, org_lvl 이 있고 로직이 순회 순서에 의존한다 — dept_cd = 'TOP' 이고 org_lvl = 0 인 행을 만나면 새 회사가 시작되는 구조이므로, 정렬이 깨지면 트리 구성 자체가 틀어진다.
Spark DataFrame 은 명시적 정렬 없이는 순서를 보장하지 않는다. 앞 단계가 정렬된 결과를 담고 있더라도 collect() 시점에 그 순서가 유지된다는 보장이 없다.
rows = (
df.orderBy(F.col("comp_cd").asc(), F.col("org_lvl").asc())
.collect()
)
정렬을 SQL 안에 넣어도 되지만, 어느 쪽이든 collect() 직전에 정렬이 명시돼 있어야 한다.
원본은 맵 구성 · DFS 순회 · 결과 조합을 한 함수가 처리했다. 세 가지로 나누면 검증도 쉬워진다.
| 함수 | 책임 |
|---|---|
parse_hierarchy_by_company |
행 목록 → 회사별 {상위부서: [하위부서, ...]} 맵 |
assign_nested_set |
맵을 DFS 로 순회하며 l_no · r_no 부여 |
| 메인 흐름 | 쿼리 실행 → 파싱 → 번호 부여 → DataFrame 생성 → 원본과 조인 |
전역 변수 대신 클로저 지역 변수(nonlocal)로 카운터를 감싸면 함수 밖으로 상태가 새지 않는다.
원본 pandas 코드는 comp · dept · l_no · r_no 네 컬럼만 만들었지만, 실제 적재 대상 테이블은 조회한 컬럼 전체에 번호 두 개가 붙은 형태인 경우가 많다. 번호 결과를 원본 DataFrame 과 comp_cd + dept_cd 기준으로 조인해 대상 스키마 순서대로 select 한다.
조인을 inner 로 두면 번호가 부여되지 않은 행 — 키 컬럼에 null 이 있어 순회에서 건너뛴 행 — 이 결과에서 빠진다. 원본 pandas 코드도 같은 동작이므로 동일 출력이 유지되지만, 건수가 줄어드는 이유를 모르면 나중에 오해한다. 변환 보고서에 적어 둔다.
from pyspark.sql import functions as F
from pyspark.sql import types as T
PART_YMD = "20260311"
TARGET_TBL = "silver.dept_nested_set"
SRC_QUERY = f"""
SELECT comp_cd, dept_cd, upr_dept_cd, org_lvl, dept_nm
FROM bronze.eem_torg
WHERE part_ymd = '{PART_YMD}'
ORDER BY comp_cd ASC, org_lvl ASC
"""
def parse_hierarchy_by_company(rows):
"""행 목록을 회사별 {상위부서: [하위부서]} 맵 목록으로 나눈다."""
result, current, comp = [], {}, ""
for row in rows:
comp_cd = row["comp_cd"]
dept_cd = row["dept_cd"]
upr = row["upr_dept_cd"]
lvl = row["org_lvl"]
if None in (comp_cd, dept_cd, upr, lvl):
continue
if dept_cd == "TOP" and lvl == 0:
if comp:
result.append((comp, current))
comp, current = comp_cd, {}
current.setdefault(upr, []).append(dept_cd)
if comp and current:
result.append((comp, current))
return result
def assign_nested_set(comp_cd, tree, root="TOP"):
"""DFS 로 순회하며 (comp_cd, dept_cd, l_no, r_no) 를 만든다."""
out = []
idx = 1
def walk(key):
nonlocal idx
for child in tree.get(key, []):
left = idx
idx += 1
walk(child)
right = idx
idx += 1
out.append((comp_cd, child, left, right))
walk(root)
return out
schema = T.StructType([
T.StructField("comp_cd", T.StringType()),
T.StructField("dept_cd", T.StringType()),
T.StructField("l_no", T.IntegerType()),
T.StructField("r_no", T.IntegerType()),
])
df = spark.sql(SRC_QUERY)
rows = df.collect()
records = []
for comp_cd, tree in parse_hierarchy_by_company(rows):
records.extend(assign_nested_set(comp_cd, tree))
nested = spark.createDataFrame(records, schema=schema)
result = (
df.alias("a")
.join(nested.alias("b"), ["comp_cd", "dept_cd"], "inner")
.select("a.comp_cd", "a.dept_cd", "a.upr_dept_cd", "a.org_lvl",
"a.dept_nm", "b.l_no", "b.r_no")
)
result.write.mode("overwrite").saveAsTable(TARGET_TBL)
재귀 깊이가 깊은 조직도라면 파이썬 기본 재귀 한도에 걸릴 수 있다. 명시적 스택으로 바꾸거나 한도를 올린다.
collect() 가 전제이므로 계층 데이터가 커지면 드라이버가 죽는다. 조직도 수준(수천~수만 행)이 아니면 다른 접근이 필요하다.dept_cd == "TOP" 같은 값 비교는 그대로 둔다.r_no 최대값은 노드 수 × 2 가 된다.