RabbitMQ superstream 수집

RabbitMQ superstream 수집

druid-rabbit-indexing-service 커뮤니티 확장을 이용하면 RabbitMQ super-stream에서 데이터를 수집할 수 있어요. Overlord에 supervisor를 구성해 RabbitMQ 인덱싱 태스크의 생성과 수명을 관리하는 구조예요.

출처: 문서

본문

Rabbit 스트림 인덱싱 서비스는 Overlord에 supervisor를 구성해서 RabbitMQ 인덱싱 태스크의 생성과 수명을 관리할 수 있게 해줘요. 이 인덱싱 태스크들은 rabbit super-stream에서 이벤트를 읽어요. supervisor는 인덱싱 태스크의 상태를 감독해서 다음을 수행해요.

  • handoff 조정
  • 실패 관리
  • Druid가 확장성과 복제 요구 사항을 유지하도록 보장

Rabbit 스트림 인덱싱 서비스를 사용하려면 druid-rabbit-indexing-service 커뮤니티 druid 확장을 로드해 주세요. 자세한 내용은 Loading community extensions를 참고해요.

Submitting a supervisor spec

Rabbit 스트림 인덱싱 서비스를 사용하려면 Overlord와 Middle Manager 둘 다에 druid-rabbit-indexing-service 확장을 로드해 주세요. supervisor spec을 제출하면 Druid가 dataSource에 대한 supervisor를 시작해요. supervisor spec을 다음 엔드포인트에 제출해 주세요.

http://<OVERLORD_IP>:<OVERLORD_PORT>/druid/indexer/v1/supervisor

예를 들어:

curl -X POST -H 'Content-Type: application/json' -d @supervisor-spec.json http://localhost:8090/druid/indexer/v1/supervisor

여기서 supervisor-spec.json 파일은 rabbit supervisor spec을 담고 있어요.

{
  "type": "rabbit",
  "spec": {
    "dataSchema": {
      "dataSource": "metrics-rabbit",
      "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",
     "inputFormat": {
       "type": "json"
      },
      "uri": "rabbitmq-stream://localhost:5552",
      "taskCount": 1,
      "replicas": 1,
      "taskDuration": "PT1H"
   },
   "tuningConfig": {
     "type": "rabbit",
     "maxRowsPerSegment": 5000000
   }
  }
}

Supervisor spec

필드 설명 필수
type supervisor 타입. 항상 rabbit이어야 해요. yes
spec supervisor 구성의 컨테이너 오브젝트. yes
dataSchema rabbit 인덱싱 태스크가 수집 중에 사용할 스키마. dataSchema 참고. yes
ioConfig supervisor와 인덱싱 태스크의 rabbit super stream 연결 및 I/O 관련 설정을 구성하는 ioConfig 오브젝트. yes
tuningConfig supervisor와 인덱싱 태스크의 성능 관련 설정을 구성하는 tuningConfig 오브젝트. no

ioConfig

필드 타입 설명 필수
stream String 읽을 RabbitMQ super stream. yes
inputFormat Object 입력 데이터를 파싱하는 방법을 지정하는 입력 포맷. 자세한 내용은 inputFormat 참고. yes
uri String RabbitMQ에 연결할 URI. yes
replicas Integer replica 세트의 수. 1은 단일 태스크 세트(복제 없음)를 의미해요. Replica 태스크는 항상 다른 워커에 할당되어 프로세스 실패에 대한 복원력을 제공해요. no (기본값 == 1)
taskCount Integer replica 세트 안의 최대 읽기 태스크 수. 최대 읽기 태스크 수는 taskCount * replicas가 되고 총 태스크(읽기 + 발행) 수는 이보다 많아요. no (기본값 == 1)
taskDuration ISO8601 Period 태스크가 읽기를 멈추고 세그먼트를 발행하기 시작하기까지의 시간. no (기본값 == PT1H)
startDelay ISO8601 Period supervisor가 태스크 관리를 시작하기 전에 기다리는 시간. no (기본값 == PT5S)
period ISO8601 Period supervisor가 관리 로직을 실행하는 주기. supervisor는 특정 이벤트(태스크 성공, 실패, taskDuration 도달 등)에 대한 응답으로도 실행되므로, 이 값은 반복 사이의 최대 시간을 지정해요. no (기본값 == PT30S)
useEarliestSequenceNumber Boolean supervisor가 dataSource를 처음 관리할 때 RabbitMQ에서 시작 시퀀스 번호 집합을 얻어요. 이 플래그는 스트림에서 가장 이른 시퀀스 번호를 가져올지 최신 시퀀스 번호를 가져올지를 결정해요. 정상 상황에서는 후속 태스크가 이전 세그먼트가 끝난 지점에서 시작하므로 이 플래그는 첫 실행에만 사용돼요. no (기본값 == false)
completionTimeout ISO8601 Period 발행 태스크를 실패로 선언하고 종료하기까지 기다리는 시간. 너무 낮게 설정하면 태스크가 발행하지 못할 수 있어요. 태스크의 발행 시계는 대략 taskDuration이 경과한 후에 시작돼요. no (기본값 == PT6H)
lateMessageRejectionPeriod ISO8601 Period 태스크가 생성되기 전 이 기간보다 이른 timestamp를 가진 메시지를 거부하도록 구성해요. 예를 들어 PT1H로 설정하고 supervisor가 2016-01-01T12:00Z에 태스크를 만들었다면, 2016-01-01T11:00Z보다 이른 timestamp의 메시지는 버려져요. 데이터 스트림에 늦은 메시지가 있고 같은 세그먼트에서 동작해야 하는 파이프라인이 여러 개(예: 실시간과 야간 배치 수집 파이프라인)라면 동시성 문제를 막는 데 도움이 될 수 있어요. no (기본값 == none)
earlyMessageRejectionPeriod ISO8601 Period 태스크가 taskDuration에 도달한 후 이 기간보다 늦은 timestamp의 메시지를 거부하도록 구성해요. 예를 들어 PT1H로 설정되고 taskDuration이 PT1H이며 supervisor가 2016-01-01T12:00Z에 태스크를 만들었다면, 2016-01-01T14:00Z보다 늦은 timestamp의 메시지는 버려져요. 참고: 태스크는 supervisor failover 같은 경우 원래 구성된 task duration을 넘어 실행되기도 해요. earlyMessageRejectionPeriod를 너무 낮게 설정하면 태스크가 원래 구성된 task duration을 넘어 실행될 때마다 메시지가 예기치 않게 버려질 수 있어요. no (기본값 == none)
Consumer Properties Object 제공하는 데 사용되는 동적 map. no (기본값 == none)

tuningConfig

tuningConfig는 선택 항목이에요. tuningConfig를 지정하지 않으면 기본 파라미터가 사용돼요.

필드 타입 설명 필수
type String 인덱싱 태스크 타입. 항상 rabbit이어야 해요. yes
maxRowsInMemory Integer persist하기 전에 집계할 행 수. 이 숫자는 post-aggregation 행이므로 입력 이벤트 수와 동일하지 않고, 그 이벤트들이 만들어내는 집계된 행 수예요. 필요한 JVM 힙 크기를 관리하는 데 사용돼요. 인덱싱의 최대 힙 메모리 사용량은 maxRowsInMemory * (2 + maxPendingPersists)에 비례해요. no (기본값 == 100000)
maxBytesInMemory Long persist하기 전에 힙 메모리에 집계할 바이트 수. 메모리 사용량의 대략적인 추정치에 기반하며 실제 사용량은 아니에요. 보통 내부적으로 계산되므로 사용자가 설정할 필요는 없어요. 인덱싱의 최대 힙 메모리 사용량은 maxBytesInMemory * (2 + maxPendingPersists)예요. no (기본값 == 최대 JVM 메모리의 1/6)
maxRowsPerSegment Integer 세그먼트 하나로 집계할 행 수. 이 숫자는 post-aggregation 행이에요. Handoff는 maxRowsPerSegment 또는 maxTotalRows에 도달하거나 매 intermediateHandoffPeriod마다 중 먼저 발생하는 때 일어나요. no (기본값 == 5000000)
maxTotalRows Long 모든 세그먼트에 걸쳐 집계할 행 수. 이 숫자는 post-aggregation 행이에요. Handoff는 maxRowsPerSegment 또는 maxTotalRows에 도달하거나 매 intermediateHandoffPeriod마다 중 먼저 발생하는 때 일어나요. no (기본값 == unlimited)
intermediatePersistPeriod ISO8601 Period 중간 persist가 발생하는 비율을 결정하는 기간. no (기본값 == PT10M)
maxPendingPersists Integer 보류(pending) 상태이지만 시작되지 않은 persist의 최대 수. 새 중간 persist가 이 한도를 초과하면, 현재 실행 중인 persist가 끝날 때까지 수집이 블록돼요. 인덱싱의 최대 힙 메모리 사용량은 maxRowsInMemory * (2 + maxPendingPersists)에 비례해요. no (기본값 == 0, 즉 persist 하나가 수집과 동시에 실행될 수 있고 대기 중인 것은 없음)
indexSpec Object 데이터가 인덱싱되는 방식을 튜닝해요. 자세한 내용은 IndexSpec 참고. no
indexSpecForIntermediatePersists Object 중간 persist된 임시 세그먼트에 대해 인덱싱 시 사용할 세그먼트 저장 포맷 옵션을 정의해요. 최종 병합에 필요한 메모리를 줄이기 위해 중간 세그먼트에서 dimension/metric 압축을 비활성화하는 데 사용할 수 있어요. 다만 중간 세그먼트에서 압축을 비활성화하면 최종 세그먼트로 병합되기 전에 사용되는 동안 페이지 캐시 사용량이 늘어날 수 있어요. 가능한 값은 IndexSpec 참고. no (기본값 = indexSpec과 동일)
reportParseExceptions Boolean true면 파싱 중 발생한 예외가 던져져서 수집이 중단되고, false면 파싱할 수 없는 행과 필드는 건너뛰어져요. no (기본값 == false)
handoffConditionTimeout Long 세그먼트 handoff를 기다리는 밀리초. 반드시 >= 0이며, 0은 영원히 기다린다는 뜻이에요. no (기본값 == 0)
resetOffsetAutomatically Boolean Druid가 더 이상 사용할 수 없는 RabbitMQ 메시지를 읽어야 할 때의 동작을 제어해요. 지원되지 않아요. no (기본값 == false)
skipSequenceNumberAvailabilityCheck Boolean 특정 RabbitMQ 스트림에서 현재 시퀀스 번호를 사용할 수 있는지 확인을 활성화할지 여부. false로 설정하면 인덱싱 태스크는 resetOffsetAutomatically 값에 따라 현재 시퀀스 번호 재설정을 시도하거나(또는 안 하거나) 해요. no (기본값 == false)
workerThreads Integer supervisor가 워커 태스크의 요청/응답과 그 외 내부 비동기 작업을 처리하는 데 사용하는 스레드 수. no (기본값 == min(10, taskCount))
chatRetries Integer 인덱싱 태스크에 대한 HTTP 요청이 태스크를 응답 없음으로 간주하기 전에 재시도되는 횟수. no (기본값 == 8)
httpTimeout ISO8601 Period 인덱싱 태스크의 HTTP 응답을 기다리는 시간. no (기본값 == PT10S)
shutdownTimeout ISO8601 Period supervisor가 종료 전에 태스크의 정상 종료(graceful shutdown)를 시도하는 시간. no (기본값 == PT80S)
recordBufferSize Integer RabbitMQ consumer와 메인 수집 스레드 사이에서 사용되는 버퍼의 크기(이벤트 수). no (기본값 == 100MB 또는 사용 가능한 힙의 대략 10% 중 더 작은 값)
recordBufferOfferTimeout Integer 타임아웃 전에 버퍼에 공간이 생기기를 기다리는 시간(밀리초). no (기본값 == 5000)
segmentWriteOutMediumFactory Object 세그먼트를 만들 때 사용할 세그먼트 write-out medium. 자세한 내용은 아래 참고. no (기본적으로 지정되지 않으며 druid.peon.defaultSegmentWriteOutMediumFactory.type 값이 사용됨)
intermediateHandoffPeriod ISO8601 Period 태스크가 세그먼트를 hand off해야 하는 주기. Handoff는 maxRowsPerSegment 또는 maxTotalRows에 도달하거나 매 intermediateHandoffPeriod마다 중 먼저 발생하는 때 일어나요. no (기본값 == P2147483647D)
logParseExceptions Boolean true면 파싱 예외가 발생할 때 오류가 발생한 행에 대한 정보를 담은 오류 메시지를 기록해요. no, 기본값 == false
maxParseExceptions Integer 태스크가 수집을 중단하고 실패하기 전에 발생할 수 있는 파싱 예외의 최대 수. reportParseExceptions가 설정되면 덮어써져요. no, 기본값 unlimited
maxSavedParseExceptions Integer 파싱 예외가 발생하면 Druid는 가장 최근의 파싱 예외를 추적할 수 있어요. maxSavedParseExceptions는 Druid가 저장하는 예외 인스턴스 수를 제한해요. 저장된 예외는 태스크 완료 후 태스크 완료 보고서에서 확인할 수 있어요. reportParseExceptions가 설정되면 덮어써져요. no, 기본값 == 0
maxRecordsPerPoll Integer poll당 버퍼에서 가져올 최대 record/event 수. 실제 최대값은 Max(maxRecordsPerPoll, Max(bufferSize, 1)). no, 기본값 = 100
repartitionTransitionDuration ISO8601 Period shard가 분할되거나 병합되면 supervisor는 shard->태스크 그룹 매핑을 다시 계산하고, 이전 매핑 아래에서 생성된 실행 중인 태스크에게 (현재 시간 + repartitionTransitionDuration)에 일찍 중지하라고 신호를 보내요. 태스크를 일찍 중지하면 Druid가 새 shard에서 더 빨리 읽기 시작할 수 있어요. 이 속성이 제어하는 repartition 전환 대기 시간은 스트림에 분할/병합 후 새 shard에 레코드를 쓸 추가 시간을 줘서, https://github.com/apache/druid/issues/7600에 설명된 빈 shard 처리 문제를 피하는 데 도움이 돼요. no, (기본값 == PT2M)
offsetFetchPeriod ISO8601 Period supervisor가 RabbitMQ와 인덱싱 태스크에 쿼리해서 현재 offset을 가져오고 lag를 계산하는 주기. 사용자가 지정한 값이 최소값(PT5S)보다 낮으면 supervisor는 그 값을 무시하고 최소값을 사용해요. no (기본값 == PT30S, 최소 == PT5S)

IndexSpec

필드 타입 설명 필수
bitmap Object 비트맵 인덱스의 압축 포맷. JSON 오브젝트여야 해요. 옵션은 Bitmap types 참고. no (기본값 Roaring)
dimensionCompression String dimension 열의 압축 포맷. LZ4, LZF, uncompressed 중 선택. no (기본값 == LZ4)
metricCompression String 기본 타입 메트릭 열의 압축 포맷. LZ4, LZF, uncompressed, none 중 선택. no (기본값 == LZ4)
longEncoding String long 타입의 메트릭 및 dimension 열의 인코딩 포맷. auto 또는 longs. auto는 열 카디널리티에 따라 시퀀스 번호나 lookup table로 값을 인코딩하고 가변 크기로 저장해요. longs는 값을 그대로 8바이트씩 저장해요. no (기본값 == longs)

Bitmap types

Roaring 비트맵의 경우:

필드 타입 설명 필수
type String 반드시 roaring이어야 해요. yes

Concise 비트맵의 경우:

필드 타입 설명 필수
type String 반드시 concise여야 해요. yes

SegmentWriteOutMediumFactory

필드 타입 설명 필수
type String 설명 및 사용 가능한 옵션은 Additional Peon configuration: SegmentWriteOutMediumFactory 참고. yes

Operations

이 섹션은 Rabbit Stream Indexing Service에서 일부 supervisor API가 어떻게 동작하는지 설명해요. 모든 supervisor API에 대해서는 Supervisor APIs를 확인하세요.

RabbitMQ authentication

RabbitMQ에 안전하게 인증하려면 username과 password를 제공하고, 표준 인증서 제공자를 사용하지 않는다면 인증서도 구성해야 해요.

이들을 구성하려면 ioConfig의 동적 구성 제공자(dynamic configuration provider)를 사용해요.

  "ioConfig": {
    "type": "rabbit",
    "stream": "api-audit",
    "uri": "rabbitmq-stream://localhost:5552",
    "taskCount": 1,
    "replicas": 1,
    "taskDuration": "PT1H",
    "consumerProperties": {
        "druid.dynamic.config.provider" : {
            "type": "environment",
            "variables": {
                "username": "RABBIT_USERNAME",
                "password": "RABBIT_PASSWORD"
            }
        }
    }
  },

더 알아보기 (Learn more)