레거시 시스템의 MariaDB SQL · Python · Scala 전처리 코드를 Cloudera Data Engineering(CDE) 위의 PySpark 코드로 옮기는 작업에서, 변환 결과가 원본과 같은 출력을 내도록 정한 규칙이다. Claude 프로젝트의 지침 문서(CONVERSION_INSTRUCTION.md)로 두고 코드 변환을 시킬 때 적용했다. 원본 테이블은 이미 Hive Parquet 테이블로 이관돼 있고, 코드만 레거시에서 가져오는 상황을 전제로 한다.
result 변수에 담는다. 전처리 로직의 최종 출력도 result 에 담는다.PART_YMD, TARGET_TBL, BASE_YMD 같은 사용자 입력 변수는 덮어쓰지 않으며, 값을 바꿔 시험하기 쉽도록 코드 맨 위에 둔다.변환 코드는 다음 기본 설정을 전제로 작성한다.
from pyspark.sql import functions as F
from pyspark.sql import types as T
spark = SparkSession.builder.appName(app_name).enableHiveSupport().getOrCreate()
spark.conf.set("spark.sql.parquet.enableVectorizedReader", "false")
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
partitionOverwriteMode=dynamic 은 INSERT OVERWRITE 시 대상 파티션만 갈아끼우게 한다. 벡터화 리더를 끄는 것은 레거시 Parquet 파일의 타입 불일치(예: decimal · timestamp 표현 차이)로 읽기 오류가 나는 것을 피하기 위한 보수적 설정이다.
input[번호] 는 데이터 소스를 뜻하며 번호는 0 부터 시작한다.input[번호] 는 해당 Hive Parquet 테이블을 읽는 spark.sql(SRC_QUERY) 로 바꾼다.WHERE part_ymd = '{PART_YMD}')를 반드시 넣어 파티션 프루닝이 되게 하고 전체 스캔을 피한다.cde job create · cde job run), CDE Resource(파일 · Python 환경 · 커스텀 이미지), Virtual Cluster, CDE Airflow Operator, 고유 REST API 를 갖는다. 코드와 제출 명령은 CDE 와 Spark 버전을 명시하고 공식 문서로 확인한 기능만 쓴다.AnalysisException 으로 드러나므로 변환 단계에서 일괄 소문자화한다.