MinIO AIStor의 PUT/DELETE 실시간 이벤트와 Daily Inventory API 스냅샷을 StarRocks로 통합할 때 핵심은 실시간 이벤트 처리 중 Daily Inventory가 덮어쓰는 레이스 조건(Race Condition)을 방지하는 것입니다.
StarRocks의 Primary Key Table Sequence Column (ORDER BY event_time) 기능을 활용하면, 이벤트 생성 시점(event_time)을 기준으로 이전 시점의 Inventory 데이터가 들어오더라도 최신 실시간 이벤트를 자동으로 보존할 수 있습니다.
1. Kafka Topic 및 MinIO Notification 설정
Kafka Broker에 이벤트 수신용 Topic을 생성합니다.
# Kafka Topic 생성
kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic minio-bucket-events \
--partitions 3 --replication-factor 1
MinIO AIStor Dashboard 또는 mc CLI를 통해 PUT/DELETE 이벤트가 위 Kafka topic (minio-bucket-events)으로 JSON 형식 발행되도록 설정합니다.
2. StarRocks DB, Table 및 Routine Load DDL
StarRocks에 데이터베이스와 Primary Key 모델 테이블을 생성합니다. ORDER BY event_time을 설정하여 이전 타임스탬프의 무효한 Upsert를 엔진 레벨에서 자동 차단합니다.
CREATE DATABASE IF NOT EXISTS minio_metadata;
USE minio_metadata;
-- 1) 오브젝트 현황 관리 테이블 (Primary Key Model)
CREATE TABLE IF NOT EXISTS object_current_state (
bucket VARCHAR(128) NOT NULL,
object_key VARCHAR(1024) NOT NULL,
prefix VARCHAR(512) NOT NULL,
size_bytes BIGINT DEFAULT "0",
etag VARCHAR(64) DEFAULT "",
event_time DATETIME NOT NULL,
__deleted TINYINT DEFAULT "0" -- 1일 경우 삭제된 객체 처리
)
PRIMARY KEY (bucket, object_key)
ORDER BY event_time -- 핵심: event_time이 더 오래된 레코드의 Upsert는 자동으로 무시됨
DISTRIBUTED BY HASH(bucket, object_key) BUCKETS 8
PROPERTIES (
"replication_num" = "1",
"enable_persistent_index" = "true"
);
-- 2) Kafka Event 실시간 수집 (Routine Load)
CREATE ROUTINE LOAD minio_event_load ON object_current_state
COLUMNS (
event_time_str,
event_name,
bucket,
object_key,
size_bytes,
etag,
event_time = str_to_date(event_time_str, '%Y-%m-%dT%H:%i:%sZ'),
prefix = if(locate('/', object_key) > 0, regexp_extract(object_key, '^(.*)/[^/]+$', 1), ''),
__deleted = if(event_name LIKE 's3:ObjectRemoved%', 1, 0)
)
PROPERTIES (
"format" = "json",
"jsonpaths" = "[\"$.eventTime\", \"$.eventName\", \"$.s3.bucket.name\", \"$.s3.object.key\", \"$.s3.object.size\", \"$.s3.object.eTag\"]",
"strip_outer_array" = "true"
)
FROM KAFKA (
"kafka_broker_list" = "localhost:9092",
"kafka_topic" = "minio-bucket-events",
"property.group.id" = "starrocks_minio_events",
"property.kafka_default_offsets" = "OFFSET_BEGINNING"
);
-- 3) Prefix별 실시간 집계 뷰 (Instant Query용)
CREATE VIEW IF NOT EXISTS v_prefix_summary AS
SELECT
bucket,
prefix,
COUNT(1) AS total_object_count,
SUM(size_bytes) AS total_size_bytes
FROM object_current_state
GROUP BY bucket, prefix;
3. Daily Inventory Upsert & 레이스 조건 방지 파이프라인
Inventory API를 통해 추출된 스냅샷(Parquet 또는 CSV)을 StarRocks로 동기화하는 Python 스크립트입니다.
snapshot_time)을 레코드의 event_time으로 지정하여 StarRocks Stream Load로 Upsert합니다.event_time > snapshot_time이므로, StarRocks가 Inventory 데이터를 로딩할 때 해당 키의 최신 이벤트를 과거 스냅샷 데이터로 덮어쓰지 않고 자동 보호합니다.import os
import requests
import pandas as pd
from datetime import datetime
# 1. 설정 정보
STARROCKS_FE_HOST = "localhost"
STARROCKS_HTTP_PORT = "8030"
DB_NAME = "minio_metadata"
TABLE_NAME = "object_current_state"
USER = "root"
PASSWORD = ""
def run_daily_inventory_upsert(inventory_parquet_path: str, snapshot_time_str: str):
"""
snapshot_time_str 예시: '2026-08-07 00:00:00'
MinIO Inventory 결과 파일(Parquet)을 읽어 StarRocks에 안전하게 Stream Load
"""
# Inventory 파일 읽기 (bucket, key, size, etag 컬럼 존재 가정)
df = pd.read_parquet(inventory_parquet_path)
# Prefix 추출 및 snapshot_time 부여
df['prefix'] = df['key'].apply(lambda k: os.path.dirname(k) if '/' in k else '')
df['event_time'] = snapshot_time_str
df['__deleted'] = 0
# CSV 변환 (Stream Load용)
csv_data = df[['bucket', 'key', 'prefix', 'size', '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": "bucket, object_key, prefix, size_bytes, etag, event_time, __deleted",
"Expect": "100-continue"
}
# Stream Load 실행
response = requests.put(
url,
data=csv_data.encode('utf-8'),
headers=headers,
auth=(USER, PASSWORD)
)
print("Stream Load Status Code:", response.status_code)
print("Response:", response.json())
if __name__ == "__main__":
# Daily 배치 실행 예시
# Snapshot 시점을 정확히 맞춰 지정
SNAPSHOT_TIME = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
run_daily_inventory_upsert("inventory_20260807.parquet", SNAPSHOT_TIME)
4. Prefix 조회 Query 사용법
StarRocks Primary Key 테이블 기반으로 구축되었기 때문에, 백만 개 이상의 객체가 존재하더라도 sub-second 수준으로 Prefix별 수량 및 용량을 즉시 응답할 수 있습니다.
-- 특정 Bucket 및 Prefix 하위 전체 집계 조회
SELECT
bucket,
'images/2026/' AS prefix,
COUNT(1) AS total_objects,
SUM(size_bytes) AS total_bytes
FROM object_current_state
WHERE bucket = 'my-aistor-bucket'
AND object_key LIKE 'images/2026/%';
-- Pre-aggregated View 조회
SELECT * FROM v_prefix_summary
WHERE bucket = 'my-aistor-bucket' AND prefix = 'images/2026';