26Z07a

Young-Kyoo Kim·2026년 8월 7일

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 스크립트입니다.

  • 동작 원리: Inventory 생성 시점(snapshot_time)을 레코드의 event_time으로 지정하여 StarRocks Stream Load로 Upsert합니다.
  • 이벤트 보호 Mechanism: Inventory가 배치로 로딩되는 동안 Kafka로 유입된 실시간 PUT/DELETE 이벤트는 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';

0개의 댓글