26Y04a

Young-Kyoo Kim·2026년 8월 4일

MinIO AIStor(또는 Enterprise MinIO)의 이벤트 스트리밍 기반 실시간 집계와 Inventory 기반의 주기적 Reconciliation(정산/동기화) 구조에서 발생하는 동시성(Concurrency) 및 일관성(Consistency) 문제를 정확하게 짚으셨습니다.

질문하신 두 가지 핵심 사항에 대해 해결 방안과 아키텍처를 제시해 드립니다.


1. Inventory 파일의 StarRocks Upsert 방법

Inventory API가 생성하는 리포트 파일(일반적으로 Parquet, CSV, ORC 형태)을 StarRocks로 수집하여 기존 데이터와 Reconciliation하는 과정은 다음과 같이 설계하는 것이 가장 효율적입니다.

① Primary Key 및 데이터 모델 설정

StarRocks에서 특정 Prefix/Bucket 단위의 최신 상태를 유지하려면 Primary Key 테이블 모델을 사용해야 합니다.

CREATE TABLE minio_object_inventory (
    bucket VARCHAR(255),
    object_key VARCHAR(1024),
    size_bytes BIGINT,
    last_modified DATETIME,
    inventory_timestamp DATETIME -- Inventory가 생성된 snapshot 시점
) ENGINE=OLAP
PRIMARY KEY(bucket, object_key)
DISTRIBUTED BY HASH(bucket, object_key) BUCKETS 32;

② 수집 (Ingestion) 방식: Broker Load 또는 S3 TVF

StarRocks는 MinIO(S3 호환)에 저장된 Inventory 파일을 직접 읽어 빠르게 Upsert할 수 있습니다.

  • S3 Table-Value Function (TVF) 사용 (추천):
    별도의 ETL 커넥터 없이 SQL 쿼리로 직접 MinIO Inventory Parquet/CSV 파일을 읽어 Upsert합니다.
INSERT INTO minio_object_inventory
SELECT 
    bucket,
    key AS object_key,
    size AS size_bytes,
    last_modified,
    '2026-08-04 00:00:00' AS inventory_timestamp
FROM S3(
    "uri" = "s3://minio-inventory-bucket/2026-08-04/inventory.parquet",
    "s3.endpoint" = "http://minio.internal:9000",
    "s3.access_key" = "<ACCESS_KEY>",
    "s3.secret_key" = "<SECRET_KEY>",
    "format" = "parquet"
);

③ Prefix별 사용량/Quota Aggregation View 생성

객체 단위 수집 테이블 위에 Prefix 단위 집계 뷰를 생성하거나, 별도의 Rollup 테이블(Primary Key 모델)을 구성합니다.

-- Prefix 추출 규칙예시 (depth 2 기준)
CREATE VIEW v_prefix_usage AS
SELECT 
    bucket,
    SPLIT_PART(object_key, '/', 1) AS prefix,
    SUM(size_bytes) AS current_usage_bytes,
    COUNT(1) AS object_count
FROM minio_object_inventory
GROUP BY bucket, SPLIT_PART(object_key, '/', 1);

2. Inventory 생성 중(스냅샷 기간) Update 발생 시 누락/덮어쓰기 방지 방안

MinIO Inventory API는 특정 시점(Snapshot Timestamp TsnapT_{snap})의 객체 상태를 덤프합니다.

Inventory 수집 작업(10분~60분 소요) 도중에도 Live Event(Kafka -> Routine Load)는 계속 들어오기 때문에, 오래된 Inventory 데이터가 최신 Kafka 이벤트를 덮어씌우는 문제(Dirty Overwrite)가 발생할 수 있습니다.

이를 보정하기 위한 3가지 핵심 전략입니다.


전략 A. Event Timestamp / Sequence 기반 Upsert 조건 제어 (가장 표준적)

StarRocks의 Primary Key 테이블 모델은 ORDER BY 컬럼 조건 기반의 Conditional Update(또는 Delete Vector)를 지원합니다.

즉, '더 최신의 이벤트 시점'을 가진 데이터만 덮어쓰도록 설정합니다.

  1. 테이블 정의 시 Sequence Column 지정:
    StarRocks Primary Key 테이블 생성 시 sequence_col 옵션을 객체의 last_modified (또는 Kafka event time / inventory snapshot time)으로 설정합니다.
CREATE TABLE minio_object_inventory (
    bucket VARCHAR(255),
    object_key VARCHAR(1024),
    size_bytes BIGINT,
    last_modified DATETIME,
    is_deleted BOOLEAN DEFAULT "false"
) ENGINE=OLAP
PRIMARY KEY(bucket, object_key)
PROPERTIES (
    "primary_key_encoding" = "BINARY",
    "sequence_col" = "last_modified" -- 핵심: last_modified가 더 높은(최신) 레코드만 업데이트 승인
);
  1. 동작 원리:
  • Inventory 데이터가 작성되는 중(예: 01:00 Snapshot) 01:15에 새로운 객체 생성/삭제 이벤트가 Kafka -> Routine Load를 통해 last_modified = 01:15로 들어옴.
  • 01:30에 Inventory 덤프 완료 후 StarRocks로 last_modified = 01:00 인 Inventory 레코드를 Upsert 시도.
  • StarRocks는 이미 last_modified = 01:15인 데이터가 존재함을 확인하고, 01:00의 Inventory 레코드를 무시(Ignore)함.

전략 B. Two-Tier Staging & Snapshot Cutoff Reconciliation (배치 보정 레이어)

만약 Kafka 이벤트에 삭제(ObjectRemoved) 이벤트가 누락되어 Inventory로 전체 동기화를 해야 하는 경우, Live 테이블을 바로 덮어쓰지 않고 Staging + Delta Apply 방식을 사용합니다.

[MinIO Inventory] ──(T_snap 시점 Snapshot)──> [Staging Table]
                                                      │
[Kafka Event]      ──(Routine Load)──────────> [Live Table]
                                                      │
                                           [Reconciliation SQL Query]
                                                      ↓
                                           [Final Master State]

처리 절차

  1. Inventory 수집: Inventory 파일을 staging_inventory 테이블에 로드합니다. 이때 Inventory의 Snapshot 생성 시작 시점인 TsnapT_{snap} (예: 2026-08-04 01:00:00)을 함께 기록합니다.
  2. Reconciliation Target 결정: Live 테이블에서 last_modified < T_snap 인 데이터만 Staging 데이터로 대체합니다.
  3. Delete Detection (선택 사항): Staging에는 없는데 Live 테이블(last_modified < T_snap)에는 존재하는 객체는 Kafka 삭제 이벤트를 놓친 것이므로 Delete 처리합니다.
-- 1. Snapshot 시점(T_snap) 이전의 데이터 중 Inventory와 차이나는 부분 보정/삭제
DELETE FROM live_object_state 
WHERE last_modified < '2026-08-04 01:00:00'
  AND (bucket, object_key) NOT IN (
      SELECT bucket, object_key FROM staging_inventory
  );

-- 2. Snapshot 시점(T_snap) 이전 기준 Inventory 최신 데이터 반영
INSERT INTO live_object_state
SELECT bucket, object_key, size_bytes, last_modified, false
FROM staging_inventory
WHERE last_modified < '2026-08-04 01:00:00';

전략 C. Routine Load의 Kafka Offset replay 활용 (Event Sourcing 접근법)

Inventory가 실행되는 동안 실시간 이벤트를 계속 처리하되, Reconciliation 시점에 Kafka Offset을 정산 시점으로 감아 재처리를 방지하는 방법입니다.

  1. Inventory Snapshot 생성 시점 TsnapT_{snap}의 Kafka Offset 저장 (또는 Timestamp 기준 Offset 조회).
  2. Inventory 전체 데이터를 StarRocks에 Upsert (해당 시점 상태로 정률화).
  3. Routine Load의 소비 위치를 TsnapT_{snap} 시점의 Kafka Offset으로 Pause -> Alter Offset -> Resume 하여 TsnapT_{snap} 이후 발생한 이벤트를 다시 재생(Replay)하여 덧씌움.

💡 요약 및 권장 가이드

고려 사항추천 방식
StarRocks 수집S3 Table-Value Function (S3())으로 Parquet 직접 인제스처
StarRocks 모델PRIMARY KEY 모델 (sequence_col = last_modified)
동시성/누락 보정전략 A (Sequence Column) 적용으로 과거 snapshot의 덮어쓰기를 DB 차원에서 차단.


추가로 삭제 이벤트 누락 검출이 필요할 경우 전략 B (Staging Delete Sweep) 병행. |

이 구조를 적용하면 Kafka 실시간 유실을 Inventory로 주기적 보정하면서도, Reconciliation 과정 중 발생하는 실시간 변경 사항이 유실되는 현상을 완벽히 방지할 수 있습니다.

0개의 댓글