26Z07f4

Young-Kyoo Kim·2026년 8월 7일
import io
import sys
import pandas as pd
from minio import Minio


def fetch_inventory_parquet(
    minio_client: Minio, inv_bucket: str, inv_key_or_prefix: str
) -> pd.DataFrame:
    """MinIO AIStor에서 단일 Parquet 파일 또는 디렉토리(Prefix) 하위의 모든 Parquet을 읽어 하나의 DataFrame으로 병합"""
    print(
        f"[1/4] Fetching Inventory Parquet from AIStor: s3://{inv_bucket}/{inv_key_or_prefix}"
    )

    parquet_dfs = []

    try:
        # 1. 전달된 인자가 단일 .parquet 파일인 경우
        if inv_key_or_prefix.endswith('.parquet'):
            target_keys = [inv_key_or_prefix]
        else:
            # 2. 디렉토리/Prefix 경로인 경우 (예: '2026-08-07T07:28Z/files' 또는 '2026-08-07T07:28Z/files/')
            prefix = (
                inv_key_or_prefix
                if inv_key_or_prefix.endswith('/')
                else f'{inv_key_or_prefix}/'
            )
            objects = minio_client.list_objects(
                inv_bucket, prefix=prefix, recursive=True
            )
            target_keys = [
                obj.object_name
                for obj in objects
                if obj.object_name.endswith('.parquet')
            ]

        if not target_keys:
            raise FileNotFoundError(
                f"No .parquet files found under s3://{inv_bucket}/{inv_key_or_prefix}"
            )

        print(
            f' -> Found {len(target_keys)} parquet file(s). Downloading into memory...'
        )

        # 각 Parquet 파일을 메모리로 로드 후 DataFrame 병합
        for key in target_keys:
            response = minio_client.get_object(inv_bucket, key)
            parquet_bytes = response.read()
            response.close()
            response.release_conn()

            df_single = pd.read_parquet(io.BytesIO(parquet_bytes))
            parquet_dfs.append(df_single)

        # 전체 DataFrame 하나로 합치기
        df = pd.concat(parquet_dfs, ignore_index=True)

        # AIStor Inventory 컬럼명 유연화 매핑 (name -> key, e_tag -> etag 등)
        col_rename = {
            'name': 'key',
            'object_name': 'key',
            'e_tag': 'etag',
            'last_modified_date': 'last_modified',
        }
        df = df.rename(columns=col_rename)

        if 'key' not in df.columns or 'size' not in df.columns:
            raise ValueError(
                f'Parquet missing required columns. Found: {list(df.columns)}'
            )

        print(f' -> Successfully loaded total {len(df)} rows from Parquet(s).')
        return df

    except Exception as e:
        print(f'[ERROR] Failed to read Parquet from AIStor: {e}', file=sys.stderr)
        sys.exit(1)

0개의 댓글