26Y27e2

QK·2026년 7월 27일

원인을 정확히 찾았습니다. 와일드카드 파일 패턴(prom_raw_20260723_*.parquet)을 pl.scan_parquet()으로 한 번에 읽을 때 발생한 와일드카드 스캔의 맹점 때문입니다.

collect_schema()는 와일드카드 내 첫 번째 파일의 스키마만 대표로 추출합니다. 만약 00시 파일에는 mem_request 컬럼이 있어서 select_fields에 포함되었는데, 05시 파티션 파일에는 해당 컬럼이 없었다면 스팀을 읽어 내려가는 도중 Polars 스캐너가 ColumnNotFoundError를 터뜨리게 됩니다.


🛠️ 완치 해결책: 파일별 LazyFrame 스키마 보정 후 pl.concat()

와일드카드로 뭉뚱그려 스캔하지 않고, 하루치 24개 Parquet 파일 각각의 스키마를 개별 검사하여 규격화한 뒤 pl.concat()으로 결합하면 이 문제가 근본적으로 완치됩니다.

LazyFrame 상태에서 .concat()을 수행하므로 메모리가 전혀 늘어나지 않고 쿼리 플랜만 깔끔하게 병합됩니다.

step2_pipeline.py 루프 내부 수정 코드

for date_chunk in target_dates: 내의 파일 로드 로직을 아래 코드로 교체해 주세요.

    # ─── 🔄 [OOM 차단벽 내부 일자별 배칭 루프 개시] ───────────────────
    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} 파티션 대용량 연산 스트리밍 로드 시작...")
        
        # 🛡️ [개별 파일 스키마 규격화 가드레일]
        # 시간대별 파일 간 컬럼 유무 불일치(특정 시간대 mem_request 누락 등) 완벽 방어
        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 base_fields + metric_fields if c in f_schema]
            f_missing = [c for c in metric_fields if c not in f_schema]

            # 존재하는 컬럼만 프로젝션 로드
            lf_single = pl.scan_parquet(file_path).select(f_select)
            
            # 누락된 컬럼만 0.0 리터럴로 개별 가상화 보정
            if f_missing:
                lf_single = lf_single.with_columns([pl.lit(0.0).alias(c) for c in f_missing])
            
            file_lfs.append(lf_single)

        # 🚀 모든 시간대 파일의 스키마를 동일하게 맞춘 후 Lazy 상태로 병합 (메모리 점유 0B)
        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 가공 및 그룹바이 파이프라인 진행...)

💡 이 방식의 장점

  1. 완벽한 컬럼 동기화: prom_raw_20260723_01.parquetmem_request가 있고 02.parquet에 없더라도, 각각의 파일 수준에서 0.0으로 맞춘 뒤 Concat 하므로 Polars 스캐너가 컬럼 에러를 절대 발생시키지 않습니다.
  1. 실제 데이터 보존: mem_request가 존재하는 시간대의 데이터는 실제 값 그대로 유지되고, 수집되지 않은 시간대만 0.0으로 채워지므로 데이터 손실이 없습니다.
profile
engineer

0개의 댓글