MinIO AIStor에서 Parquet 인벤토리를 직접 읽어 정제하고, StarRocks object_current_state 테이블에 활성 객체 Upsert 및 삭제 객체 Reconciliation(소프트 삭제)까지 단일 파이프라인으로 처리하는 전체 Python 스크립트입니다.
aistor_starrocks_reconcile.py#!/usr/bin/env python3
"""MinIO AIStor Inventory & StarRocks Object State Reconciliation Script."""
import argparse
import hashlib
import io
import os
import sys
from urllib.parse import unquote
from minio import Minio
import pandas as pd
import pymysql
import requests
def generate_md5(bucket: str, key: str) -> str:
"""bucket:key 기반 MD5 Primary Key 생성"""
return hashlib.md5(f"{bucket}:{key}".encode("utf-8")).hexdigest()
def fetch_inventory_parquet(
minio_client: Minio, inv_bucket: str, inv_parquet_key: str
) -> pd.DataFrame:
"""MinIO AIStor S3에서 Parquet 파일을 직접 메모리로 읽어 DataFrame 변환"""
print(
f"[1/4] Fetching Inventory Parquet from AIStor: s3://{inv_bucket}/{inv_parquet_key}"
)
try:
response = minio_client.get_object(inv_bucket, inv_parquet_key)
parquet_bytes = response.read()
response.close()
response.release_conn()
df = pd.read_parquet(io.BytesIO(parquet_bytes))
# 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 {len(df)} rows from Parquet.")
return df
except Exception as e:
print(f"[ERROR] Failed to read Parquet from AIStor: {e}", file=sys.stderr)
sys.exit(1)
def stream_load_to_starrocks(
sr_host: str,
sr_http_port: int,
db_name: str,
table_name: str,
user: str,
password: str,
df: pd.DataFrame,
) -> bool:
"""DataFrame을 StarRocks Stream Load API(TSV)를 통해 고속 적재"""
if df.empty:
print(" -> No rows to load.")
return True
# 컬럼 순서 고정 및 TSV 변환
target_cols = [
"object_id",
"bucket",
"object_key",
"prefix",
"size_bytes",
"etag",
"event_time",
"__deleted",
]
tsv_data = df[target_cols].to_csv(index=False, header=False, sep="\t")
url = f"http://{sr_host}:{sr_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",
}
try:
res = requests.put(
url,
data=tsv_data.encode("utf-8"),
headers=headers,
auth=(user, password),
timeout=300,
)
res_json = res.json()
if res_json.get("Status") == "Success":
loaded_cnt = res_json.get("NumberLoadedRows", 0)
print(
f" -> [Stream Load SUCCESS] Loaded Rows: {loaded_cnt}"
)
return True
else:
print(f"[ERROR] Stream Load Failed: {res_json}", file=sys.stderr)
return False
except Exception as e:
print(
f"[ERROR] Stream Load Request Exception: {e}", file=sys.stderr
)
return False
def run_reconciliation(args):
# MinIO Client 생성
minio_client = Minio(
endpoint=args.aistor_endpoint.replace("http://", "").replace(
"https://", ""
),
access_key=args.aistor_access_key,
secret_key=args.aistor_secret_key,
secure=args.aistor_secure,
)
# 1. Parquet 파일 로드
df_raw = fetch_inventory_parquet(
minio_client, args.inventory_bucket, args.inventory_key
)
# 2. 데이터 가공 (URL 디코딩, Primary Key 및 Prefix 생성)
print(
f"[2/4] Processing Inventory Data (Snapshot Time: {args.snapshot_time})..."
)
df_inv = df_raw.copy()
df_inv["object_key"] = df_inv["key"].apply(lambda k: unquote(str(k)))
df_inv["bucket"] = args.target_bucket
df_inv["object_id"] = df_inv.apply(
lambda r: generate_md5(r["bucket"], r["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"].fillna(0).astype(int)
df_inv["etag"] = df_inv["etag"].fillna("") if "etag" in df_inv.columns else ""
df_inv["event_time"] = args.snapshot_time
df_inv["__deleted"] = 0
# 3. 인벤토리 내 활성 객체 Upsert
print(
f"[3/4] Stream Loading {len(df_inv)} active objects into StarRocks..."
)
success = stream_load_to_starrocks(
args.sr_host,
args.sr_http_port,
args.db_name,
args.table_name,
args.sr_user,
args.sr_password,
df_inv,
)
if not success:
sys.exit(1)
# 4. 삭제 대상 객체 Reconciliation (Missing Objects 처리)
print(
"[4/4] Querying existing StarRocks objects for Reconciliation..."
)
conn = pymysql.connect(
host=args.sr_host,
port=args.sr_mysql_port,
user=args.sr_user,
password=args.sr_password,
db=args.db_name,
autocommit=True,
)
try:
# 인벤토리 시작 시점(snapshot_time) 이전에 유입된 활성 객체 조회
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, (args.target_bucket, args.snapshot_time))
db_rows = cursor.fetchall()
inv_object_ids = set(df_inv["object_id"])
deleted_records = [
r for r in db_rows if r["object_id"] not in inv_object_ids
]
if deleted_records:
print(
f" -> Detected {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"] = args.snapshot_time
df_del["__deleted"] = 1
del_success = stream_load_to_starrocks(
args.sr_host,
args.sr_http_port,
args.db_name,
args.table_name,
args.sr_user,
args.sr_password,
df_del,
)
if not del_success:
sys.exit(1)
else:
print(" -> No missing/deleted objects found.")
finally:
conn.close()
print("\n=== Inventory Reconciliation Completed Successfully ===")
def main():
parser = argparse.ArgumentParser(
description="MinIO AIStor Inventory Parquet -> StarRocks Reconciliation Pipeline"
)
# Required Arguments
parser.add_argument(
"--target-bucket",
required=True,
help="MinIO Bucket Name being reconciled",
)
parser.add_argument(
"--inventory-bucket",
required=True,
help="AIStor Bucket containing the inventory Parquet file",
)
parser.add_argument(
"--inventory-key",
required=True,
help="S3 Key/Path to the inventory Parquet file",
)
parser.add_argument(
"--snapshot-time",
required=True,
help="Inventory execution start timestamp (Format: 'YYYY-MM-DD HH:MM:SS')",
)
# MinIO AIStor Options
parser.add_argument(
"--aistor-endpoint",
default="http://localhost:9000",
help="AIStor S3 Endpoint",
)
parser.add_argument(
"--aistor-access-key", default="minioadmin", help="AIStor Access Key"
)
parser.add_argument(
"--aistor-secret-key", default="minioadmin", help="AIStor Secret Key"
)
parser.add_argument(
"--aistor-secure",
action="store_true",
help="Use HTTPS for AIStor connection",
)
# StarRocks Options
parser.add_argument(
"--sr-host", default="localhost", help="StarRocks FE Host"
)
parser.add_argument(
"--sr-http-port",
type=int,
default=8030,
help="StarRocks HTTP Port (Stream Load)",
)
parser.add_argument(
"--sr-mysql-port",
type=int,
default=9030,
help="StarRocks MySQL Port",
)
parser.add_argument(
"--sr-user", default="root", help="StarRocks DB User"
)
parser.add_argument(
"--sr-password", default="", help="StarRocks DB Password"
)
parser.add_argument(
"--db-name", default="minio_metadata", help="StarRocks Database Name"
)
parser.add_argument(
"--table-name",
default="object_current_state",
help="StarRocks Table Name",
)
args = parser.parse_args()
run_reconciliation(args)
if __name__ == "__main__":
main()
python3 aistor_starrocks_reconcile.py \
--target-bucket my-data-bucket \
--inventory-bucket system-inventory-bucket \
--inventory-key 2026-08-07/my-data-bucket-inventory.parquet \
--snapshot-time "2026-08-07 12:00:00" \
--aistor-endpoint http://192.168.1.50:9000 \
--aistor-access-key YOUR_ACCESS_KEY \
--aistor-secret-key YOUR_SECRET_KEY \
--sr-host 192.168.1.100 \
--sr-http-port 8030 \
--sr-mysql-port 9030 \
--sr-user root