JSON 기반 배치

JSON 기반 배치 (JSON-based batch)

정보

이 페이지는 ingestion spec을 사용한 JSON 기반 배치 수집을 설명해요. druid-multi-stage-query 엔진을 사용한 SQL 기반 배치 수집은 SQL-based ingestion을 참고하세요. 어떤 수집 방법이 적합한지는 ingestion methods 표를 참고해 주세요.

Apache Druid는 다음 타입의 JSON 기반 배치 indexing 태스크를 지원합니다.

  • 여러 indexing 태스크를 동시에 실행할 수 있는 Parallel task 인덱싱(index_parallel). Parallel task는 프로덕션 수집 태스크에 잘 맞아요.
  • 한 번에 단일 indexing 태스크를 실행하는 Simple task 인덱싱(index). Simple task 인덱싱은 개발·테스트 환경에 적합합니다.

이 주제는 index_parallel ingestion spec의 구성을 다룹니다.

배치 인덱싱 관련 정보는 다음을 참고하세요.

  • 배치 수집 방법 비교: Batch ingestion method comparison table
  • 튜토리얼: JSON 기반 배치 수집 튜토리얼은 Tutorial: Loading a file
  • 가능한 input source: Input sources
  • 가능한 input format: Source input formats

출처: 문서

본문

indexing 태스크 제출하기

두 종류의 JSON 기반 배치 indexing 태스크를 실행하려면 다음 중 하나를 할 수 있어요.

  • 웹 콘솔의 Load Data UI를 사용해서 ingestion spec을 정의·제출
  • 배치 인덱싱 예시와 레퍼런스 주제를 기반으로 ingestion spec을 JSON으로 정의한 뒤, 그 ingestion spec을 Overlord 서비스의 Tasks API 엔드포인트인 /druid/indexer/v1/task에 POST. 또는 Druid에 포함된 bin/post-index-task 스크립트를 사용할 수 있음

Parallel task 인덱싱

Parallel task 타입 index_parallel은 멀티스레드 배치 인덱싱을 위한 태스크예요.

index_parallel 태스크는 전체 인덱싱 프로세스를 조율하는 supervisor 태스크입니다. supervisor 태스크는 입력 데이터를 분할하고 데이터의 개별 부분들을 처리할 worker 태스크를 만듭니다.

Druid는 worker 태스크를 Overlord에 발행합니다. Overlord는 Middle Manager나 Indexer에서 worker를 예약·실행합니다. worker 태스크가 할당된 입력 부분을 성공적으로 처리하면, 결과 세그먼트 목록을 Supervisor 태스크에 보고합니다.

Supervisor 태스크는 주기적으로 worker 태스크의 상태를 확인합니다. 태스크가 실패하면 Supervisor는 재시도 횟수가 구성된 한도에 도달할 때까지 태스크를 재시도해요. 모든 worker 태스크가 성공하면 보고된 세그먼트를 한 번에 게시하고 수집을 마무리합니다.

Parallel task의 세부 동작은 partitionsSpec에 따라 달라집니다. 자세한 내용은 partitionsSpec을 참고하세요.

Parallel task에는 다음이 필요합니다.

  • ioConfig에 분할 가능한 inputSource. 지원되는 분할 가능 input format 목록은 Splittable input sources 참고
  • tuningConfig의 maxNumConcurrentSubTasks가 1보다 큼. 그렇지 않으면 태스크가 순차 실행됨. index_parallel 태스크는 각 입력 파일을 하나씩 읽고 스스로 세그먼트를 만듭니다.

지원되는 압축 형식

JSON 기반 배치 수집은 다음 압축 형식을 지원합니다.

  • bz2
  • gz
  • xz
  • zip
  • sz (Snappy)
  • zst (ZSTD)

구현 고려사항

이 섹션은 parallel task 수집을 구현할 때 고려할 구현 세부 사항을 다룹니다.

worker 태스크의 볼륨 제어

병렬 수집의 단계에 따라 다른 구성으로 각 worker 태스크가 처리하는 입력 데이터 양을 제어할 수 있어요. 파티셔닝이 태스크의 데이터 볼륨에 어떻게 영향을 주는지에 대한 자세한 내용은 partitionsSpec 참고.

inputSource에서 데이터를 읽는 태스크에는 tuningConfig에서 Split hint spec을 설정할 수 있고, 셔플된 세그먼트를 병합하는 태스크에는 tuningConfig에서 totalNumMergeTasks를 설정할 수 있습니다.

실행 중인 태스크 수

tuningConfig의 maxNumConcurrentSubTasks는 병렬로 실행되는 동시 worker 태스크 수를 결정합니다. Supervisor 태스크는 현재 실행 중인 worker 태스크 수를 확인하고, 그것이 maxNumConcurrentSubTasks보다 작으면 사용 가능한 태스크 슬롯 수와 무관하게 더 만듭니다. 이는 다른 수집 성능에 영향을 줄 수 있어요. 자세한 내용은 Capacity planning 섹션 참고.

데이터 대체 또는 append

기본적으로 JSON 기반 배치 수집은 쓰는 모든 세그먼트에 대해 granularitySpec의 interval에 있는 모든 데이터를 대체합니다. 대신 세그먼트에 추가하고 싶다면 ioConfig에서 appendToExisting 플래그를 설정하세요. JSON 기반 배치 수집은 데이터를 적극적으로 추가하는 세그먼트에서만 데이터를 대체합니다. granularitySpec의 interval에 태스크에서 온 데이터가 없는 세그먼트는 변경되지 않은 채 남습니다. 기존 세그먼트가 granularitySpec의 interval과 부분적으로 겹치면, 새 spec의 interval 밖에 있는 그 세그먼트의 부분은 계속 보입니다.

동시 append와 replace 태스크도 수행할 수 있어요. 자세한 내용은 Concurrent append and replace를 참고하세요.

tombstone으로 기존 세그먼트 완전 대체

ioConfig에서 dropExisting 플래그를 true로 설정하면 수집 태스크가 granularitySpec의 interval 안에서 시작하고 끝나는 모든 기존 세그먼트를 대체할 수 있어요. 이는 새 데이터가 모든 기존 세그먼트를 덮는지 여부와 무관합니다. dropExisting은 appendToExisting이 false이고 granularitySpec에 interval이 포함된 경우에만 적용됩니다.

다음 예시들은 ioConfig에서 dropExisting 속성을 true로 설정해야 하는 경우를 보여줘요.

2020-01-01부터 2021-01-01까지의 interval과 YEAR segmentGranularity를 가진 기존 세그먼트를 고려해 보세요. 더 세밀한 MONTH segmentGranularity로 2020-01-01부터 2021-01-01까지의 전체 interval을 새 데이터로 덮어쓰고 싶습니다. 대체 데이터가 2020-01-01부터 2021-01-01까지의 매 달마다 레코드를 가지지 않으면, Druid는 모든 대체 데이터를 포함하더라도 원래 YEAR 세그먼트를 버릴 수 없어요. 이 경우 원래 YEAR segmentGranularity 세그먼트가 더 이상 필요 없으므로 dropExisting을 true로 설정해서 대체하세요.

datasource를 다시 수집하거나 덮어쓰고 싶은데 새 데이터에 datasource에 존재하는 일부 시간 interval이 없는 경우를 상상해 보세요. 예를 들어 datasource가 MONTH segmentGranularity로 다음 데이터를 포함합니다.

  • 1월: 레코드 1개
  • 2월: 레코드 10개
  • 3월: 레코드 10개

다음과 같이 다시 수집하고 새 데이터로 덮어쓰고 싶습니다.

  • 1월: 레코드 0개
  • 2월: 레코드 10개
  • 3월: 레코드 9개

dropExisting을 true로 설정하지 않으면, 같은 MONTH segmentGranularity로 덮어쓰며 수집한 결과는 다음과 같습니다.

  • 1월: 레코드 1개
  • 2월: 레코드 10개
  • 3월: 레코드 9개

새 데이터가 1월에 0개 레코드를 가지므로 이것은 기대와 다를 수 있어요. 불필요한 1월 세그먼트를 tombstone으로 대체하려면 dropExisting을 true로 설정하세요.

Parallel indexing 예시

다음 예시는 parallel indexing 태스크의 구성을 보여줍니다.

예시 보기

{  "type": "index_parallel",  "spec": {    "dataSchema": {      "dataSource": "wikipedia_parallel_index_test",      "timestampSpec": {        "column": "timestamp"      },      "dimensionsSpec": {        "dimensions": [          "country",          "page",          "language",          "user",          "unpatrolled",          "newPage",          "robot",          "anonymous",          "namespace",          "continent",          "region",          "city"        ]      },      "metricsSpec": [        {          "type": "count",          "name": "count"        },        {          "type": "doubleSum",          "name": "added",          "fieldName": "added"        },        {          "type": "doubleSum",          "name": "deleted",          "fieldName": "deleted"        },        {          "type": "doubleSum",          "name": "delta",          "fieldName": "delta"        }      ],      "granularitySpec": {        "segmentGranularity": "DAY",        "queryGranularity": "second",        "intervals": [          "2013-08-31/2013-09-02"        ]      }    },    "ioConfig": {      "type": "index_parallel",      "inputSource": {        "type": "local",        "baseDir": "examples/indexing/",        "filter": "wikipedia_index_data*"      },      "inputFormat": {        "type": "json"      }    },    "tuningConfig": {      "type": "index_parallel",      "partitionsSpec": {        "type": "single_dim",        "partitionDimension": "country",        "targetRowsPerSegment": 5000000      },      "maxNumConcurrentSubTasks": 2    }  }}

Parallel indexing 구성

다음 표는 input spec의 주요 섹션을 정의합니다.

| Property | Description | Required | | type | 태스크 타입. parallel task 인덱싱에서는 값을 index_parallel로 설정 | yes | | id | 태스크 ID. 생략하면 Druid가 태스크 타입, 데이터 소스 이름, interval, 날짜-시간 스탬프로 태스크 ID를 생성 | no | | spec | data schema, IO config, tuning config를 정의하는 ingestion spec | yes | | context | 다양한 태스크 구성 파라미터를 지정하는 컨텍스트. 자세한 내용은 Task context parameters 참고 | no |

dataSchema

이 필드는 필수입니다. 일반적으로 Druid가 데이터를 저장하는 방식을 정의합니다: 기본 타임스탬프 컬럼, dimensions, metrics, 그리고 어떤 변환. 개요는 Ingestion Spec DataSchema를 참고하세요.

index parallel용 granularitySpec을 정의할 때, 데이터의 시간 범위를 안다면 intervals를 명시적으로 정의하는 것을 고려하세요. 이렇게 하면 잠금 실패가 더 빨리 일어나고, 일부 행이 예상 밖의 타임스탬프를 가지고 있어도 Druid가 interval 범위 밖의 데이터를 실수로 대체하지 않습니다. 이유는 다음과 같습니다.

  • intervals를 명시적으로 정의하면 JSON 기반 배치 수집은 시작할 때 지정된 모든 interval을 잠급니다. 여러 수집·인덱싱 태스크가 같은 interval에 잠금을 얻으려 하면 잠금 문제가 빠르게 드러납니다. 예를 들어 Kafka 수집 태스크가 잠긴 interval에 잠금을 얻으려 하면 수집 태스크가 실패합니다. 또한 지정된 interval 밖의 행이 있으면 Druid가 그것들을 버려 예상 밖의 interval과의 충돌을 피합니다.
  • intervals를 정의하지 않으면 JSON 기반 배치 수집은 interval이 발견될 때 각 interval을 잠급니다. 이 경우 태스크가 더 높은 우선순위 태스크와 겹치면, 잠금 충돌 문제가 수집 프로세스 후반에 발생합니다. 소스 데이터에 예상 밖의 타임스탬프를 가진 행이 포함되면, 예상 밖의 interval 잠금을 유발할 수 있습니다.

ioConfig

다음 표는 ioConfig 객체의 속성을 나열합니다.

| Property | Description | Default | Required | | type | 태스크 타입. 값을 index_parallel로 설정 | none | yes | | inputFormat | 입력 데이터를 어떻게 파싱할지 지정하는 inputFormat | none | yes | | appendToExisting | 세그먼트를 최신 버전의 추가 샤드로 만들어, 세그먼트 집합을 대체하는 대신 효과적으로 append함. 즉 원래 파티셔닝 방식과 무관하게 어떤 datasource에도 새 세그먼트를 append할 수 있음. append된 세그먼트에는 반드시 dynamic 파티셔닝 타입을 사용해야 함. 다른 타입을 지정하면 태스크가 오류로 실패함 | false | no | | dropExisting | true이고 appendToExisting이 false이고 granularitySpec에 interval이 포함되어 있으면, 수집 태스크가 새 세그먼트를 게시할 때 지정된 interval에 완전히 포함된 모든 기존 세그먼트를 대체함. 수집이 실패하면 Druid는 어떤 기존 세그먼트도 변경하지 않음. appendToExisting이 true이거나 granularitySpec에 interval이 지정되지 않은 오설정의 경우, dropExisting이 true여도 Druid는 어떤 세그먼트도 대체하지 않음 | false | no |

tuningConfig

tuningConfig는 선택 사항입니다. tuningConfig를 지정하지 않으면 Druid는 기본 파라미터를 사용해요.

다음 표는 tuningConfig 객체의 속성을 나열합니다.

| Property | Description | Default | Required | | type | 태스크 타입. 값을 index_parallel로 설정 | none | yes | | maxRowsInMemory | Druid가 디스크에 중간 persist를 수행해야 하는 시점을 결정. 보통 설정할 필요 없음. 데이터 성격에 따라 행이 바이트로 짧을 때, 예를 들어 백만 행을 메모리에 저장하고 싶지 않을 때 이 값을 설정 | 1000000 | no | | maxBytesInMemory | Druid가 디스크에 중간 persist를 수행해야 하는 시점을 결정하는 데 사용. 보통 Druid가 내부적으로 계산하므로 설정할 필요 없음. 이 값은 persist 전에 힙 메모리에 모을 바이트 수를 나타냄. 메모리 사용의 대략적 추정치에 기반하고 실제 사용량은 아님. 인덱싱의 최대 힙 메모리 사용량은 maxBytesInMemory * (2 + maxPendingPersists). 참고로 maxBytesInMemory는 중간 persist에서 만들어진 아티팩트의 힙 사용도 포함. 즉 매 persist 후 다음 persist까지 사용 가능한 maxBytesInMemory 양이 줄어듦. 모든 중간 persisted 아티팩트의 바이트 합이 maxBytesInMemory를 초과하면 태스크는 실패 | max JVM 메모리의 1/6 | no | | maxColumnsToMerge | 게시를 위해 세그먼트를 병합할 때 단일 단계에서 병합할 세그먼트 수의 한도. 이 한도는 병합할 세그먼트 집합에 존재하는 총 컬럼 수에 영향을 줌. 한도를 초과하면 세그먼트 병합이 여러 단계로 발생. Druid는 이 설정과 무관하게 단계당 최소 2개 세그먼트를 병합 | -1 (무제한) | no | | maxTotalRows | Deprecated. 대신 partitionsSpec 사용. 푸시를 기다리는 세그먼트의 총 행 수. 중간 push가 언제 발생할지 결정하는 데 사용 | 20000000 | no | | numShards | Deprecated. 대신 partitionsSpec 사용. hashed partitionsSpec을 사용할 때 만들 샤드 수를 직접 지정. 이 값이 지정되고 granularitySpec에 intervals가 지정되면 index task는 데이터를 훑어 interval/파티션을 결정하는 패스를 건너뜀 | null | no | | splitHintSpec | 각 첫 단계 태스크가 읽는 데이터 양을 제어하는 힌트. input source의 구현에 따라 Druid가 힌트를 무시할 수 있음. 자세한 내용은 Split hint spec 참고 | size-based split hint spec | no | | partitionsSpec | 각 timeChunk에서 데이터를 어떻게 파티셔닝할지 정의. PartitionsSpec 참고 | forceGuaranteedRollup = false면 dynamic, forceGuaranteedRollup = true면 hashed 또는 single_dim | no | | indexSpec | 인덱싱 시점에 사용할 세그먼트 저장 형식 옵션 정의. IndexSpec 참고 | null | no | | indexSpecForIntermediatePersists | 인덱싱 시점에 중간 persisted 임시 세그먼트에 사용할 세그먼트 저장 형식 옵션 정의. 이 구성으로 중간 세그먼트의 dimension/metric 압축을 비활성화해서 최종 병합에 필요한 메모리를 줄일 수 있음. 다만 중간 세그먼트의 압축을 끄면, Druid가 그것들을 최종 게시 세그먼트로 병합하기 전 중간 세그먼트가 사용되는 동안 페이지 캐시 사용이 늘어날 수 있음. 가능한 값은 IndexSpec 참고 | indexSpec과 동일 | no | | maxPendingPersists | 시작되지 않은 채 대기할 수 있는 pending persist 최대 수. 새 중간 persist가 이 한도를 초과하면 현재 실행 중인 persist가 끝날 때까지 수집이 블록됨. 인덱싱의 최대 힙 메모리 사용량은 maxRowsInMemory * (2 + maxPendingPersists)에 비례 | 0 (수집과 동시에 실행될 수 있는 persist는 1개이고, 큐에 쌓일 수 있는 것은 없음을 의미) | no | | forceGuaranteedRollup | perfect rollup을 강제함. perfect rollup은 생성된 세그먼트의 총 크기와 쿼리 시간을 최적화하지만 인덱싱 시간은 늘어남. true면 granularitySpec에 intervals를 지정하고 partitionsSpec에 hashed 또는 single_dim을 사용해야 함. 이 플래그는 IOConfig의 appendToExisting과 함께 사용할 수 없음. 자세한 내용은 Segment pushing modes 참고 | false | no | | reportParseExceptions | true면 Druid가 파싱 중 발생한 예외를 던져 수집을 중단. false면 Druid가 파싱할 수 없는 행·필드를 건너뜀 | false | no | | pushTimeout | 세그먼트 푸시를 기다리는 밀리초. >= 0이어야 하고, 0은 영원히 기다림을 의미 | 0 | no | | segmentWriteOutMediumFactory | 세그먼트를 만들 때 사용할 세그먼트 write-out 매체. SegmentWriteOutMediumFactory 참고 | 지정하지 않으면 druid.peon.defaultSegmentWriteOutMediumFactory.type 값을 사용 | no | | maxNumConcurrentSubTasks | 동시에 병렬로 실행될 수 있는 worker 태스크 최대 수. supervisor 태스크는 현재 사용 가능한 태스크 슬롯과 무관하게 maxNumConcurrentSubTasks까지 worker 태스크를 생성. 이 값이 1이면 supervisor 태스크는 worker 태스크를 생성하는 대신 스스로 데이터 수집을 처리. 이 값이 너무 크게 설정되면 supervisor가 다른 수집 태스크를 막는 너무 많은 worker 태스크를 만들 수 있음. 자세한 내용은 Capacity planning 참고 | 1 | no | | maxRetry | 태스크 실패 시 최대 재시도 횟수 | 3 | no | | maxNumSegmentsToMerge | 두 번째 단계에서 단일 태스크가 동시에 병합할 수 있는 세그먼트 수의 최대 한도. forceGuaranteedRollup가 true일 때만 사용됨 | 100 | no | | totalNumMergeTasks | partitionsSpec이 hashed 또는 single_dim으로 설정됐을 때 병합 단계에서 세그먼트를 병합하는 태스크의 총 수 | 10 | no | | taskStatusCheckPeriodMs | 실행 중인 태스크 상태를 확인하는 폴링 주기(밀리초) | 1000 | no | | chatHandlerTimeout | worker 태스크에서 푸시된 세그먼트를 보고하는 타임아웃 | PT10S | no | | chatHandlerNumRetries | worker 태스크에서 푸시된 세그먼트를 보고하는 재시도 횟수 | 5 | no | | awaitSegmentAvailabilityTimeoutMillis | 수집 완료 후 새로 인덱싱된 세그먼트가 쿼리 가능해질 때까지 기다리는 밀리초. <= 0이면 대기하지 않음. > 0이면 태스크가 새 세그먼트가 쿼리 가능하다는 Coordinator 표시를 기다림. 타임아웃이 경과하면 태스크는 성공으로 종료되지만 세그먼트가 쿼리 가능하다는 것은 확인되지 않음 | Long | no (기본 = 0) |

Split Hint Spec

split hint spec은 supervisor 태스크가 input source를 나누는 데 도움을 주기 위해 사용됩니다. 각 worker 태스크는 단일 입력 분할을 처리합니다. 첫 단계에서 각 worker 태스크가 읽는 데이터 양을 제어할 수 있어요.

Size-based Split Hint Spec

size-based split hint spec은 HTTP input source와 SQL input source를 제외한 모든 분할 가능 input source에 영향을 줍니다.

| Property | Description | Default | Required | | type | 값을 maxSize로 설정 | none | yes | | maxSplitSize | 단일 서브태스크에서 처리할 입력 파일의 최대 바이트 수. 단일 파일이 이 한도보다 크면 Druid는 그 파일을 단일 서브태스크에서 혼자 처리. Druid는 파일을 태스크 사이에 나누지 않음. 하나의 서브태스크는 총 크기가 maxSplitSize보다 작아도 maxNumFiles보다 많은 파일을 처리하지 않음. Human-readable format 지원 | 1GiB | no | | maxNumFiles | 단일 서브태스크에서 처리할 입력 파일의 최대 수. 이 한도는 ingestion spec이 너무 길 때 태스크 실패를 방지. 직렬화된 ingestion spec의 최대 크기에 두 가지 알려진 한도가 있음: ZooKeeper의 최대 ZNode 크기(jute.maxbuffer)와 MySQL의 최대 패킷 크기(max_allowed_packet). 직렬화된 ingestion spec 크기가 그 중 하나에 도달하면 이 한도들이 수집 태스크를 실패시킬 수 있음. 하나의 서브태스크는 총 파일 수가 maxNumFiles보다 작아도 maxSplitSize보다 많은 데이터를 처리하지 않음 | 1000 | no |

Segments Split Hint Spec

segments split hint spec은 DruidInputSource에만 사용됩니다.

| Property | Description | Default | Required | | type | 값을 segments로 설정 | none | yes | | maxInputSegmentBytesPerTask | 단일 서브태스크에서 처리할 입력 세그먼트의 최대 바이트 수. 단일 세그먼트가 이 숫자보다 크면 Druid는 그 세그먼트를 단일 서브태스크에서 혼자 처리. Druid는 입력 세그먼트를 태스크 사이에 나누지 않음. 단일 서브태스크는 총 크기가 maxInputSegmentBytesPerTask보다 작아도 maxNumSegments보다 많은 세그먼트를 처리하지 않음. Human-readable format 지원 | 1GiB | no | | maxNumSegments | 단일 서브태스크에서 처리할 입력 세그먼트의 최대 수. 이 한도는 ingestion spec이 너무 길어서 발생하는 실패를 방지. 직렬화된 ingestion spec의 최대 크기에 두 가지 알려진 한도가 있음: ZooKeeper의 최대 ZNode 크기(jute.maxbuffer)와 MySQL의 최대 패킷 크기(max_allowed_packet). 직렬화된 ingestion spec 크기가 그 중 하나에 도달하면 이 한도들이 수집 태스크를 실패시킬 수 있음. 단일 서브태스크는 총 세그먼트 수가 maxNumSegments보다 작아도 maxInputSegmentBytesPerTask보다 많은 데이터를 처리하지 않음 | 1000 | no |

partitionsSpec

Druid의 기본 파티션은 시간입니다. partitions spec에서 보조 파티셔닝 방법을 정의할 수 있어요. 롤업 방법에 적용되는 partitionsSpec 타입을 사용하세요.

perfect rollup에는 다음을 사용할 수 있습니다.

  • 각 행의 지정된 dimension들의 해시 값에 기반한 hashed 파티셔닝
  • 단일 dimension의 값 범위에 기반한 single_dim
  • 여러 dimension의 값 범위에 기반한 range

best-effort rollup에는 dynamic을 사용하세요.

개요는 Partitioning을 참고하세요.

partitionsSpec 타입들은 각각 다른 특징이 있습니다.

| PartitionsSpec | Ingestion speed | Partitioning method | Supported rollup mode | Secondary partition pruning at query time | | dynamic | Fastest | 세그먼트의 행 수에 기반한 동적 파티셔닝 | Best-effort rollup | N/A | | hashed | Moderate | 여러 dimension 해시 기반 파티셔닝은 데이터 locality를 개선해서 datasource 크기와 쿼리 지연 시간을 줄일 수 있음. 자세한 내용은 Partitioning 참고 | Perfect rollup | broker가 파티션 정보를 사용해서 세그먼트를 일찍 잘라내 쿼리를 빠르게 할 수 있음. broker는 partitionDimensions 값을 해시해서 세그먼트를 찾는 방법을 알므로, 모든 partitionDimensions에 대한 필터를 포함한 쿼리가 주어지면 broker는 partitionDimensions 필터를 만족하는 행을 가진 세그먼트만 골라 쿼리 처리할 수 있음. 참고로 쿼리 시점의 보조 파티션 프루닝을 활성화하려면 수집 시점에 partitionDimensions를 설정해야 함 | | single_dim | Slower | 단일 dimension 범위 파티셔닝은 데이터 locality를 개선해서 datasource 크기와 쿼리 지연을 줄일 수 있음. 자세한 내용은 Partitioning 참고 | Perfect rollup | broker가 각 세그먼트의 partitionDimension 값 범위를 알므로, partitionDimension에 대한 필터를 포함한 쿼리가 주어지면 broker는 그 필터를 만족하는 행을 가진 세그먼트만 골라 쿼리 처리할 수 있음 | | range | Slowest | 여러 dimension 범위 파티셔닝은 데이터 locality를 개선해서 datasource 크기와 쿼리 지연을 줄일 수 있음. 자세한 내용은 Partitioning 참고 | Perfect rollup | broker가 각 세그먼트 안 partitionDimensions 값의 범위를 알므로, partitionDimensions의 첫 번째에 대한 필터를 포함한 쿼리가 주어지면 broker는 첫 번째 파티션 dimension의 필터를 만족하는 행을 가진 세그먼트만 골라 쿼리 처리할 수 있음 |

Dynamic partitioning

| Property | Description | Default | Required | | type | 값을 dynamic으로 설정 | none | yes | | maxRowsPerSegment | 샤딩에 사용. 각 세그먼트에 몇 행이 들어갈지 결정. 값은 0보다 커야 함 | 5000000 | no | | maxTotalRows | 푸시를 기다리는 모든 세그먼트의 총 행 수. 중간 세그먼트 push가 언제 발생할지 결정하는 데 사용. 값은 0보다 커야 함 | 20000000 | no |

dynamic 파티셔닝에서 parallel index 태스크는 여러 worker 태스크(타입 single_phase_sub_task)를 생성하는 단일 단계로 실행되며, 각 태스크가 세그먼트를 만듭니다.

worker 태스크가 세그먼트를 만드는 방식:

  • 현재 세그먼트의 행 수가 maxRowsPerSegment를 초과할 때마다
  • 모든 시간 청크에 걸친 모든 세그먼트의 총 행 수가 maxTotalRows에 도달할 때. 이 시점에 태스크는 지금까지 만든 모든 세그먼트를 딥 스토리지로 푸시하고 새 것을 만듭니다.

Hash 기반 파티셔닝

| Property | Description | Default | Required | | type | 값을 hashed로 설정 | none | yes | | numShards | 만들 샤드 수를 직접 지정. 이 값이 지정되고 granularitySpec에 intervals가 지정되면 index task는 데이터를 훑어 interval/파티션을 결정하는 패스를 건너뜀. 이 속성과 targetRowsPerSegment는 둘 다 설정할 수 없음 | none | no | | targetRowsPerSegment | 각 파티션의 목표 행 수. numShards가 미지정이면 Parallel task가 각 파티션이 목표에 가까운 행 수를 가지도록 파티션 수를 자동 결정(입력 데이터의 키가 고르게 분포돼 있다고 가정). numShards와 targetRowsPerSegment가 둘 다 null이면 세그먼트당 목표 행 수 500만이 사용됨 | null (또는 numShards와 targetRowsPerSegment가 둘 다 null이면 5,000,000) | no | | partitionDimensions | 파티셔닝할 dimension들. 비워두면 모든 dimension 선택 | null | no | | partitionFunction | 파티션 dimension의 해시를 계산하는 함수. Hash partition function 참고 | murmur3_32_abs | no |

hash 기반 파티셔닝의 Parallel task는 MapReduce와 비슷합니다. 태스크는 최대 세 단계로 실행됩니다: partial dimension cardinality, partial segment generation, partial segment merge.

partial dimension cardinality 단계는 numShards가 지정되지 않은 경우에만 실행되는 선택적 단계입니다. Parallel task는 split hint spec에 따라 입력 데이터를 분할하고 worker 태스크에 할당합니다. 각 worker 태스크(타입 partial_dimension_cardinality)는 각 시간 청크에 대한 파티셔닝 dimension 카디널리티의 추정치를 수집합니다. Parallel task는 이 추정치들을 worker 태스크에서 집계해서 입력 데이터의 모든 시간 청크에 걸친 최고 카디널리티를 결정하고, 이 카디널리티를 targetRowsPerSegment로 나눠 numShards를 자동 결정합니다.

partial segment generation 단계에서는 MapReduce의 Map 단계처럼, Parallel task가 split hint spec에 기반해 입력 데이터를 분할하고 각 분할을 worker 태스크에 할당합니다. 각 worker 태스크(타입 partial_index_generate)는 할당된 분할을 읽고, granularitySpec의 시간 청크 segmentGranularity(기본 파티션 키)로, 그다음 partitionsSpec의 partitionDimensions의 해시 값(보조 파티션 키)으로 행을 파티셔닝합니다. 파티셔닝된 데이터는 Middle Manager 또는 Indexer의 로컬 스토리지에 저장됩니다.

partial segment merge 단계는 MapReduce의 Reduce 단계와 비슷합니다. Parallel task는 새 worker 태스크 집합(타입 partial_index_generic_merge)을 만들어 이전 단계에서 만든 파티셔닝된 데이터를 병합합니다. 여기서 파티셔닝된 데이터는 병합되도록 시간 청크와 partitionDimensions의 해시 값에 따라 셔플됩니다. 각 worker 태스크는 여러 Middle Manager/Indexer 프로세스에서 같은 시간 청크와 같은 해시 값에 들어가는 데이터를 읽고 병합해서 최종 세그먼트를 만듭니다. 마지막으로 최종 세그먼트를 한 번에 딥 스토리지로 푸시합니다.

Hash partition function

hash 파티셔닝에서 partition function은 파티션 dimension의 해시를 계산하는 데 사용됩니다. 파티션 dimension 값들은 먼저 전체적으로 바이트 배열로 직렬화된 다음, partition function이 적용되어 바이트 배열의 해시를 계산합니다. Druid는 현재 하나의 partition function만 지원합니다.

| name | description | | murmur3_32_abs | murmur3_32의 결과에 절댓값 함수를 적용 |

단일 dimension 범위 파티셔닝

정보

단일 dimension 범위 파티셔닝은 index_parallel 태스크 타입의 순차(sequential) 모드에서 지원되지 않습니다.

범위 파티셔닝은 저장 공간(footprint)과 쿼리 성능과 관련해 여러 이점이 있습니다.

maxNumConcurrentSubTasks를 1로 설정하면 Parallel task는 하나의 서브태스크를 사용합니다.

이 기법으로 데이터를 파티셔닝할 때, partitionDimension의 데이터가 고르지 않게 분포돼 있으면 세그먼트 크기도 고르지 않게 분포될 수 있어요. 따라서 데이터 배치의 불균형을 피하려면 파티셔닝 전략을 결정하기 전에 소스 데이터의 값 분포를 검토하세요.

범위 파티셔닝은 multi-value dimension에서는 불가능합니다. 제공된 partitionDimension이 multi-value이면 수집 작업이 오류를 보고할 것입니다.

| Property | Description | Default | Required | | type | 값을 single_dim으로 설정 | none | yes | | partitionDimension | 파티셔닝할 dimension. 단일 dimension 값을 가진 행만 허용됨 | none | yes | | targetRowsPerSegment | 파티션에 포함할 목표 행 수. 500MB~1GB 세그먼트를 목표로 하는 숫자여야 함 | none | 이 값 또는 maxRowsPerSegment 둘 중 하나 | | maxRowsPerSegment | 파티션에 포함할 행 수의 소프트 최대값 | none | 이 값 또는 targetRowsPerSegment 둘 중 하나 | | assumeGrouped | 입력 데이터가 이미 시간과 dimension으로 그룹화돼 있다고 가정. 가정이 위반되면 수집이 더 빠르지만 최적이 아닌 파티션을 선택할 수 있음 | false | no |

single-dim 파티셔닝에서 Parallel task는 3단계로 실행됩니다: partial dimension distribution, partial segment generation, partial segment merge. 첫 단계는 최상의 파티셔닝을 찾기 위한 통계를 수집하고, 나머지 2단계는 hash 기반 파티셔닝에서처럼 각각 partial 세그먼트를 만들고 병합합니다.

partial dimension distribution 단계에서 Parallel task는 split hint spec에 따라 입력 데이터를 분할하고 worker 태스크에 할당합니다. 각 worker 태스크(타입 partial_dimension_distribution)는 할당된 분할을 읽고 partitionDimension에 대한 히스토그램을 만듭니다. Parallel task는 worker 태스크에서 히스토그램을 수집하고 partitionDimension에 기반해 파티션들 사이에 행을 고르게 분배하는 최상의 범위 파티셔닝을 찾습니다. targetRowsPerSegment 또는 maxRowsPerSegment 중 하나가 최상의 파티셔닝을 찾는 데 사용됩니다.

partial segment generation 단계에서 Parallel task는 새 worker 태스크(타입 partial_range_index_generate)를 생성해서 파티셔닝된 데이터를 만듭니다. 각 worker 태스크는 이전 단계에서 만든 분할을 읽고, granularitySpec의 시간 청크 segmentGranularity(기본 파티션 키)로, 그다음 이전 단계에서 찾은 범위 파티셔닝으로 행을 파티셔닝합니다. 파티셔닝된 데이터는 Middle Manager 또는 Indexer의 로컬 스토리지에 저장됩니다.

partial segment merge 단계에서 parallel index 태스크는 새 worker 태스크 집합(타입 partial_index_generic_merge)을 생성해서 이전 단계에서 만든 파티셔닝된 데이터를 병합합니다. 여기서 파티셔닝된 데이터는 시간 청크와 partitionDimension의 값에 따라 셔플됩니다. 각 worker 태스크는 여러 Middle Manager/Indexer 프로세스에서 같은 범위의 같은 파티션에 들어가는 세그먼트를 읽고 병합해서 최종 세그먼트를 만듭니다. 마지막으로 최종 세그먼트를 딥 스토리지로 푸시합니다.

정보

단일 dimension 범위 파티셔닝 태스크는 partial dimension distribution과 partial segment generation 단계에서 입력을 두 번 통과하기 때문에, 두 통과 사이에 입력이 변경되면 태스크가 실패할 수 있습니다.

다중 dimension 범위 파티셔닝

정보

다중 dimension 범위 파티셔닝은 index_parallel 태스크 타입의 순차 모드에서 지원되지 않습니다.

범위 파티셔닝은 저장 공간과 쿼리 성능과 관련해 여러 이점이 있습니다. 다중 dimension 범위 파티셔닝은 Druid가 세그먼트 크기를 더 고르게 분배하고 더 많은 dimension에서 프루닝할 수 있게 해서 단일 dimension 범위 파티셔닝을 개선합니다.

범위 파티셔닝은 multi-value dimension에서는 불가능합니다. 제공된 partitionDimensions 중 하나가 multi-value이면 수집 작업이 오류를 보고할 것입니다.

| Property | Description | Default | Required | | type | 값을 range로 설정 | none | yes | | partitionDimensions | 파티셔닝할 dimension 배열. 가장 자주 쿼리되는 dimension부터 가장 적게 쿼리되는 dimension 순서로 정렬. 최상의 결과를 위해 dimension 수를 35개로 제한할 것 | none | yes | | targetRowsPerSegment | 파티션에 포함할 목표 행 수. 500MB1GB 세그먼트를 목표로 하는 숫자여야 함 | none | 이 값 또는 maxRowsPerSegment 둘 중 하나 | | maxRowsPerSegment | 파티션에 포함할 행 수의 소프트 최대값 | none | 이 값 또는 targetRowsPerSegment 둘 중 하나 | | assumeGrouped | 입력 데이터가 이미 시간과 dimension으로 그룹화돼 있다고 가정. 가정이 위반되면 수집이 더 빠르지만 최적이 아닌 파티션을 선택할 수 있음 | false | no |

범위 파티셔닝의 이점

single_dim 또는 range 범위 파티셔닝은 여러 이점이 있습니다.

  • 비슷한 데이터를 같은 세그먼트에 결합해 압축률을 개선하기 때문에 저장 공간이 더 작음
  • 쿼리 필터와 일치하는 데이터를 절대 포함할 수 없는 세그먼트는 고려에서 제거하는 Broker-level 세그먼트 프루닝 덕분에 쿼리 성능이 더 좋음

Broker-level 세그먼트 프루닝이 효과적이려면 WHERE 절에 파티션 dimension을 포함해야 해요. 이전 파티션 dimension들(왼쪽에 있는 것들)이 참여하고 있고, 쿼리가 프루닝을 지원하는 필터를 사용한다면 각 파티션 dimension이 프루닝에 참여할 수 있습니다.

프루닝을 지원하는 필터에는 다음이 포함됩니다.

  • 문자열 리터럴에 대한 동등성, 예: x = 'foo'와 x IN ('foo', 'bar')(여기서 x는 문자열)
  • 문자열 컬럼과 문자열 리터럴 사이의 비교, 예: x < 'foo' 또는 <, >, <=, >=를 포함한 다른 비교

예를 들어 수집 중 다음 range 파티셔닝을 구성했다면:

"partitionsSpec": {  "type": "range",  "partitionDimensions": ["countryName", "cityName"],  "targetRowsPerSegment": 5000000}

WHERE countryName = 'United States' 또는 WHERE countryName = 'United States' AND cityName = 'New York' 같은 필터는 프루닝을 활용할 수 있습니다. 그러나 WHERE cityName = 'New York'는 countryName이 관여하지 않으므로 프루닝을 활용할 수 없어요. WHERE cityName LIKE 'New%' 절도 LIKE 필터가 프루닝을 지원하지 않으므로 프루닝을 활용할 수 없습니다.

HTTP 상태 엔드포인트

Supervisor 태스크는 실행 상태를 얻을 수 있는 몇 가지 HTTP 엔드포인트를 제공합니다.

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/mode — indexing 태스크가 병렬로 실행 중이면 parallel을 반환합니다. 그렇지 않으면 sequential을 반환합니다.

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/phase — 태스크가 병렬 모드로 실행 중이면 현재 단계의 이름을 반환합니다.

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/progress — supervisor 태스크가 병렬 모드로 실행 중이면 현재 단계의 추정 진행률을 반환합니다.

결과의 예는 다음과 같습니다.

{  "running":10,  "succeeded":0,  "failed":0,  "complete":0,  "total":10,  "estimatedExpectedSucceeded":10}

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/subtasks/running — 실행 중인 worker 태스크의 태스크 ID를 반환하거나, supervisor 태스크가 순차 모드로 실행 중이면 빈 목록을 반환합니다.

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/subtaskspecs — 모든 worker 태스크 스펙을 반환하거나, supervisor 태스크가 순차 모드로 실행 중이면 빈 목록을 반환합니다.

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/subtaskspecs/running — 실행 중인 worker 태스크 스펙을 반환하거나, supervisor 태스크가 순차 모드로 실행 중이면 빈 목록을 반환합니다.

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/subtaskspecs/complete — 완료된 worker 태스크 스펙을 반환하거나, supervisor 태스크가 순차 모드로 실행 중이면 빈 목록을 반환합니다.

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/subtaskspec/{SUB_TASK_SPEC_ID} — 주어진 id의 worker 태스크 스펙을 반환하거나, supervisor 태스크가 순차 모드로 실행 중이면 HTTP 404 Not Found 오류를 반환합니다.

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/subtaskspec/{SUB_TASK_SPEC_ID}/state — 주어진 id의 worker 태스크 스펙 상태를 반환하거나, supervisor 태스크가 순차 모드로 실행 중이면 HTTP 404 Not Found 오류를 반환합니다. 반환된 결과는 worker 태스크 스펙, 있으면 현재 태스크 상태, 태스크 시도 이력을 포함합니다.

응답 보기

{  "spec": {    "id": "index_parallel_lineitem_2018-04-20T22:12:43.610Z_2",    "groupId": "index_parallel_lineitem_2018-04-20T22:12:43.610Z",    "supervisorTaskId": "index_parallel_lineitem_2018-04-20T22:12:43.610Z",    "context": null,    "inputSplit": {      "split": "/path/to/data/lineitem.tbl.5"    },    "ingestionSpec": {      "dataSchema": {        "dataSource": "lineitem",        "timestampSpec": {          "column": "l_shipdate",          "format": "yyyy-MM-dd"        },        "dimensionsSpec": {          "dimensions": [            "l_orderkey",            "l_partkey",            "l_suppkey",            "l_linenumber",            "l_returnflag",            "l_linestatus",            "l_shipdate",            "l_commitdate",            "l_receiptdate",            "l_shipinstruct",            "l_shipmode",            "l_comment"          ]        },        "metricsSpec": [          {            "type": "count",            "name": "count"          },          {            "type": "longSum",            "name": "l_quantity",            "fieldName": "l_quantity",            "expression": null          },          {            "type": "doubleSum",            "name": "l_extendedprice",            "fieldName": "l_extendedprice",            "expression": null          },          {            "type": "doubleSum",            "name": "l_discount",            "fieldName": "l_discount",            "expression": null          },          {            "type": "doubleSum",            "name": "l_tax",            "fieldName": "l_tax",            "expression": null          }        ],        "granularitySpec": {          "type": "uniform",          "segmentGranularity": "YEAR",          "queryGranularity": {            "type": "none"          },          "rollup": true,          "intervals": [            "1980-01-01T00:00:00.000Z/2020-01-01T00:00:00.000Z"          ]        },        "transformSpec": {          "filter": null,          "transforms": []        }      },      "io...

http://{PEON_IP}:{PEON_PORT}/druid/worker/v1/chat/{SUPERVISOR_TASK_ID}/subtaskspec/{SUB_TASK_SPEC_ID}/history — 주어진 id의 worker 태스크 스펙의 태스크 시도 이력을 반환하거나, supervisor 태스크가 순차 모드로 실행 중이면 HTTP 404 Not Found 오류를 반환합니다.

세그먼트 푸시 모드 (Segment pushing modes)

parallel task 인덱싱으로 데이터를 수집하는 동안 Druid는 입력 데이터에서 세그먼트를 만들어 푸시합니다. parallel task index는 롤업 타입에 따라 다음 세그먼트 푸시 모드를 지원합니다.

  • 벌크 푸시 모드(Bulk pushing mode): perfect rollup에 사용. Druid는 index task의 맨 끝에서 모든 세그먼트를 푸시함. 그때까지 Druid는 생성된 세그먼트를 index task를 실행하는 서비스의 메모리와 로컬 스토리지에 저장. 벌크 푸시 모드를 활성화하려면 tuning config에서 forceGuaranteedRollup을 true로 설정. 벌크 푸시는 IOConfig의 appendToExisting과 함께 사용할 수 없음.
  • 증분 푸시 모드(Incremental pushing mode): best-effort rollup에 사용. Druid는 인덱싱 태스크가 진행되는 동안 세그먼트를 점진적으로 푸시함. index task는 수집된 총 행 수가 maxTotalRows를 초과할 때까지 태스크를 실행하는 서비스의 메모리·디스크에 데이터를 수집하고 생성된 세그먼트를 저장. 그 시점에 index task는 지금까지 만든 모든 세그먼트를 즉시 푸시하고, 푸시된 세그먼트를 정리한 뒤 남은 데이터 수집을 계속함.

용량 계획 (Capacity planning)

Supervisor 태스크는 현재 사용 가능한 태스크 슬롯 수와 무관하게 최대 maxNumConcurrentSubTasks개의 worker 태스크를 만들 수 있습니다. 결과적으로 동시에 실행될 수 있는 총 태스크 수는 (maxNumConcurrentSubTasks + 1)(Supervisor 태스크 포함)입니다. 이는 태스크 슬롯 총 수(모든 worker의 용량 합)보다 더 클 수도 있음에 유의하세요. maxNumConcurrentSubTasks가 n (사용 가능한 태스크 슬롯)보다 크면 supervisor 태스크에 의해 maxNumConcurrentSubTasks개의 태스크가 생성되지만, 시작되는 것은 n개의 태스크뿐입니다. 나머지는 실행 중인 태스크가 끝날 때까지 pending 상태로 기다립니다.

Parallel Index Task를 스트림 수집과 함께 사용할 때는, 배치 수집이 스트림 수집을 막지 않도록 배치 수집의 최대 용량을 제한할 것을 권장합니다. 동시에 실행할 Parallel Index Task가 t개 있고 배치 수집의 최대 태스크 수를 b로 제한하고 싶다고 가정해 보세요. 그러면 (모든 Parallel Index Task의 maxNumConcurrentSubTasks 합 + t(supervisor 태스크용))이 b보다 작아야 합니다.

다른 태스크보다 우선순위가 높은 태스크가 있다면, 그 maxNumConcurrentSubTasks를 낮은 우선순위 태스크보다 높은 값으로 설정할 수 있습니다. 이는 높은 우선순위 태스크에 더 많은 태스크 슬롯을 할당해서 낮은 우선순위 태스크보다 더 일찍 끝나게 도울 수 있어요.

분할 가능한 input sources (Splittable input sources)

inputSource 객체를 사용해서 index가 데이터를 읽을 위치를 정의해요. 네이티브 parallel task와 simple task만 input source를 지원합니다.

사용 가능한 input source에 대한 자세한 내용은 다음을 참고하세요.

  • S3 input source(s3)는 Amazon S3 스토리지에서 데이터를 읽음
  • Google Cloud Storage input source(gs)는 Google Cloud Storage에서 데이터를 읽음
  • Azure input source(azure)는 Azure Blob Storage와 Azure Data Lake에서 데이터를 읽음
  • HDFS input source(hdfs)는 HDFS 스토리지에서 데이터를 읽음
  • HTTP input Source(http)는 HTTP 서버에서 데이터를 읽음
  • Inline input Source는 웹 콘솔에 붙여 넣는 데이터를 읽음
  • Local input Source(local)는 로컬 스토리지에서 데이터를 읽음
  • Druid input Source(druid)는 Druid datasource에서 데이터를 읽음
  • SQL input Source(sql)는 RDBMS 소스에서 데이터를 읽음

input source를 결합하는 방법은 Combining input source를 참고하세요.

segmentWriteOutMediumFactory

| Property | Type | Description | Required | | type | String | 설명과 사용 가능한 옵션은 Additional Peon Configuration: SegmentWriteOutMediumFactory 참고 | yes |

더 알아보기 (Learn more)