Supervisor

Supervisor

Apache Druid는 supervisor를 사용해서 외부 스트리밍 소스에서 Druid로의 스트리밍 수집을 관리해요. supervisor는 indexing task의 상태를 감독하면서 hand-off를 조정하고, 실패를 관리하며, 확장성·복제 요구사항이 유지되도록 보장합니다. 데이터가 수집된 후 자동 컴팩션(automatic compaction)을 수행하는 데에도 사용할 수 있어요.

이 주제에서는 파티션 안 레코드의 식별자를 지칭할 때 Apache Kafka의 용어인 offset을 사용합니다. Amazon Kinesis를 쓴다면 이에 해당하는 용어는 sequence number예요.

출처: 문서

본문

Supervisor spec

Druid는 흔히 supervisor spec이라 부르는 JSON 명세를 사용해서 스트리밍 수집이나 자동 컴팩션에 사용할 태스크를 정의해요. supervisor spec은 Druid가 외부 스트림이나 Druid 자신의 데이터를 어떻게 소비·처리·인덱싱할지 지정합니다.

supervisor spec의 상위 수준 구성 옵션은 다음 표와 같아요.

| Property | Type | Description | Required | | id | String | supervisor id. supervisor를 식별하게 해 줄 고유 ID여야 함. 지정하지 않으면 기본값은 spec.dataSchema.dataSource | No | | type | String | supervisor 타입. 스트리밍 수집에서는 kafka, kinesis, rabbit 중 하나. 자동 컴팩션에서는 타입을 autocompact로 설정 | Yes | | spec | Object | supervisor 구성의 컨테이너 객체. 자동 컴팩션에서는 컴팩션 구성과 동일 | Yes | | spec.dataSchema | Object | 수집 중 indexing task가 사용할 스키마. 자세한 내용은 dataSchema 참고 | Yes | | spec.ioConfig | Object | supervisor와 indexing task의 연결·I/O 관련 설정을 정의하는 I/O 구성 객체 | Yes | | spec.tuningConfig | Object | supervisor와 indexing task의 성능 관련 설정을 정의하는 튜닝 구성 객체 | No | | context | Object | supervisor와 supervisor가 만드는 태스크 양쪽의 추가 구성을 허용 | No | | suspended | Boolean | supervisor를 일시 중지(suspended) 상태로 둠 | No |

I/O 구성

다음 표는 Apache Kafka와 Amazon Kinesis 수집 방법 양쪽에 적용되는 ioConfig 구성 속성을 정리한 것이에요. Kafka와 Kinesis에 특화된 속성은 각각 Kafka I/O configuration과 Kinesis I/O configuration을 참고하세요.

| Property | Type | Description | Required | Default | | inputFormat | Object | 입력 데이터 파싱을 정의하는 input format | Yes | | | autoScalerConfig | Object | 수집 태스크의 자동 확장(autoscaling) 동작 정의. 자세한 내용은 Task autoscaler 참고 | No | null | | taskCount | Integer | replica set 안의 읽기 태스크 최대 수. taskCount와 replicas를 곱하면 읽기 태스크의 최대 수를 구할 수 있음. 읽기·게시 태스크를 합친 전체 태스크 수는 최대 읽기 태스크 수보다 큼. 자세한 내용은 Capacity planning 참고. taskCount가 Kafka 파티션 또는 Kinesis shard 수보다 크면 실제 읽기 태스크 수는 taskCount 값보다 작음 | No | 1 | | replicas | Integer | replica set 수. 1은 단일 태스크 집합(복제 없음). Druid는 항상 태스크 복제본을 다른 worker에 할당해서 프로세스 실패에 대한 회복탄력성을 제공. 태스크 복제본에 서버 우선순위를 할당하려면 serverPriorityToReplicas 참고 | No | 1 | | taskDuration | ISO 8601 period | 태스크가 읽기를 멈추고 세그먼트 게시를 시작하기까지의 시간 길이 | No | PT1H | | startDelay | ISO 8601 period | supervisor가 태스크 관리를 시작하기 전에 기다리는 기간 | No | PT5S | | period | ISO 8601 period | supervisor가 관리 로직을 얼마나 자주 실행할지 결정. supervisor는 또한 태스크 성공·실패·task duration 도달 같은 특정 이벤트에 반응해서 실행됨. period 값은 반복 사이의 최대 시간을 지정 | No | PT30S | | completionTimeout | ISO 8601 period | 게시 태스크를 실패로 선언하고 종료하기 전에 기다리는 시간 길이. 값이 너무 낮으면 태스크가 게시를 못 할 수도 있음. 태스크의 게시 클록(clock)은 대략 taskDuration이 경과한 뒤 시작 | No | PT30M | | lateMessageRejectionStartDateTime | ISO 8601 date time | 이 날짜-시간보다 이른 타임스탬프의 메시지를 거부하도록 태스크 구성. 예를 들어 이 속성이 2016-01-01T11:00Z로 설정되고 supervisor가 2016-01-01T12:00Z에 태스크를 만들면, Druid는 2016-01-01T11:00Z보다 이른 타임스탬프의 메시지를 버림. 데이터 스트림에 늦은 메시지가 있고 같은 세그먼트에서 동작해야 하는 파이프라인이 여러 개(예: 실시간과 야간 배치 수집 파이프라인) 있다면 동시성 문제를 예방할 수 있음 | No | | | lateMessageRejectionPeriod | ISO 8601 period | 태스크가 생성되기 전의 이 기간보다 이른 타임스탬프의 메시지를 거부하도록 태스크 구성. 예를 들어 이 속성이 PT1H로 설정되고 supervisor가 2016-01-01T12:00Z에 태스크를 만들면, Druid는 2016-01-01T11:00Z보다 이른 타임스탬프의 메시지를 버림. 데이터 스트림에 늦은 메시지가 있고 같은 세그먼트에서 동작해야 하는 파이프라인이 여러 개(예: 스트리밍과 야간 배치 수집 파이프라인) 있다면 동시성 문제를 예방하는 데 도움이 될 수 있음. 늦은 메시지 거부 속성은 하나만 지정할 수 있음 | No | | | earlyMessageRejectionPeriod | ISO 8601 period | 태스크가 task duration에 도달한 후 이 기간보다 늦은 타임스탬프의 메시지를 거부하도록 태스크 구성. 예를 들어 이 속성이 PT1H로 설정되고 task duration이 PT1H로 설정되며 supervisor가 2016-01-01T12:00Z에 태스크를 만들면, Druid는 2016-01-01T14:00Z보다 늦은 타임스탬프의 메시지를 버림. 태스크는 supervisor 장애 조치(failover) 같은 경우에 task duration을 넘겨 실행되기도 함 | No | | | stopTaskCount | Integer | Druid가 한 번에 순환(cycle)할 수 있는 수집 태스크 수를 제한. 설정하지 않으면 Druid는 모든 태스크를 동시에 순환할 수 있음. taskCount보다 작은 값으로 설정하면 supervisor를 실행하는 데 필요한 가용 슬롯이 더 적어짐. 수집 티어를 축소해서 비용을 절약할 수 있지만, 이는 더 느린 순환 시간과 lag를 초래할 수 있음. 자세한 내용은 stopTaskCount 참고 | No | taskCount 값 | | serverPriorityToReplicas | Object ( Map<Integer, Integer> ) | 서버 우선순위별 복제본 수 매핑. 설정하면 각 태스크 복제본에 Peon 프로세스의 druid.server.priority에 해당하는 서버 우선순위가 할당되어, query routing strategies로 혼합 워크로드에 대해 쿼리 격리를 활성화함. 구성하지 않으면 replicas 설정이 적용되고 모든 태스크 복제본에 기본 우선순위 0이 할당됨. | No | null |

예를 들어 serverPriorityToReplicas를 {"1": 2, "0": 1}로 설정하면 태스크 그룹당 druid.server.priority=1인 태스크 복제본 2개와 druid.server.priority=0인 태스크 복제본 1개가 만들어져요. 이 구성은 taskCount에 비례해서 확장됩니다. 예를 들어 taskCount가 5로 설정되면 총 15개의 태스크가 생성됩니다 — 우선순위 1 태스크 10개와 우선순위 0 태스크 5개. replicas와 serverPriorityToReplicas 둘 다 설정하면 serverPriorityToReplicas의 복제본 합이 replicas와 같아야 해요.

Task autoscaler

ioConfig 객체의 autoScalerConfig 속성으로 수집 태스크의 자동 확장 동작을 선택적으로 구성할 수 있어요.

autoScalerConfig 구성 속성은 다음 표와 같습니다.

| Property | Description | Required | Default | | enableTaskAutoScaler | autoscaler 활성화. 지정하지 않으면 autoScalerConfig가 null이 아니어도 Druid는 autoscaler를 비활성화함 | No | false | | taskCountMax | 수집 태스크의 최대 수. taskCountMin보다 크거나 같아야 함. taskCountMax가 Kafka 파티션 또는 Kinesis shard 수보다 크면 Druid는 최대 읽기 태스크 수를 Kafka 파티션 또는 Kinesis shard 수로 설정하고 taskCountMax를 무시함 | Yes | | | taskCountMin | 수집 태스크의 최소 수. autoscaler를 활성화하면 Druid는 다음 순서로 config를 확인해서 시작 태스크 수를 계산함: taskCountStart, 그다음 taskCount(ioConfig의), 그다음 taskCountMin | Yes | | | taskCountStart | 시작할 수집 태스크 수를 지정하는 선택적 config. autoscaler를 활성화하면 Druid는 다음 순서로 config를 확인해서 시작 태스크 수를 계산함: taskCountStart, 그다음 taskCount(ioConfig의), 그다음 taskCountMin | No | taskCount 또는 taskCountMin | | minTriggerScaleActionFrequencyMillis | 두 스케일 동작 사이의 최소 시간 간격 | No | 600000 | | autoScalerStrategy | autoscaler의 알고리즘. Druid는 lagBased 전략만 지원. 자세한 내용은 Autoscaler strategy 참고 | No | lagBased | | stopTaskCountRatio | 유효 범위 (0.0, 1.0]인 ioConfig.stopTaskCount의 변수 버전. 정상 상태에서 중지 가능한 태스크 최대 수가 현재 실행 중인 태스크 수에 비례하도록 허용 | No | |

Autoscaler strategy
  1. Lag 기반 autoscaler 전략

정보

Kafka indexing service와 달리, Kinesis는 lag 메트릭을 메시지 수가 아니라 현재 시퀀스 번호와 최신 시퀀스 번호 사이의 밀리초 시간 차이로 보고합니다.

lagBased autoscaler 전략 관련 구성 속성은 다음 표와 같아요.

| Property | Description | Required | Default | | lagCollectionIntervalMillis | Druid가 lag 메트릭 포인트를 수집하는 시간 주기 | No | 30000 | | lagCollectionRangeMillis | lag 수집의 전체 시간 윈도우. lagCollectionIntervalMillis와 함께 사용해서 lag 메트릭 포인트를 수집할 간격을 지정 | No | 600000 | | scaleOutThreshold | 확장(scale out) 동작의 임계값 | No | 6000000 | | triggerScaleOutFractionThreshold | lag 포인트의 triggerScaleOutFractionThreshold 퍼센트가 scaleOutThreshold보다 높으면 확장 동작 활성화 | No | 0.3 | | scaleInThreshold | 축소(scale in) 동작의 임계값 | No | 1000000 | | triggerScaleInFractionThreshold | lag 포인트의 triggerScaleInFractionThreshold 퍼센트가 scaleOutThreshold보다 낮으면 축소 동작 활성화 | No | 0.9 | | scaleActionStartDelayMillis | supervisor가 시작된 후 첫 스케일 로직 검사 전까지 지연할 밀리초 수 | No | 300000 | | scaleActionPeriodMillis | 스케일 동작이 트리거됐는지 확인하는 빈도(밀리초) | No | 60000 | | scaleInStep | 축소할 때 한 번에 줄일 태스크 수 | No | 1 | | scaleOutStep | 확장할 때 한 번에 추가할 태스크 수 | No | 2 | | lagAggregate | 스케일링 결정을 위해 lag 메트릭을 계산하는 데 사용하는 집계 함수. 가능한 값은 MAX, SUM, AVERAGE | No | SUM |

다음은 lagBased autoscaler가 있는 supervisor spec 예시예요.

예시 보기

{  "type": "kinesis",  "dataSchema": {    "dataSource": "metrics-kinesis",    "timestampSpec": {      "column": "timestamp",      "format": "auto"    },    "dimensionsSpec": {      "dimensions": [],      "dimensionExclusions": [        "timestamp",        "value"      ]    },    "metricsSpec": [      {        "name": "count",        "type": "count"      },      {        "name": "value_sum",        "fieldName": "value",        "type": "doubleSum"      },      {        "name": "value_min",        "fieldName": "value",        "type": "doubleMin"      },      {        "name": "value_max",        "fieldName": "value",        "type": "doubleMax"      }    ],    "granularitySpec": {      "type": "uniform",      "segmentGranularity": "HOUR",      "queryGranularity": "NONE"    }  },  "ioConfig": {    "stream": "metrics",    "autoScalerConfig": {      "enableTaskAutoScaler": true,      "taskCountMax": 6,      "taskCountMin": 2,      "minTriggerScaleActionFrequencyMillis": 600000,      "autoScalerStrategy": "lagBased",      "lagCollectionIntervalMillis": 30000,      "lagCollectionRangeMillis": 600000,      "scaleOutThreshold": 600000,      "triggerScaleOutFractionThreshold": 0.3,      "scaleInThreshold": 100000,      "triggerScaleInFractionThreshold": 0.9,      "scaleActionStartDelayMillis": 300000,      "scaleActionPeriodMillis": 60000,      "scaleInStep": 1,      "scaleOutStep": 2    },    "inputFormat": {      "type": "json"    },    "endpoint": "kinesis.us-east-1.amazonaws.com",    "taskCount": 1,    "replicas": 1,    "taskDuration": "PT1H"  },  "tuningConfig": {    "type": "kinesis",    "maxRowsPerSegment": 5000000  }}
  1. 비용 기반 autoscaler 전략 (experimental)

수집 lag와 poll-to-idle 비율에 기반한 비용 함수로 필요한 supervisor 태스크 수를 계산하는 autoscaler예요. 태스크 수는 분할 수의 인자·약수에서 엄격히 도출되는 게 아니라, 현재 파티션-대-태스크 비율에서 파생된 제한된 범위에서 선택됩니다. 이 제한된 파티션-당-태스크 윈도우는 큰 점프를 피하면서 점진적인 스케일링을 가능하게 하고, 필요할 때 약수가 아닌 태스크 수도 허용합니다.

이것은 실험적이며 구현 세부 사항과 비용 함수 파라미터는 변경될 수 있어요.

참고: Kinesis는 아직 지원되지 않으며 지원이 진행 중입니다.

costBased autoscaler 전략 관련 구성 속성은 다음 표와 같아요.

| Property | Description | Required | Default | | scaleActionPeriodMillis | 스케일 동작이 트리거됐는지 확인하는 빈도(밀리초) | No | 600000 | | lagWeight | 비용 함수에서 추출된 lag 값의 가중치 | No | 0.25 | | idleWeight | 비용 함수에서 추출된 poll idle 값의 가중치 | No | 0.75 | | useTaskCountBoundaries | 태스크 수를 선택할 때 제한된 파티션-대-태스크 윈도우 활성화 | No | false | | highLagThreshold | 0보다 큰 값으로 설정하면 버스트 스케일업을 트리거하는 평균 파티션 lag 임계값. 음수로 설정하면 버스트 스케일업 비활성화 | No | -1 | | minScaleDownDelay | 성공적인 스케일 동작 사이의 최소 기간. ISO-8601 기간 문자열로 지정 | No | PT30M | | scaleDownDuringTaskRolloverOnly | 태스크 축소가 태스크 롤오버 중에만으로 제한되는지 여부 | No | false |

다음은 lagBased autoscaler가 있는 supervisor spec 예시예요.

예시 보기

{  "ioConfig": {    "stream": "metrics",    "autoScalerConfig": {      "enableTaskAutoScaler": true,      "autoScalerStrategy": "costBased",      "taskCountMin": 1,      "taskCountMax": 10,      "minTriggerScaleActionFrequencyMillis": 600000,      "lagWeight": 0.1,      "idleWeight": 0.9,    }  }}

stopTaskCount

stopTaskCount를 설정하기 전에 다음을 주의하세요.

  • 일부 작업은 모든 태스크가 동시에 순환해야 합니다. 예를 들어 supervisor spec 변경과 Kafka 파티션 수 변경이 그렇습니다. 이런 작업은 충분한 태스크 슬롯 용량이 없으면 lag를 유발할 수 있어요.
  • 태스크 autoscaler는 태스크 수 변경에 반응해 태스크를 종료할 때 stopTaskCount를 무시합니다. autoscaler는 파티션을 태스크들 사이에 재분배해야 하며, 그러려면 모든 태스크가 종료되어야 해요.
  • stopTaskCount를 taskCount보다 작은 값으로 설정하면 Druid는 가장 오래 실행된 태스크부터 순환하고, 그다음 설정된 값까지 다른 태스크를 순환합니다.

Tuning 구성

tuningConfig 객체는 선택 사항이에요. tuningConfig 객체를 지정하지 않으면 Druid는 기본 구성 설정을 사용합니다.

다음 표는 Kafka와 Kinesis 수집 방법 양쪽에 적용되는 tuningConfig 구성 속성을 정리한 것이에요. Kafka와 Kinesis에 특화된 속성은 각각 Kafka tuning configuration과 Kinesis tuning configuration을 참고하세요.

| Property | Type | Description | Required | Default | | type | String | 수집 방법의 튜닝 타입 코드. kafka 또는 kinesis 중 하나 | Yes | | | maxRowsInMemory | Integer | persist 전에 누적할 행 수. 이 숫자는 사후 집계(post-aggregation) 행을 나타냄. 입력 이벤트 수와 동일하지 않고, 결과적으로 집계된 행 수임. Druid는 maxRowsInMemory로 필요한 JVM 힙 크기를 관리. 인덱싱의 최대 힙 메모리 사용량은 maxRowsInMemory * (2 + maxPendingPersists). 보통 설정할 필요는 없지만, 데이터 성격에 따라 행이 바이트로 짧다면 백만 행을 메모리에 저장하고 싶지 않을 수 있고 이때 값을 설정해야 함 | No | 150000 | | maxBytesInMemory | Long | persist 전에 힙 메모리에 누적할 바이트 수. 메모리 사용의 대략적 추정치에 기반하고 실제 사용량은 아님. 보통 Druid가 내부적으로 계산. 인덱싱의 최대 힙 메모리 사용량은 maxBytesInMemory * (2 + maxPendingPersists) | No | 최대 JVM 메모리의 1/6 | | skipBytesInMemoryOverheadCheck | Boolean | maxBytesInMemory 계산은 수집 중 생성되는 오버헤드 객체와 각 중간 persist를 고려함. 이 오버헤드 객체들의 바이트를 maxBytesInMemory 검사에서 제외하려면 skipBytesInMemoryOverheadCheck를 true로 설정 | No | false | | maxRowsPerSegment | Integer | 세그먼트에 저장할 행 수. 이 숫자는 사후 집계 행임. maxRowsPerSegment 또는 maxTotalRows에 도달하거나 매 intermediateHandoffPeriod마다, 어느 것이든 먼저 일어날 때 hand-off가 발생 | No | 5000000 | | maxTotalRows | Long | 모든 세그먼트에 걸쳐 집계할 행 수. 이 숫자는 사후 집계 행임. maxRowsPerSegment 또는 maxTotalRows에 도달하거나 매 intermediateHandoffPeriod마다, 어느 것이든 먼저 일어날 때 hand-off가 발생 | No | 20000000 | | intermediateHandoffPeriod | ISO 8601 period | 태스크가 세그먼트를 얼마나 자주 hand-off 하는지 결정하는 기간. maxRowsPerSegment 또는 maxTotalRows에 도달하거나 매 intermediateHandoffPeriod마다, 어느 것이든 먼저 일어날 때 hand-off 발생 | No | P2147483647D | | intermediatePersistPeriod | ISO 8601 period | 중간 persist가 발생하는 비율을 결정하는 기간 | No | PT10M | | maxPendingPersists | Integer | pending 상태로 시작되지 않은 채 대기할 수 있는 persist 최대 수. 새 중간 persist가 이 한도를 초과하면 Druid는 현재 실행 중인 persist가 끝날 때까지 수집을 블록함. 하나의 persist가 수집과 동시에 실행될 수 있고, 큐에 쌓인 것은 없음. 인덱싱의 최대 힙 메모리 사용량은 maxRowsInMemory * (2 + maxPendingPersists) | No | 0 | | indexSpec | Object | 인덱싱 시점에 사용할 세그먼트 저장 형식 옵션 정의. 자세한 내용은 IndexSpec 참고 | No | | | indexSpecForIntermediatePersists | Object | 인덱싱 시점에 중간 persisted 임시 세그먼트에 사용할 세그먼트 저장 형식 옵션 정의. indexSpecForIntermediatePersists로 중간 세그먼트의 dimension/metric 압축을 비활성화해서 최종 병합에 필요한 메모리를 줄일 수 있음. 다만 중간 세그먼트의 압축을 끄면 최종 세그먼트로 병합되기 전 사용 중에 페이지 캐시 사용이 늘어날 수 있음 | No | | | reportParseExceptions | Boolean | DEPRECATED. true면 Druid가 파싱 중 발생한 예외를 던져 수집을 중단시킴. false면 Druid가 파싱할 수 없는 행·필드를 건너뜀. reportParseExceptions를 true로 설정하면 maxParseExceptions와 maxSavedParseExceptions의 기존 구성을 덮어써서 maxParseExceptions를 0으로, maxSavedParseExceptions를 1 이하로 제한 | No | false | | handoffConditionTimeout | Long | 세그먼트 hand-off를 기다리는 밀리초 수. >= 0 값으로 설정, 0은 무기한 기다림을 의미 | No | Kafka는 900000 (15분). Kinesis는 0 | | resetOffsetAutomatically | Boolean | offset을 사용할 수 없을 때 파티션을 재설정. true로 설정하면 Druid는 useEarliestOffset 또는 useEarliestSequenceNumber 값에 따라 파티션을 가장 이른 또는 최신 offset으로 재설정(true면 가장 이른, false면 최신). false로 설정하면 Druid가 예외를 표면화해서 태스크가 실패하고 수집이 중단됨. 이런 일이 발생하면 supervisor 재설정을 통해 상황을 바로잡기 위한 수동 개입이 필요할 수 있음 | No | false | | workerThreads | Integer | supervisor가 worker 태스크의 요청/응답과 다른 내부 비동기 작업을 처리하는 데 사용하는 스레드 수 | No | min(10, taskCount) | | chatRetries | Integer | Druid가 indexing task에 대한 HTTP 요청을 태스크가 응답하지 않는 것으로 간주하기 전에 재시도하는 횟수 | No | 8 | | httpTimeout | ISO 8601 period | indexing task로부터 HTTP 응답을 기다리는 시간 | No | PT10S | | shutdownTimeout | ISO 8601 period | supervisor가 종료 전에 태스크의 정상 종료(graceful shutdown)를 시도하는 시간 | No | PT80S | | offsetFetchPeriod | ISO 8601 period | supervisor가 스트리밍 소스와 indexing task를 조회해서 현재 offset을 가져오고 lag를 계산하는 빈도 결정. 사용자 지정 값이 최소값 PT5S보다 낮으면 supervisor는 그 값을 무시하고 최소값을 사용 | No | PT30S | | segmentWriteOutMediumFactory | Object | 세그먼트를 만들 때 사용할 세그먼트 write-out 매체. 설명과 사용 가능한 옵션은 Additional Peon configuration: SegmentWriteOutMediumFactory 참고 | No | 지정하지 않으면 Druid는 druid.peon.defaultSegmentWriteOutMediumFactory.type 값을 사용 | | logParseExceptions | Boolean | true면 파싱 예외 발생 시 오류가 발생한 행 정보를 포함한 오류 메시지를 Druid가 로깅 | No | false | | maxParseExceptions | Integer | 태스크가 수집을 중단하고 실패하기 전에 발생할 수 있는 파싱 예외 최대 수. reportParseExceptions를 설정하면 이 한도를 무시 | No | unlimited | | maxSavedParseExceptions | Integer | 파싱 예외 발생 시 Druid는 가장 최근 파싱 예외들을 추적. maxSavedParseExceptions는 저장된 예외 인스턴스 수를 제한. 저장된 예외는 태스크 완료 후 task completion report에서 확인 가능. reportParseExceptions를 설정하면 이 한도를 무시 | No | 0 | | maxColumnsToMerge | Integer | 게시를 위해 세그먼트를 병합할 때 단일 단계에서 병합할 세그먼트 수의 한도. 이 한도는 병합할 세그먼트 집합에 존재하는 총 컬럼 수에 영향을 줌. 한도를 초과하면 세그먼트 병합이 여러 단계로 발생. Druid는 이 설정과 무관하게 단계당 최소 2개의 세그먼트를 병합 | No | -1 |

supervisor 시작하기

supervisor spec을 제출하면 Druid가 새 supervisor를 시작해요. supervisor spec은 Druid 웹 콘솔의 데이터 로더나 Supervisor API로 제출할 수 있습니다.

다음 스크린샷은 supervisor 두 개가 있는 클러스터의 웹 콘솔 Supervisors 뷰예요.

시작되면 supervisor는 구성된 메타데이터 데이터베이스에 유지됩니다. 여러 supervisor가 같은 datasource로 수집할 수 있어요. 자세한 내용은 Multi-Supervisor Support를 참고하세요. 기존 supervisor ID로 supervisor spec을 제출하면 이전 것을 덮어씁니다.

Overlord가 리더십을 얻을 때(시작으로 인해, 또는 다른 Overlord가 실패한 결과로), 메타데이터 데이터베이스의 각 supervisor spec에 대해 supervisor를 생성해요. 그다음 supervisor는 실행 중인 indexing task를 발견하고, 그것들이 supervisor 구성과 호환되면 채택(adopt)하려고 시도합니다. 호환되지 않으면 태스크는 종료되고 supervisor는 새 태스크 집합을 만듭니다. 이렇게 해서 supervisor 수집 태스크는 Overlord 재시작과 장애 조치를 거쳐도 유지됩니다.

스키마와 구성 변경

스키마나 구성을 변경하려면 새 supervisor spec을 제출해야 해요. Overlord는 기존 supervisor의 정상 종료(graceful shutdown)를 시작합니다. 실행 중인 supervisor는 태스크에 읽기를 멈추고 게시를 시작하라고 신호를 보내고 스스로 종료됩니다. 그다음 Druid는 새 구성을 사용해서 새 supervisor를 만듭니다. Druid는 기존 게시 태스크를 유지하면서 업데이트된 스키마를 제출하고, 이전 태스크 offset에서 새 태스크를 시작해요. 이렇게 해서 수집에 일시 중지 없이 구성 변경을 적용할 수 있습니다.

상태 보고서 (Status report)

supervisor 상태 보고서는 supervisor 태스크의 상태와 recentErrors로 보고되는 최근 발생한 예외 배열을 담고 있어요. 예외의 최대 크기는 druid.supervisor.maxStoredExceptionEvents 구성으로 제어할 수 있습니다.

웹 콘솔에서 supervisor 상태를 보려면 Supervisors 뷰로 이동하고 supervisor ID를 클릭해서 Supervisor 대화상자를 엽니다. 왼쪽 탐색 창에서 Status를 클릭하면 상태가 표시됩니다.

다음은 social_media라는 이름의 supervisor 상태 예시예요.

예시 보기

{  "dataSource": "social_media",  "stream": "social_media",  "partitions": 1,  "replicas": 1,  "durationSeconds": 3600,  "activeTasks": [    {      "id": "index_kafka_social_media_8ff3096f21fe448_jajnddno",      "startingOffsets": {        "0": 0      },      "startTime": "2024-01-30T21:21:41.696Z",      "remainingSeconds": 479,      "type": "ACTIVE",      "currentOffsets": {        "0": 50000      },      "lag": {        "0": 0      }    }  ],  "publishingTasks": [],  "latestOffsets": {    "0": 50000  },  "minimumLag": {    "0": 0  },  "aggregateLag": 0,  "offsetsLastUpdated": "2024-01-30T22:13:19.335Z",  "suspended": false,  "healthy": true,  "state": "RUNNING",  "detailedState": "RUNNING",  "recentErrors": []}

상태 보고서에는 supervisor의 상태에 해당하는 두 가지 속성이 있어요: state와 detailedState. state 속성은 어떤 타입의 supervisor에도 적용되는 소수의 일반 상태를 담습니다. detailedState 속성은 더 설명적인 구현-특정 상태를 담아서 supervisor의 활동에 대한 더 많은 통찰을 제공할 수 있어요.

가능한 state 값은 PENDING, RUNNING, SUSPENDED, STOPPING, UNHEALTHY_SUPERVISOR, UNHEALTHY_TASKS입니다.

detailedState 값과 그에 해당하는 state 매핑은 다음 표와 같아요.

| detailedState | state | Description | | UNHEALTHY_SUPERVISOR | UNHEALTHY_SUPERVISOR | supervisor가 이전 druid.supervisor.unhealthinessThreshold 반복에서 오류를 겪음 | | UNHEALTHY_TASKS | UNHEALTHY_TASKS | 마지막 druid.supervisor.taskUnhealthinessThreshold 태스크가 모두 실패함 | | UNABLE_TO_CONNECT_TO_STREAM | UNHEALTHY_SUPERVISOR | supervisor가 스트림과 연결 문제를 겪고 있으며 과거에 성공적으로 연결한 적이 없음 | | LOST_CONTACT_WITH_STREAM | UNHEALTHY_SUPERVISOR | supervisor가 스트림과 연결 문제를 겪고 있지만 과거에 성공적으로 연결한 적이 있음 | | PENDING (첫 반복에만) | PENDING | supervisor가 초기화됐지만 아직 스트림에 연결을 시작하지 않음 | | CONNECTING_TO_STREAM (첫 반복에만) | RUNNING | supervisor가 스트림에 연결하고 파티션 데이터를 갱신하려 시도 중 | | DISCOVERING_INITIAL_TASKS (첫 반복에만) | RUNNING | supervisor가 이미 실행 중인 태스크를 발견하는 중 | | CREATING_TASKS (첫 반복에만) | RUNNING | supervisor가 태스크를 만들고 상태를 발견하는 중 | | RUNNING | RUNNING | supervisor가 태스크를 시작했고 taskDuration이 경과하기를 기다리는 중 | | IDLE | IDLE | 입력 스트림이 새 데이터를 받지 않았고 모든 기존 데이터가 읽혀서 supervisor가 태스크를 만들지 않음 | | SUSPENDED | SUSPENDED | supervisor가 일시 중지됨 | | STOPPING | STOPPING | supervisor가 중지 중 |

supervisor의 실행 루프 각 반복에서 supervisor는 순서대로 다음 작업을 완료해요.

  • 파티션 목록을 검색하고 각 파티션의 시작 offset을 결정. 계속 진행 중이면 Druid는 마지막으로 처리된 offset을 사용. 새 스트림이면 Druid는 useEarliestOffset 속성에 따라 스트림의 시작 또는 끝에서 시작
  • supervisor의 datasource에 쓰는 실행 중인 indexing task를 발견하고, supervisor 구성과 일치하면 채택하고, 그렇지 않으면 중지하도록 신호를 보냄
  • 각 감독 태스크에 상태 요청을 보내 감독 중인 태스크의 상태 보기를 갱신
  • taskDuration을 초과했고 읽기에서 게시로 전환해야 하는 태스크 처리
  • 게시를 마친 태스크를 처리하고 중복된 replica 태스크에 중지 신호를 보냄
  • 실패한 태스크를 처리하고 supervisor 내부 상태를 정리
  • 정상 태스크 목록을 요청된 taskCount와 replicas 구성과 비교하고 필요하면 추가 태스크를 생성

detailedState 속성은 supervisor가 시작 후 또는 일시 중지에서 재개된 후 이 실행 루프를 처음 실행할 때 추가 값들(위 표에서 "첫 반복에만"으로 표시)을 보여줍니다. 이것은 supervisor가 안정 상태에 도달하지 못하는 초기화-유형 문제를 표면화하기 위한 것입니다. 예를 들어 supervisor가 스트림에 연결할 수 없거나, 스트림에서 읽을 수 없거나, 기존 태스크와 통신할 수 없는 경우입니다. supervisor가 안정적이 되면 — 즉 문제 없이 전체 실행을 완료하면 — detailedState는 중지·일시 중지되거나 실패 임계값에 도달해 비정상 상태로 전환될 때까지 RUNNING 상태를 보여줍니다.

정보

Kafka indexing service에서 supervisor가 Kafka로부터 최신 offset 응답을 받지 못하면 Druid는 파티션별 consumer lag를 음수 값으로 보고할 수 있어요. aggregate lag 값은 항상 >= 0입니다.

SUPERVISORS 시스템 테이블

Druid는 특별한 시스템 스키마를 통해 시스템 정보를 노출해요. sys.supervisors 테이블을 쿼리해서 supervisor 내부에 대한 정보를 검색할 수 있습니다. 다음 예시는 건강 상태로 필터링된 supervisor 태스크 정보를 검색하는 방법을 보여줘요.

SELECT * FROM sys.supervisors WHERE healthy=0;

supervisors 시스템 테이블에 대한 자세한 내용은 SUPERVISORS table을 참고하세요.

supervisor 관리하기

웹 콘솔이나 Supervisor API로 supervisor를 관리할 수 있어요. 웹 콘솔에서는 Supervisors 뷰로 이동해서 Actions 컬럼의 생략 부호(…)를 클릭합니다. 나타난 메뉴에서 원하는 동작을 선택하세요.

이 중 일부 동작은 supervisor가 실행 중일 때만 사용 가능합니다.

일시 중지 (Suspend)

Suspend는 실행 중인 supervisor를 일시 중지합니다. 일시 중지된 supervisor는 계속해서 로그와 메트릭을 내보내요. indexing task는 supervisor를 재개할 때까지 일시 중지된 상태로 유지됩니다. API로 supervisor를 일시 중지하는 방법은 Supervisors: Suspend a running supervisor를 참고하세요.

오프셋 설정 (Set offsets)

정보

이 동작은 메시지가 건너뛰어져 데이터 손실이나 중복 데이터가 발생할 수 있으므로 주의해서 수행하세요.

Set offsets는 supervisor 파티션의 offset을 재설정합니다. 이 동작은 저장된 offset을 지우고 supervisor가 지정된 offset부터 데이터 읽기를 재개하도록 지시해요. 저장된 offset이 없으면 Druid는 지정된 offset을 메타데이터 저장소에 저장합니다. Set offsets는 지정된 파티션에 대한 활성 태스크를 종료하고 재생성해서 재설정된 offset부터 읽기를 시작합니다. 이 작업에서 지정되지 않은 파티션은 supervisor가 마지막으로 저장된 offset부터 재개합니다.

API로 offset을 재설정하는 방법은 Supervisors: Reset offsets for a supervisor를 참고하세요.

하드 리셋 (Hard reset)

정보

이 동작은 메시지가 건너뛰어져 데이터 손실이나 중복 데이터가 발생할 수 있으므로 주의해서 수행하세요.

Hard reset은 supervisor 메타데이터를 지워서, useEarliestOffset 설정에 따라 supervisor가 가장 이른 또는 최신 사용 가능 위치부터 데이터 읽기를 재개하게 합니다. Hard reset은 활성 태스크를 종료·재생성하므로 태스크가 유효한 위치부터 읽기를 시작합니다.

이 동작은 오프셋 누락으로 인한 중지 상태에서 복구하는 데 사용하세요.

API로 supervisor를 재설정하는 방법은 Supervisors: Reset a supervisor를 참고하세요.

종료 (Terminate)

Terminate는 supervisor와 그 indexing task를 중지하고, 그 세그먼트의 게시를 트리거합니다. supervisor를 종료하면 Druid는 메타데이터 저장소에 tombstone 마커를 배치해서 재시작 시 재로딩을 방지합니다. 종료된 supervisor는 여전히 메타데이터 저장소에 존재하며 그 이력은 검색할 수 있습니다.

API로 supervisor를 종료하는 방법은 Supervisors: Terminate a supervisor를 참고하세요.

용량 계획 (Capacity planning)

indexing task는 Middle Manager에서 실행되며 Middle Manager 클러스터에서 사용 가능한 리소스에 의해 제한됩니다. 특히 druid.worker.capacity 속성으로 구성되는 충분한 worker 용량이 있어서 supervisor spec의 구성을 처리할 수 있는지 확인해야 해요. worker 용량은 모든 타입의 indexing task(배치 처리, 스트리밍 태스크, 병합 태스크 같은)에 걸쳐 공유되므로, 총 indexing 부하를 처리하도록 worker 용량을 계획해야 합니다. worker 용량이 바닥나면 indexing task는 큐에 쌓이고 다음 사용 가능한 worker를 기다립니다. 이는 쿼리가 부분 결과를 반환하게 할 수 있지만, 태스크가 스트림이 그 offset을 정리하기 전에 실행된다고 가정하면 데이터 손실은 발생하지 않습니다.

실행 중인 태스크는 읽기(reading) 또는 게시(publishing)의 두 상태 중 하나로 있을 수 있어요. 태스크는 taskDuration에 정의된 기간 동안 읽기 상태로 있다가 게시 상태로 전환됩니다. 태스크는 세그먼트를 생성하고, 세그먼트를 딥 스토리지로 푸시하고, Historical 서비스가 그것을 로드·제공할 때까지, 또는 completionTimeout이 경과할 때까지 게시 상태로 유지됩니다.

읽기 태스크 수는 replicas와 taskCount로 제어됩니다. 일반적으로 replicas * taskCount개의 읽기 태스크가 있어요. taskCount가 Kinesis의 shard 수나 Kafka의 파티션 수를 초과하면 예외가 발생하는데, 이 경우 Druid는 shard 수 또는 파티션 수를 사용합니다. taskDuration이 경과하면 이 태스크들은 게시 상태로 전환되고 replicas * taskCount개의 새 읽기 태스크가 생성됩니다. 읽기 태스크와 게시 태스크가 동시에 실행될 수 있도록 최소 용량은 다음과 같아야 합니다.

workerCapacity = 2 * replicas * taskCount

이 값은 이상적인 상황, 즉 최대 한 집합의 태스크가 게시하는 동안 다른 집합이 읽는 상황에 대한 값이에요. 어떤 상황에서는 여러 집합의 태스크가 동시에 게시할 수도 있습니다. 이는 게시 시간(세그먼트 생성, 딥 스토리지로 푸시, Historical에 로드)이 taskDuration보다 길 때 발생합니다. 이는 타당하고 올바른 시나리오이지만 추가 worker 용량이 필요해요. 일반적으로 이전 태스크 집합이 현재 집합이 시작하기 전에 게시를 마칠 수 있도록 taskDuration을 충분히 크게 두는 것이 좋습니다.

Multi-Supervisor 지원

Druid는 여러 스트림 supervisor가 같은 datasource로 수집하는 것을 지원해요. 즉 언제든지 원하는 수의 스트림 supervisor(Kafka, Kinesis 등)가 같은 datasource로 동시에 수집할 수 있다는 뜻입니다. 여러 supervisor가 있는 수집 태스크 사이의 적절한 동기화를 보장하려면 supervisor spec의 context 필드에 useConcurrentLocks=true를 설정하는 것이 중요합니다. 더 읽기.

더 알아보기 (Learn more)

관련 주제는 다음을 참고하세요.

  • Supervisor API — API로 supervisor를 관리·모니터링하는 방법
  • Apache Kafka ingestion — Apache Kafka 스트림에서 데이터 수집하기
  • Amazon Kinesis ingestion — Amazon Kinesis 스트림에서 데이터 수집하기