MinIO AIStor Inventory Parquet 데이터를 StarRocks object_current_state 테이블에 동기화할 때, Inventory 실행 시점(inventory_snapshot_time) 기준의 삭제 처리(Reconciliation)까지 완벽하게 수행하는 가장 깔끔하고 안전한 전략을 안내해 드립니다.
StarRocks의 Primary Key 모델과 sequence_col (조건부 업데이트) 특성을 활용합니다.
event_time = inventory_snapshot_time, __deleted = 0으로 생성하여 Upsert합니다.inventory_snapshot_time) 이전에 생성/수정되었으나 이번 인벤토리 Parquet에 존재하지 않는 객체는 인벤토리 추출 직전에 삭제된 객체입니다.object_id)를 추출하여 event_time = inventory_snapshot_time, __deleted = 1, size_bytes = 0으로 Upsert합니다.event_time > inventory_snapshot_time)는 StarRocks의 merge_condition에 의해 인벤토리 동기화 작업이 덮어쓰지 않고 자동으로 보존됩니다.아래 스크립트는 Parquet 읽기 인벤토리 내 객체 Stream Load 삭제 대상(Missing Objects) 추출 삭제 처리 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으로 Upsert | DB에서 소프트 삭제 처리되어 COUNT/SUM에서 제외됨 |
| 3. 인벤토리 시작 후 유입된 객체 | DB에 event_time > snapshot_time으로 저장됨 | Stream Load 시 설정된 sequence_col: event_time에 의해 인벤토리 동기화 건이 자동 무시(기존 실시간 데이터 유지)됨 |
| 4. 인벤토리 시작 후 삭제된 객체 | Kafka를 통해 s3:ObjectRemoved 이벤트 유입 | Kafka 이벤트 타임스탬프가 snapshot_time보다 미래이므로 삭제 이벤트가 보존됨 |