26Z10c4

Young-Kyoo Kim·5일 전

StarRocks에서 generate_seriesunnest 같은 테이블 함수(Table Function/UDTF)를 transform() 같은 단일 스칼라 표현식 내부에 중첩해 사용하려 했기 때문에 발생한 에러입니다.

StarRocks의 Async Materialized View에서는 테이블 함수(UDTF)의 표현식 내 중첩을 허용하지 않으므로, 숫자 시퀀스 CTE(Numbers Table)와 표준 배열 스칼라 함수(split, slice, array_to_string)의 CROSS JOIN 조합으로 변경해야 합니다.


1. 수정된 Async Materialized View DDL

아래 쿼리는 UDTF 대신 표준 배열 함수만 사용하므로 구문 오류 없이 정상 생성됩니다.

CREATE MATERIALIZED VIEW mv_prefix_summary
COMMENT "S3 Prefix Level Recursive Summary MV"
DISTRIBUTED BY HASH(bucket, prefix)
REFRESH ASYNC EVERY (INTERVAL 5 MINUTE)
AS
WITH numbers AS (
    -- 최대 10단계 깊이의 Prefix 지원 (필요 시 UNION ALL로 확장 가능)
    SELECT 1 AS pos UNION ALL SELECT 2 UNION ALL SELECT 3 UNION ALL 
    SELECT 4 UNION ALL SELECT 5 UNION ALL SELECT 6 UNION ALL 
    SELECT 7 UNION ALL SELECT 8 UNION ALL SELECT 9 UNION ALL SELECT 10
),
recursive_prefix AS (
    SELECT 
        o.bucket,
        -- split으로 분할된 배열에서 1부터 pos까지 자른 뒤 다시 '/'로 결합
        array_to_string(slice(split(o.prefix, '/'), 1, n.pos), '/') AS prefix,
        o.size_bytes
    FROM object_current_state o
    CROSS JOIN numbers n
    WHERE o.__deleted = 0 
      AND o.prefix IS NOT NULL 
      AND o.prefix != ''
      AND n.pos <= array_length(split(o.prefix, '/'))
)
SELECT 
    bucket,
    prefix,
    COUNT(1) AS total_object_count,
    SUM(size_bytes) AS total_size_bytes
FROM recursive_prefix
GROUP BY bucket, prefix;

2. 작동 원리 (예시: prefix = 'a/b/c')

  1. split('a/b/c', '/') \rightarrow ['a', 'b', 'c'] (배열 길이: 3)
  2. numbers CTE와 조인 시 pos <= 3 조건에 의해 pos 1, 2, 3행 생성
  3. pos = 1: slice(..., 1, 1) \rightarrow ['a'] \rightarrow 'a'
  4. pos = 2: slice(..., 1, 2) \rightarrow ['a', 'b'] \rightarrow 'a/b'
  5. pos = 3: slice(..., 1, 3) \rightarrow ['a', 'b', 'c'] \rightarrow 'a/b/c'

3. 대안: 물리 테이블 + 파이프라인 배치 갱신 (추천)

만약 Async MV의 백그라운드 리프레시 부하를 완전히 없애고 Reconciliation 배치 완료 직후에만 즉시 갱신하고 싶다면, 물리 Aggregate 테이블을 만들고 파이썬 스크립트 끝단에서 INSERT OVERWRITE를 날리는 방식이 더 깔끔하고 확실합니다.

테이블 생성 (Aggregate Key 모델):

CREATE TABLE prefix_summary (
    bucket             VARCHAR(64) NOT NULL,
    prefix             VARCHAR(1024) NOT NULL,
    total_object_count BIGINT SUM DEFAULT "0",
    total_size_bytes   BIGINT SUM DEFAULT "0"
) ENGINE=OLAP
AGGREGATE KEY(bucket, prefix)
DISTRIBUTED BY HASH(bucket, prefix);

Python 스크립트 맨 뒤에 추가할 갱신 쿼리:

# Reconciliation 완료 후 집계 테이블 일괄 갱신 (1~2초 소요)
refresh_sql = """
    INSERT OVERWRITE prefix_summary
    WITH numbers AS (
        SELECT 1 AS pos UNION ALL SELECT 2 UNION ALL SELECT 3 UNION ALL 
        SELECT 4 UNION ALL SELECT 5 UNION ALL SELECT 6 UNION ALL 
        SELECT 7 UNION ALL SELECT 8 UNION ALL SELECT 9 UNION ALL SELECT 10
    )
    SELECT 
        o.bucket,
        array_to_string(slice(split(o.prefix, '/'), 1, n.pos), '/') AS prefix,
        COUNT(1) AS total_object_count,
        SUM(o.size_bytes) AS total_size_bytes
    FROM object_current_state o
    CROSS JOIN numbers n
    WHERE o.__deleted = 0 
      AND o.prefix IS NOT NULL AND o.prefix != ''
      AND n.pos <= array_length(split(o.prefix, '/'))
    GROUP BY o.bucket, prefix;
"""
cursor.execute(refresh_sql)

===

물리 Aggregate 테이블(prefix_summary)을 하루 한 번 Reconciliation 시점에 INSERT OVERWRITE로 갱신하는 구조의 조회 속도, 데이터 신선도(Freshness), 그리고 트레이드오프는 다음과 같습니다.


1. 조회 속도: 극도로 빠름 (1~10ms 수준)

평상시 외부에서 prefix_summary 테이블을 조회할 때의 속도는 실시간에 가까운 초고속 응답을 제공합니다.

  • 이유:
  • 데이터 크기 축소 (Data Aggregation): 원천 테이블(object_current_state)이 1억 개의 객체 데이터를 갖고 있더라도, Prefix 단위로 그룹핑된 prefix_summary 테이블은 수천~수만 건에 불과합니다.
  • Direct Scan: Complex 계산(Array Unnest, Join, Group By)을 이미 배치가 다 끝내놓은 상태이므로, 단순 SELECT ... WHERE bucket = ... AND prefix = ... 쿼리는 인덱스를 타고 디스크/메모리에서 몇 개 행만 바로 읽어 반환합니다.
  • DB 부하 제로: 동시 조회 요청이 초당 수천 건(High TPS)이 들어와도 StarRocks FE/BE에 아무런 계산 부담을 주지 않습니다.

2. 고려해야 할 핵심 트레이드오프: 데이터 신선도 (Data Lag)

조회 속도는 최고 수준이지만, 데이터의 실시간성(Accuracy) 측면에서 차이가 발생합니다.

  • 상황:
  • 하루 동안 Kafka를 통해 신규 파일 업로드(s3:ObjectCreated)나 삭제(s3:ObjectRemoved) 이벤트가 지속적으로 유입됩니다.
  • object_current_state (Primary Key 테이블)에는 실시간으로 반영되어 최신 상태를 유지합니다.
  • 하지만 prefix_summary는 Daily 배치 전까지 어제 Reconciliation 시점의 스냅샷 데이터를 그대로 유지합니다. (최대 24시간의 시차 발생)

3. 요구사항별 대응 전략

사용하는 서비스의 목적에 따라 아래 2가지 방향 중 선택할 수 있습니다.

방향 A. 24시간 시차 허용 가능 (대시보드 / 용량 트렌드 / 일별 리포트)

  • 방식: 질문하신 현 구조 그대로 유지 (INSERT OVERWRITE 1회 실행)
  • 장점: 아키텍처가 가장 단순하며, 조회 성능 및 안정성이 100% 보장됨.

방향 B. 실시간 정확도 필수 (현재 시점의 정확한 Prefix 용량/객체수 조회)

Reconciliation 사이에 들어온 Kafka 변동분까지 합산하여 실시간으로 보여주고 싶다면, 하이브리드 View를 만들어 외부에 제공하면 됩니다.

-- 실시간 하이브리드 View (기존 배치 집계 + 오늘자 Kafka 변동분 합산)
CREATE VIEW v_realtime_prefix_summary AS
WITH today_changes AS (
    -- 오늘 snapshot_time 이후 변경된 객체들만 추출하여 Prefix 계산
    SELECT 
        o.bucket,
        array_to_string(slice(split(o.prefix, '/'), 1, n.pos), '/') AS prefix,
        -- 삭제된 건은 -1 / -size, 신규/갱신 건은 +1 / +size로 계산
        SUM(IF(o.__deleted = 1, -1, 1)) AS delta_count,
        SUM(IF(o.__deleted = 1, -o.size_bytes, o.size_bytes)) AS delta_size
    FROM object_current_state o
    CROSS JOIN (
        SELECT 1 AS pos UNION ALL SELECT 2 UNION ALL SELECT 3 UNION ALL 
        SELECT 4 UNION ALL SELECT 5 UNION ALL SELECT 6 UNION ALL 
        SELECT 7 UNION ALL SELECT 8 UNION ALL SELECT 9 UNION ALL SELECT 10
    ) n
    WHERE o.event_time > (SELECT MAX(event_time) FROM prefix_summary) -- 마지막 배치 시점 이후 데이터만 Filter
      AND n.pos <= array_length(split(o.prefix, '/'))
    GROUP BY o.bucket, prefix
)
SELECT 
    COALESCE(base.bucket, chg.bucket) AS bucket,
    COALESCE(base.prefix, chg.prefix) AS prefix,
    COALESCE(base.total_object_count, 0) + COALESCE(chg.delta_count, 0) AS total_object_count,
    COALESCE(base.total_size_bytes, 0) + COALESCE(chg.delta_size, 0) AS total_size_bytes
FROM prefix_summary base
FULL OUTER JOIN today_changes chg 
  ON base.bucket = chg.bucket AND base.prefix = chg.prefix;
  • 동작 특징: 99%의 과거 데이터는 pre-computed된 prefix_summary에서 즉시 읽고, 최근 몇 시간 동안 변경된 소량의 Kafka 데이터만 즉석 계산하여 합치므로 조회 속도(수십 ms)와 실시간성을 모두 충족시킵니다.

0개의 댓글