26Y27e3

QK·2026년 7월 27일

이 에러가 발생하는 이유는 파일별로 컬럼의 '순서(Order)'가 달라졌기 때문입니다.

with_columns로 누락된 컬럼(mem_request 등)을 개별 생성할 때 새로 만든 컬럼들이 맨 뒤에 붙으면서, file_afile_b 간의 컬럼 위치 순서가 일치하지 않아 pl.concat() 시 스키마 불일치(SchemaError)가 발생한 것입니다.


🛠️ 완치 해결책: 표준 컬럼 순서(Canonical Column Order) 재정렬

각 파일별로 누락 컬럼을 보정한 직후, 미리 정해둔 기준 컬럼 순서(target_columns)로 .select()하여 순서를 동일하게 맞춰주면 간단히 해결됩니다.

step2_pipeline.py 루프 내부 수정 코드

for date_chunk in target_dates: 내부의 file_lfs 생성 부분을 아래 코드로 업데이트해 주세요.

    # ─── 🔄 [OOM 차단벽 내부 일자별 배칭 루프 개시] ───────────────────
    # 컬럼 순서를 일관되게 고정할 표준 스키마 정의
    target_columns = base_fields + metric_fields

    for date_chunk in target_dates:
        day_files = sorted(list(Path(RAW_DIR).glob(f"prom_raw_{date_chunk}_*.parquet")))
        
        if not day_files:
            print(f"\n   ⚠️  [스킵] {date_chunk} 일자의 로컬 Parquet 원부가 없어 다음 날짜로 롤오버합니다.")
            continue
            
        print(f"\n📅 [일자 격리 가동] {date_chunk} 파티션 대용량 연산 스트리밍 로드 시작...")
        
        file_lfs = []
        for file_path in day_files:
            try:
                f_schema = pl.scan_parquet(file_path).collect_schema().names()
            except AttributeError:
                f_schema = pl.scan_parquet(file_path).schema.keys()

            f_select = [c for c in target_columns if c in f_schema]
            f_missing = [c for c in target_columns if c not in f_schema]

            # 1. 존재하는 컬럼 프로젝션
            lf_single = pl.scan_parquet(file_path).select(f_select)
            
            # 2. 누락된 컬럼 0.0 보정
            if f_missing:
                lf_single = lf_single.with_columns([pl.lit(0.0).alias(c) for c in f_missing])
            
            # 💡 3. [핵심 수정] 모든 파일의 컬럼 순서를 target_columns 순서로 정렬 보정!
            lf_single = lf_single.select(target_columns)

            file_lfs.append(lf_single)

        # 🚀 모든 파일의 컬럼 명칭과 '순서'가 완벽히 일치하므로 에러 없이 바로 병합됩니다.
        lf_day = pl.concat(file_lfs)

        # 🛡️ [OOM 방어 제3선]: 문자열 카테고리 캐싱 압축 주입
        lf_day = lf_day.with_columns([
            pl.col("namespace").cast(pl.Categorical),
            pl.col("node").cast(pl.Categorical),
            pl.col("pod").cast(pl.Categorical),
            pl.col("container").cast(pl.Categorical)
        ])
        
        # 15분 청크 수집 경계면 겹침 방지용 고유 가드레일 주입
        lf_day = lf_day.unique(subset=["timestamp", "namespace", "pod", "container"])
        
        # 🚀 [OOM 방어 제4선]: collect 시 streaming=True 실행
        df_day = lf_day.collect(streaming=True)
        
        # (이후 기존 df_day 집계/계산 로직 진행...)

💡 포인트

마지막에 lf_single.select(target_columns) 구문을 거치게 되면, 어떤 시간대 파일에 컬럼이 빠져있었더라도 병합 직전 모든 LazyFrame의 컬럼 명과 순서가 100% 동일해지므로 SchemaError 없이 깔끔하게 Concat 및 병렬 스트리밍 집계가 완료됩니다!

profile
engineer

0개의 댓글