26Z07f2

Young-Kyoo Kim·2026년 8월 7일

MinIO AIStor Inventory Parquet 데이터를 StarRocks object_current_state 테이블에 동기화할 때, Inventory 실행 시점(inventory_snapshot_time) 기준의 삭제 처리(Reconciliation)까지 완벽하게 수행하는 가장 깔끔하고 안전한 전략을 안내해 드립니다.


핵심 구현 아이디어

StarRocks의 Primary Key 모델sequence_col (조건부 업데이트) 특성을 활용합니다.

  1. 인벤토리 대상 추출 및 Upsert
  • 인벤토리 Parquet에 포함된 객체들은 event_time = inventory_snapshot_time, __deleted = 0으로 생성하여 Upsert합니다.
  1. 삭제된 객체 추적 및 소프트 삭제(Delete Flag) Upsert
  • 인벤토리 실행 시점(inventory_snapshot_time) 이전에 생성/수정되었으나 이번 인벤토리 Parquet에 존재하지 않는 객체는 인벤토리 추출 직전에 삭제된 객체입니다.
  • 이 객체들의 Primary Key(object_id)를 추출하여 event_time = inventory_snapshot_time, __deleted = 1, size_bytes = 0으로 Upsert합니다.
  1. 인벤토리 실행 중 유입된 신규 객체 보호
  • 인벤토리 실행 시점 이후에 유입된 실시간 Kafka 이벤트(event_time > inventory_snapshot_time)는 StarRocks의 merge_condition에 의해 인벤토리 동기화 작업이 덮어쓰지 않고 자동으로 보존됩니다.

1. Python Reconciliation 전체 코드

아래 스크립트는 Parquet 읽기 \rightarrow 인벤토리 내 객체 Stream Load \rightarrow 삭제 대상(Missing Objects) 추출 \rightarrow 삭제 처리 Stream Load 단계를 일괄 수행합니다.

import hashlib
import os
import sys
import requests
import pandas as pd
from urllib.parse import unquote

# Configuration
STARROCKS_FE_HOST = "localhost"
STARROCKS_HTTP_PORT = 8030
STARROCKS_MYSQL_PORT = 9030
DB_NAME = "minio_metadata"
TABLE_NAME = "object_current_state"
USER = "root"
PASSWORD = ""

def generate_md5(bucket: str, key: str) -> str:
    return hashlib.md5(f"{bucket}:{key}".encode('utf-8')).hexdigest()

def stream_load_to_starrocks(df: pd.DataFrame, snapshot_time_str: str) -> bool:
    """Dataframe을 StarRocks Stream Load HTTP API로 적재"""
    if df.empty:
        print("[INFO] No data to load.")
        return True

    # TSV 파싱 데이터 준비
    csv_data = df[['object_id', 'bucket', 'object_key', 'prefix', 'size_bytes', 'etag', 'event_time', '__deleted']].to_csv(
        index=False, header=False, sep='\t'
    )

    url = f"http://{STARROCKS_FE_HOST}:{STARROCKS_HTTP_PORT}/api/{DB_NAME}/{TABLE_NAME}/_stream_load"
    
    headers = {
        "format": "csv",
        "column_separator": "\t",
        "columns": "object_id, bucket, object_key, prefix, size_bytes, etag, event_time, __deleted",
        "sequence_col": "event_time", # 타임스탬프 기반 최신 데이터만 보존
        "Expect": "100-continue"
    }

    auth = (USER, PASSWORD)
    response = requests.put(url, data=csv_data.encode('utf-8'), headers=headers, auth=auth)
    res_json = response.json()
    
    if res_json.get("Status") == "Success":
        print(f"[SUCCESS] Stream Load Completed. Loaded Rows: {res_json.get('NumberLoadedRows')}")
        return True
    else:
        print(f"[ERROR] Stream Load Failed: {res_json}")
        return False

def run_inventory_reconciliation(target_bucket: str, inventory_parquet_path: str, snapshot_time_str: str):
    """
    snapshot_time_str 예시: '2026-08-07 12:00:00' (인벤토리 추출 시작 시각)
    """
    import pymysql
    
    print(f"=== Starting Inventory Reconciliation for Bucket: {target_bucket} ===")
    print(f"Snapshot Timestamp: {snapshot_time_str}")

    # Step 1. Parquet 파일 로드 및 정제
    df_inv = pd.read_parquet(inventory_parquet_path)
    
    # 컬럼명 매핑 및 URL 디코딩 (%2F -> /)
    df_inv['object_key'] = df_inv['key'].apply(lambda k: unquote(str(k)))
    df_inv['bucket'] = target_bucket
    df_inv['object_id'] = df_inv.apply(lambda row: generate_md5(row['bucket'], row['object_key']), axis=1)
    df_inv['prefix'] = df_inv['object_key'].apply(lambda k: os.path.dirname(k) if '/' in k else '')
    df_inv['size_bytes'] = df_inv['size']
    df_inv['event_time'] = snapshot_time_str
    df_inv['__deleted'] = 0

    # Step 2. 인벤토리 존재 객체 Upsert
    print(f"[Step 1/2] Upserting {len(df_inv)} active objects from Inventory...")
    if not stream_load_to_starrocks(df_inv, snapshot_time_str):
        sys.exit(1)

    # Step 3. StarRocks DB에서 삭제 대상(Inventory에 없는 기존 객체) 조회
    print("[Step 2/2] Detecting and marking deleted objects...")
    conn = pymysql.connect(
        host=STARROCKS_FE_HOST,
        port=STARROCKS_MYSQL_PORT,
        user=USER,
        password=PASSWORD,
        db=DB_NAME,
        autocommit=True
    )

    try:
        # DB에서 snapshot_time 이전에 존재했고, 삭제되지 않은 객체들의 ID 추출
        sql = """
            SELECT object_id, bucket, object_key, prefix
            FROM object_current_state
            WHERE bucket = %s
              AND event_time <= %s
              AND __deleted = 0
        """
        with conn.cursor(pymysql.cursors.DictCursor) as cursor:
            cursor.execute(sql, (target_bucket, snapshot_time_str))
            db_rows = cursor.fetchall()

        # Parquet의 object_id Set과 비교하여 삭제된 ID 도출
        inv_object_ids = set(df_inv['object_id'])
        deleted_records = [row for row in db_rows if row['object_id'] not in inv_object_ids]

        if deleted_records:
            print(f"[INFO] Found {len(deleted_records)} deleted objects. Applying soft-delete...")
            df_del = pd.DataFrame(deleted_records)
            df_del['size_bytes'] = 0
            df_del['etag'] = ''
            df_del['event_time'] = snapshot_time_str
            df_del['__deleted'] = 1

            if not stream_load_to_starrocks(df_del, snapshot_time_str):
                sys.exit(1)
        else:
            print("[INFO] No deleted objects found.")

    finally:
        conn.close()

    print("=== Inventory Reconciliation Successfully Finished ===")

if __name__ == "__main__":
    # 예시 실행
    run_inventory_reconciliation(
        target_bucket="my-bucket",
        inventory_parquet_path="./inventory_snapshot.parquet",
        snapshot_time_str="2026-08-07 12:00:00" # MinIO AIStor 인벤토리 생성 시작 시각
    )

동작 시나리오 및 안전성 검증

케이스처리 방식이유 및 결과
1. 인벤토리 내 기존 객체__deleted = 0, event_time = snapshot_time으로 Upsert최신 상태로 갱신됨
2. 인벤토리 시작 전 삭제된 객체__deleted = 1, event_time = snapshot_time으로 UpsertDB에서 소프트 삭제 처리되어 COUNT/SUM에서 제외됨
3. 인벤토리 시작 후 유입된 객체DB에 event_time > snapshot_time으로 저장됨Stream Load 시 설정된 sequence_col: event_time에 의해 인벤토리 동기화 건이 자동 무시(기존 실시간 데이터 유지)
4. 인벤토리 시작 후 삭제된 객체Kafka를 통해 s3:ObjectRemoved 이벤트 유입Kafka 이벤트 타임스탬프가 snapshot_time보다 미래이므로 삭제 이벤트가 보존됨

0개의 댓글