레거시 분석 도구의 전처리 SQL 을 CDE 위의 PySpark 코드로 옮기면서 실제로 결과가 달라질 뻔했던 지점들을 모았다. 변환 작업 전반의 규칙은 레거시 코드를 PySpark 로 변환할 때의 규칙 에 있다.
레거시 도구는 ${=VAR} 같은 자체 치환 문법을 쓴다. 이 문법은 Spark 가 모른다. 코드 최상단에 파이썬 변수로 올리고 f-string 이나 파라미터로 넘긴다.
PART_YMD = "20260317" # DAG 에서 전달받는 기준일 (yyyyMMdd)
REGN = "KR" # 지역 구분
TARGET_TBL = "silver.filtered_pt_login_hist"
원본이 CURRENT_DATE() 로 오늘을 잡고 있다면 그대로 옮기면 안 된다. 재실행할 때 결과가 달라진다. 전달받은 기준일로 바꾸고, 기준일이 원본의 어느 날짜에 해당하는지 확인한다. 수집일 기준인지 데이터 발생일 기준인지에 따라 하루가 어긋난다.
from datetime import datetime, timedelta
_part = datetime.strptime(PART_YMD, "%Y%m%d")
_from = (_part - timedelta(days=1)).strftime("%Y-%m-%d")
_to = _part.strftime("%Y-%m-%d")
레거시 원천에서 넘어온 문자열에 \u0000 이 붙어 있는 경우가 있다. 원본 SQL 이 비교할 때마다 REPLACE(COL, '\u0000', '') 를 반복 호출하고 있으면, 그 정제를 읽은 직후 한 번만 하고 이후 로직은 깨끗한 값을 쓰게 만든다.
df = df.withColumn("login_type_cd", F.regexp_replace("login_type_cd", "\u0000", ""))
정제를 여러 곳에 흩어 두면 한 군데만 빠졌을 때 조건 하나가 조용히 어긋난다.
substr 의 시작 위치Spark SQL 과 Hive 의 substr · substring 은 1 기반이다. pos 에 0 을 주면 1 과 같게 처리되므로 substr(col, 0, 4) 와 substr(col, 1, 4) 의 결과는 같다. 레거시 코드에 0 이 들어 있어도 결과는 바뀌지 않지만, 읽는 사람이 0 기반이라고 오해하지 않도록 변환하면서 1 로 명시한다.
음수 pos 는 뒤에서부터 센다. 원본이 음수를 쓰고 있으면 의미를 확인한다.
Hive 컬럼명은 소문자다. 원본 SQL 이 대문자로 써 있어도 변환 코드는 쿼리와 row["..."] 참조를 모두 소문자로 통일한다. 대소문자가 어긋나면 실행 시점에 AnalysisException 으로 드러난다.
데이터 값에는 적용하지 않는다. dept_cd = 'TOP' 같은 비교는 그대로 둔다. 다만 적재 ETL 이 값을 소문자로 바꿔 저장하는 구조라면 예외이므로, 원천에서 적재까지의 변환 규칙을 한 번 확인한다.
원본이 별칭으로 참조하던 컬럼명이 실제 Hive 테이블에서 다른 이름인 경우도 흔하다. 조인하기 전에 대상 테이블의 실제 컬럼명을 확인하고 매핑표를 변환 보고서에 남긴다.
원본 SQL 에 WHERE 가 없더라도, 대상이 Hive 파티션 테이블이면 파티션 필터를 넣는다. 빠뜨리면 전체 테이블을 읽어 가상 클러스터 자원을 소진한다.
WHERE part_ymd = '{PART_YMD}'
단, 필터를 추가하면 결과 범위가 달라질 수 있다. 원본이 전체 조회를 의도한 것이라면 필터를 넣는 순간 다른 결과가 나온다. 판단이 필요한 지점이므로 임의로 넣지 말고 확인한 뒤 보고서에 남긴다.
Kudu 테이블은 Hive 파티션 프루닝 대상이 아니므로 일반 WHERE 조건만으로 충분하다. 원천이 어느 쪽인지 먼저 확인한다.
row_number 의 정렬 기준원본이 이런 형태로 중복을 걷어내는 경우가 많다.
SELECT * FROM (
SELECT *, row_number() OVER (PARTITION BY epid, login_ts ORDER BY epid) AS rownum
FROM base
) x
WHERE x.rownum = 1
ORDER BY 가 파티션 키와 같은 컬럼이면 순서가 결정되지 않는다. 같은 그룹 안의 행이 모두 같은 정렬 값을 가지므로 어느 행이 남을지 실행마다 달라질 수 있다. 원본 도구에서는 우연히 일정하게 나왔더라도 Spark 에서는 보장되지 않는다.
같은 값이 여러 벌 들어 있어 어느 것이 남아도 상관없는 경우라면 그대로 둔다. 그렇지 않다면 결정적인 정렬 키(적재 시각 등)를 추가하고 그 변경을 보고서에 적는다.
원본이 지역 코드로 분기해 UTC 에 오프셋을 더하는 방식이면 그대로 옮기되, 정의되지 않은 지역은 보정하지 않는다는 점을 확인한다. CASE 의 ELSE 가 원본 값을 그대로 돌려주는 구조면, 나중에 지역이 추가돼도 조용히 보정 없이 지나간다.
result = df.withColumn(
"regn_ts",
F.when(F.lit(REGN) == "KR", F.from_unixtime(F.unix_timestamp("login_ts") + 9 * 3600))
.when(F.lit(REGN) == "US", F.from_unixtime(F.unix_timestamp("login_ts") - 5 * 3600))
.otherwise(F.col("login_ts")),
)
고정 오프셋은 서머타임을 반영하지 않는다. 원본이 그렇게 돼 있으면 동일 출력 원칙상 그대로 두되, 알고 있는 제약으로 기록한다.
단계가 길면 중간 결과를 stage 테이블에 저장하고 다음 단계가 그 테이블을 읽게 나눈다.
STAGE_TBL_01 = "stage.<pipeline>_01_<desc>"
STAGE_TBL_02 = "stage.<pipeline>_02_<desc>"
장점은 단계별로 원본과 건수를 대조할 수 있다는 것이다. 어느 단계에서 행이 사라졌는지 바로 드러난다. 단점은 쓰기가 늘어나는 것이므로, 검증이 끝나면 병합할 단계를 정리한다. 실제로 마지막 두 단계는 WHERE 조건을 합쳐 하나의 쿼리로 처리할 수 있었다.
같은 입력으로 원본과 변환본을 돌려 다음을 대조한다. 대조 없이 "같다" 고 보고하지 않는다.
[ ] 최종 건수
[ ] 그룹별 건수 (회사·부서 등 주요 키)
[ ] 주요 수치 컬럼의 합계
[ ] null 비율
[ ] 중복 제거 후 남은 행의 표본 비교