SQL 기반 인제스트 개념

SQL 기반 인제스트 개념 (SQL-based ingestion concepts)

MSQ(멀티 스테이지 쿼리) 태스크 엔진으로 SQL 기반 배치 인제스트를 수행하는 개념을 설명하는 문서예요. EXTERN·INSERT·REPLACE 같은 SQL 확장, 기본 타임스탬프, 파티셔닝, 클러스터링, 롤업, 그리고 태스크 실행 흐름을 다루어요.

출처: 문서

본문

이 페이지는 멀티 스테이지 쿼리(MSQ) 태스크 엔진을 사용한 SQL 기반 배치 인제스트를 설명해요. 인제스트 방법 표를 보고 어떤 인제스트 방식이 자신에게 맞는지 판단하세요.

MSQ 태스크 엔진 (Multi-stage query task engine)

MSQ 태스크 엔진은 인덱싱 서비스에서 SQL 문을 배치 태스크로 실행하며, 이 태스크는 Middle Managers에서 실행돼요. INSERT와 REPLACE 태스크는 다른 모든 배치 인제스트 방식처럼 세그먼트를 게시해요. 각 쿼리는 실행 중에 최소 두 개의 태스크 슬롯을 차지해요: 컨트롤러 태스크 하나와 워커 태스크 하나 이상이에요. 실험적 기능으로, MSQ 태스크 엔진은 SELECT 쿼리를 배치 태스크로 실행하는 것도 지원해요. INSERT나 REPLACE가 없는 순수 SELECT의 동작과 결과 형식은 변경될 수 있어요.

SQL 문은 웹 콘솔의 Query 뷰나 /druid/v2/sql/task API를 통해 MSQ 태스크 엔진으로 실행할 수 있어요.

MSQ 태스크 엔진으로 SQL 쿼리가 실행되는 방식에 대한 자세한 내용은 멀티 스테이지 쿼리 태스크 섹션을 참고하세요.

SQL 확장 (SQL extensions)

인제스트를 지원하기 위해 MSQ 태스크 엔진을 통해 추가 SQL 기능을 사용할 수 있어요.

EXTERN으로 외부 데이터 읽기 (Read external data with EXTERN)

쿼리 태스크는 EXTERN 함수를 통해 네이티브 배치 입력 소스와 입력 형식을 사용해 외부 데이터에 접근할 수 있어요.

EXTERN은 서로 다른 워커 태스크에서 여러 파일을 병렬로 읽을 수 있어요. 하지만 EXTERN은 개별 파일을 여러 워커 태스크로 분할하지는 않아요. 매우 큰 입력 파일이 몇 개뿐이라면 입력 파일을 나누어 쿼리 병렬성을 높일 수 있어요.

구문에 대한 자세한 내용은 EXTERN 문서를 참고하세요.

또한 EXTERN보다 더 편리할 수 있는 SQL 친화적인 입력 소스별 테이블 함수 집합도 확인하세요.

INSERT로 데이터 로드 (Load data with INSERT)

INSERT 문은 새 데이터 소스를 만들거나 기존 데이터 소스에 추가할 수 있어요. 표준 SQL과 달리 Druid SQL에서는 테이블을 만드는 것과 테이블에 데이터를 추가하는 것 사이에 구문상 차이가 없어요. Druid에는 CREATE TABLE 문이 없어요.

거의 모든 SELECT 기능이 INSERT ... SELECT 쿼리에서 사용 가능해요. 일부 예외는 알려진 문제 페이지에 나열되어 있어요.

INSERT 문은 대상 데이터 소스에 공유 잠금(shared lock)을 획득해요. 클러스터에 충분한 태스크 슬롯이 있다면 같은 데이터 소스에 대해 여러 INSERT 문을 동시에 실행할 수 있어요.

다른 모든 배치 인제스트 방식과 마찬가지로 각 INSERT 문은 새 세그먼트를 생성하고 실행이 끝날 때 게시해요. 이러한 이유로 더 큰 배치로 데이터를 로드하는 데 가장 적합해요. INSERT 문을 연속된 마이크로배치 시퀀스로 데이터를 로드하는 데 사용하지 마세요. 그런 경우에는 스트리밍 인제스트를 사용하세요.

REPLACE를 쓸지 INSERT를 쓸지 결정할 때, REPLACE로 생성된 세그먼트는 차원 기반 프루닝(pruning)이 가능하지만 INSERT로 생성된 것은 불가능하다는 점을 기억하세요. 차원 기반 프루닝 요구 사항에 대한 자세한 내용은 클러스터링 섹션을 참고하세요.

구문에 대한 자세한 내용은 INSERT 문서를 참고하세요.

REPLACE로 데이터 덮어쓰기 (Overwrite data with REPLACE)

REPLACE 문은 새 데이터 소스를 만들거나 기존 데이터 소스의 데이터를 덮어쓸 수 있어요. 표준 SQL과 달리 Druid SQL에서는 테이블을 만드는 것과 테이블의 데이터를 덮어쓰는 것 사이에 구문상 차이가 없어요. Druid에는 CREATE TABLE 문이 없어요.

REPLACE는 어떤 데이터를 덮어쓸지 결정하기 위해 OVERWRITE 절을 사용해요. 전체 테이블이나 테이블의 특정 시간 범위를 덮어쓸 수 있어요. 특정 시간 범위를 덮어쓸 때 그 시간 범위는 PARTITIONED BY 절에 지정된 세분성(granularity)과 정렬되어야 해요.

REPLACE 문은 대상 데이터 소스의 대상 시간 범위에 대해 배타적 쓰기 잠금을 획득해요. 태스크가 실행되는 동안 그 시간 범위에 대해 다른 인제스트나 컴팩션 작업은 진행될 수 없어요. 하지만 다른 시간 범위에 대한 인제스트와 컴팩션 작업은 진행될 수 있어요.

거의 모든 SELECT 기능이 REPLACE ... SELECT 쿼리에서 사용 가능해요. 일부 예외는 알려진 문제 페이지에 나열되어 있어요.

구문에 대한 자세한 내용은 REPLACE 문서를 참고하세요.

REPLACE를 쓸지 INSERT를 쓸지 결정할 때, REPLACE로 생성된 세그먼트는 차원 기반 프루닝이 가능하지만 INSERT로 생성된 것은 불가능하다는 점을 기억하세요. 차원 기반 프루닝 요구 사항에 대한 자세한 내용은 클러스터링 섹션을 참고하세요.

EXTERN으로 외부 대상에 쓰기 (Write to an external destination with EXTERN)

쿼리 태스크는 EXTERN 함수가 INTO 절과 함께 사용될 때(예: INSERT INTO EXTERN(...)) EXTERN 함수를 통해 외부 대상에 데이터를 쓸 수 있어요. EXTERN 함수는 파일을 쓸 위치를 지정하는 인자를 받아요. 형식은 AS 절로 지정할 수 있어요.

구문에 대한 자세한 내용은 EXTERN 문서를 참고하세요.

기본 타임스탬프 (Primary timestamp)

Druid 테이블은 항상 __time이라는 기본 타임스탬프를 포함해요.

날짜 및 시간 함수로 기본 타임스탬프를 설정하는 것이 일반적이에요. 예: TIME_FORMAT("timestamp", 'yyyy-MM-dd HH:mm:ss') AS __time.

__time 컬럼은 시간 기반 파티셔닝에 사용돼요. PARTITIONED BY ALL이나 PARTITIONED BY ALL TIME을 사용하면 시간 기반 파티셔닝이 비활성화돼요. 이 경우 INSERT 문에 __time 컬럼을 포함할 필요가 없어요. 하지만 Druid는 여전히 Druid 테이블에 __time 컬럼을 만들고 모든 타임스탬프를 1970-01-01 00:00:00으로 설정해요.

자세한 내용은 기본 타임스탬프 문서를 참고하세요.

시간 기반 파티셔닝 (Partitioning by time)

INSERT와 REPLACE 문은 PARTITIONED BY 절을 요구하며, 이 절이 시간 기반 파티셔닝이 수행되는 방식을 결정해요. Druid에서 데이터는 PARTITIONED BY 세분성에 의해 정의된 시간 청크별로 하나 이상의 세그먼트로 분할돼요.

시간 기반 파티셔닝은 세 가지 이유로 중요해요.

  1. __time(SQL)이나 intervals(네이티브)로 필터링하는 쿼리는 시간 파티셔닝을 사용해 고려할 세그먼트 집합을 프루닝할 수 있어요.
  2. 기존 데이터의 덮어쓰기와 컴팩션 같은 일부 데이터 관리 작업은 시간 파티션에 배타적 쓰기 잠금을 획득해요. 더 세분화된 파티셔닝은 더 세분화된 배타적 쓰기 잠금을 가능하게 해요.
  3. 각 세그먼트 파일은 시간 파티션 안에 완전히 포함돼요. 너무 세분화된 파티셔닝은 많은 수의 작은 세그먼트를 만들어 성능이 나빠질 수 있어요.

PARTITIONED BY HOUR와 PARTITIONED BY DAY는 이러한 고려 사항의 균형을 맞추는 가장 일반적인 선택이에요.

데이터셋에 기본 타임스탬프가 없다면 PARTITIONED BY ALL이 적합해요.

구문에 대한 자세한 내용은 PARTITIONED BY 문서를 참고하세요.

클러스터링 (Clustering)

시간 파티셔닝으로 정의된 각 시간 청크 안에서 데이터는 선택적인 CLUSTERED BY 절로 더 분할될 수 있어요.

예를 들어 PARTITIONED BY HOUR와 CLUSTERED BY hostName으로 시간당 1억 개의 행을 인제스트한다고 가정해 보아요. 인제스트 태스크는 rowsPerSegment의 기본값인 약 300만 행의 세그먼트를 생성하며, hostName의 사전식 범위가 세그먼트로 그룹화돼요.

클러스터링은 두 가지 이유로 중요해요.

  1. 향상된 지역성(locality)으로 인한 더 낮은 스토리지 용량, 따라서 더 나은 압축성.
  2. 차원 기반 세그먼트 프루닝으로 인한 더 나은 쿼리 성능. 이는 쿼리 필터와 일치하는 데이터를 절대 포함할 수 없는 세그먼트를 고려 대상에서 제거해요. 이는 x = 'foo'와 x IN ('foo', 'bar') 같은 필터를 빠르게 만들어요.

차원 기반 프루닝을 활성화하려면 다음 요구 사항이 충족되어야 해요.

  • 세그먼트가 INSERT 문이 아닌 REPLACE 문으로 생성되었을 것.
  • CLUSTERED BY가 단일 값 문자열 컬럼으로 시작할 것. 이 단일 값 문자열 컬럼이 프루닝에 사용돼요.

이 요구 사항이 충족되지 않으면 Druid는 인제스트 중에 여전히 데이터를 클러스터링하지만 쿼리 시점에는 차원 기반 세그먼트 프루닝을 수행할 수 없어요. 차원 기반 세그먼트 프루닝이 가능한지 확인하려면 sys.segments 테이블을 사용해 인제스트 쿼리가 생성한 세그먼트의 shard_spec을 검사하세요. 타입이 range나 single이면 차원 기반 세그먼트 프루닝이 가능해요. 그렇지 않으면 불가능해요. 샤드 스펙 타입은 Segments 뷰의 Partitioning 컬럼에서도 볼 수 있어요.

구문에 대한 자세한 내용은 CLUSTERED BY 문서를 참고하세요.

클러스터링의 메커니즘에 대한 자세한 내용은 보조 파티셔닝과 정렬 문서를 참고하세요.

롤업 (Rollup)

롤업은 인제스트 중에 데이터를 미리 집계해 저장되는 데이터 양을 줄이는 기법이에요. 중간 집계가 생성된 세그먼트에 저장되고, 추가 집계는 쿼리 시점에 수행돼요. 이는 스토리지 용량을 줄이고 성능을 향상시키며, 종종 크게 향상시켜요.

롤업으로 인제스트를 수행하려면:

  1. GROUP BY를 사용하세요. GROUP BY 절의 컬럼은 차원이 되고, 집계 함수는 지표(metric)가 돼요.
  2. 컨텍스트에서 finalizeAggregations: false를 설정하세요. 이렇게 하면 집계 함수가 최종 결과 대신 내부 상태를 생성된 세그먼트에 작성하고, 쿼리 시점의 추가 집계를 가능하게 해요.
  3. ARRAY 컬럼 인제스트에 대한 정보는 ARRAY 타입 문서를 참고하세요.
  4. 다중 값 VARCHAR 컬럼 인제스트에 대한 정보는 다중 값 차원 문서를 참고하세요.

이 모든 작업을 수행하면 Druid는 사용자가 롤업으로 인제스트하려는 것임을 이해하고 생성된 세그먼트에 롤업 관련 메타데이터를 기록해요. 그러면 다른 애플리케이션이 segmentMetadata 쿼리로 롤업 관련 정보를 검색할 수 있어요.

다음 집계 함수는 인제스트 시점의 롤업에 지원돼요:

COUNT(쿼리 시점에 SUM으로 전환), SUM, MIN, MAX, EARLIEST 및 EARLIEST_BY, LATEST 및 LATEST_BY, APPROX_COUNT_DISTINCT, APPROX_COUNT_DISTINCT_BUILTIN, APPROX_COUNT_DISTINCT_DS_HLL, APPROX_COUNT_DISTINCT_DS_THETA, DS_QUANTILES_SKETCH(쿼리 시점에 APPROX_QUANTILE_DS로 전환). AVG는 사용하지 마세요. 대신 인제스트 시점에 SUM과 COUNT를 사용하고 쿼리 시점에 몫을 계산하세요.

예시는 롤업이 포함된 INSERT 예제를 참고하세요.

MSQ 태스크 (Multi-stage query tasks)

실행 흐름 (Execution flow)

태스크 엔드포인트 /druid/v2/sql/task로 SQL 문을 실행하면 다음이 발생해요.

  1. Broker가 평소처럼 SQL 쿼리를 네이티브 쿼리로 계획해요.
  2. Broker가 네이티브 쿼리를 query_controller 타입의 태스크로 감싸 인덱싱 서비스에 제출해요.
  3. Broker가 태스크 ID를 반환하고 종료돼요.
  4. 컨트롤러 태스크가 maxNumTasks와 taskAssignment 컨텍스트 파라미터에 의해 결정되는 일정 수의 워커 태스크를 시작해요. 이 설정은 각 쿼리에 대해 개별적으로 설정할 수 있어요.
  5. query_worker 타입의 워커 태스크가 쿼리를 실행해요.
  6. 쿼리가 SELECT 쿼리라면 워커 태스크가 결과를 컨트롤러 태스크로 보내고, 컨트롤러는 이를 태스크 리포트에 기록해요. 쿼리가 INSERT 또는 REPLACE 쿼리라면 워커 태스크가 새 Druid 세그먼트를 생성해 제공된 데이터 소스에 게시해요.

병렬성 (Parallelism)

maxNumTasks 쿼리 파라미터는 쿼리가 사용할 최대 태스크 수를 결정하며, 하나의 query_controller 태스크를 포함해요. 일반적으로 워커가 많을수록 쿼리 성능이 좋아요. maxNumTasks의 가능한 최저값은 2(워커 하나와 컨트롤러 하나)예요. 클러스터에서 사용 가능한 슬롯 수보다 높게 설정하지 마세요. 그러면 TaskStartTimeout 오류가 발생해요.

외부 데이터를 읽을 때 EXTERN은 서로 다른 워커 태스크에서 여러 파일을 병렬로 읽을 수 있어요. 하지만 EXTERN은 개별 파일을 여러 워커 태스크로 분할하지는 않아요. 매우 큰 입력 파일이 몇 개뿐이라면 입력 파일을 나누어 쿼리 병렬성을 높일 수 있어요.

각 Middle Manager의 druid.worker.capacity 서버 프로퍼티는 각 서버에서 한 번에 실행할 수 있는 최대 워커 태스크 수를 결정해요. 워커 태스크는 단일 스레드로 실행되며, 이는 멀티 스테이지 쿼리에 기여할 수 있는 서버의 최대 프로세서 수도 결정해요.

메모리 사용 (Memory usage)

사용 가능한 메모리 양을 늘리면 특정 경우에 성능이 향상될 수 있어요.

  • 데이터가 덜 자주 디스크로 넘칠 때 세그먼트 생성이 더 효율적이 돼요.
  • 사용 가능한 메모리가 필요한 정렬 패스 수에 영향을 주므로 정렬 스테이지 출력 데이터가 더 효율적이 돼요.

워커 태스크는 JVM 힙 메모리와 오프힙("direct") 메모리를 모두 사용해요.

Middle Managers가 시작한 Peon에서 JVM 힙의 대부분(75%, lookup이 사용하는 공간을 뺀 값)은 두 개의 동일한 크기 묶음으로 분할돼요: 하나는 프로세서 묶음, 하나는 워커 묶음이에요. 각각 사용 가능한 JVM 힙의 37.5%(lookup이 사용하는 공간을 뺀 값)를 구성해요.

쿼리 유형에 따라 컨트롤러와 워커 태스크는 파티션 경계를 결정하기 위해 스케치(sketch)를 사용할 수 있어요. 이 스케치의 힙 공간(footprint)은 사용 가능한 메모리의 10% 또는 300MB 중 더 낮은 값으로 제한돼요.

프로세서 메모리 묶음은 쿼리 처리와 세그먼트 생성에 사용돼요. 각 프로세서 묶음은 또한 스테이지 간 I/O를 버퍼링할 공간을 제공해야 해요. 구체적으로 각 하위 스테이지는 각 상위 워커에 대해 1MB의 버퍼 공간을 요구해요. 예를 들어 스테이지 0에서 100개의 워커가 실행되고 스테이지 1이 스테이지 0에서 읽는다면, 스테이지 1의 각 워커는 프레임 버퍼에 1M * 100 = 100 MB의 메모리를 요구해요.

워커 메모리 묶음은 셔플 전 스테이지 출력 데이터를 정렬하는 데 사용돼요. 워커는 메모리에 맞는 것보다 더 많은 데이터를 정렬할 수 있으며, 이 경우 디스크를 사용하도록 전환해요.

워커 태스크는 또한 오프힙("direct") 메모리를 사용해요. 사용 가능한 direct 메모리(-XX:MaxDirectMemorySize)를 최소 (druid.processing.numThreads + 1) * druid.processing.buffer.sizeBytes로 설정하세요. 최소값 이상으로 direct 메모리를 늘려도 처리가 더 빨라지지는 않아요.

디스크 사용 (Disk usage)

워커 태스크는 로컬 디스크를 네 가지 용도로 사용해요.

  • 입력 데이터의 임시 복사본. 각 임시 파일은 다음 파일을 읽기 전에 삭제돼요. 태스크당 한 번에 입력 파일 하나를 저장할 충분한 임시 디스크 공간만 있으면 돼요.
  • 세그먼트 생성 관련 임시 데이터. 태스크당 한 번에 세그먼트 하나 분량의 데이터를 저장할 충분한 임시 디스크 공간만 있으면 돼요. 일반적으로 태스크당 2GB 미만이에요.
  • 셔플 전 데이터의 외부 정렬. 태스크에 대해 전체 출력 데이터셋의 압축 복사본을 저장할 충분한 공간이 필요해요.
  • 셔플 중 스테이지 출력 데이터 저장. 태스크에 대해 전체 출력 데이터셋의 압축 복사본을 저장할 충분한 공간이 필요해요.

워커는 이러한 항목에 대해 druid.indexer.task.baseDir이 제공하는 태스크 작업 디렉터리를 사용해요. 이 디렉터리에 이러한 용도를 위한 충분한 공간이 있는 것이 중요해요.

더 알아보기 (Learn more)

  • MSQ 참조 문서에서 INSERT·REPLACE·EXTERN·PARTITIONED BY·CLUSTERED BY의 파라미터를 확인해 보세요.
  • MSQ 알려진 문제에서 현재 지원되지 않는 기능을 살펴보아요.
  • MSQ 예제에서 SQL 기반 인제스트를 실제 예시로 확인해 보세요.