Apache Kafka 수집
Apache Kafka 수집
정보
Kafka indexing service를 사용하려면 Apache Kafka 버전 0.11.x 이상이어야 해요. 이전 버전을 사용 중이라면 Apache Kafka upgrade guide를 참고하세요.
Kafka indexing service를 활성화하면 Overlord에서 supervisor를 구성해서 Kafka indexing task의 생성과 수명을 관리할 수 있어요. Kafka indexing task는 Kafka 파티션과 offset 메커니즘으로 이벤트를 읽어 정확히 한 번(Exactly-once) 수집을 보장합니다. supervisor는 indexing task의 상태를 감독하면서 hand-off를 조정하고, 실패를 관리하며, 확장성·복제 요구사항이 유지되도록 보장해요.
이 주제는 Apache Druid의 Kafka indexing service supervisor에 대한 구성 정보를 담고 있습니다.
출처: 문서
본문
설정 (Setup)
Kafka indexing service를 사용하려면 먼저 Overlord와 Middle Manager 양쪽에 druid-kafka-indexing-service 확장을 로드해야 해요. 자세한 내용은 Loading extensions를 참고하세요.
Supervisor spec 구성
이 섹션은 Apache Kafka 스트리밍 수집 방법에 특화된 구성 속성을 다룹니다. Druid가 지원하는 모든 스트리밍 수집 방법에 공통으로 적용되는 속성은 Supervisor spec을 참고해 주세요.
다음은 Kafka indexing service의 supervisor spec 예시예요.
예시 보기
{ "type": "kafka", "spec": { "dataSchema": { "dataSource": "metrics-kafka", "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": { "topic": "metrics", "inputFormat": { "type": "json" }, "consumerProperties": { "bootstrap.servers": "localhost:9092" }, "taskCount": 1, "replicas": 1, "taskDuration": "PT1H" }, "tuningConfig": { "type": "kafka", "maxRowsPerSegment": 5000000 } }}
I/O 구성
다음 표는 Kafka에 특화된 ioConfig 구성 속성을 정리한 것이에요. 모든 스트리밍 수집 방법에 공통인 속성은 Supervisor I/O configuration을 참고하세요.
| Property | Type | Description | Required | Default |
| topic | String | 읽을 Kafka 토픽. 참고로 이 값은 supervisor에 한 번 정해지면 업데이트가 지원되지 않음. 여러 토픽에서 수집하려면 topicPattern 사용 | topicPattern이 설정되지 않았다면 Yes | |
| topicPattern | String | 정규식 패턴으로 전달하는, 읽을 여러 Kafka 토픽. 자세한 내용은 Ingest from multiple topics 참고 | topic이 설정되지 않았다면 Yes | |
| consumerProperties | String, Object | Kafka consumer에 전달할 속성 맵. 자세한 내용은 Consumer properties 참고. 최소한 Kafka 클러스터에 대한 초기 연결을 수립하기 위해 bootstrap.servers 속성은 반드시 설정해야 함 | Yes | |
| pollTimeout | Long | Kafka consumer가 레코드를 폴링할 때까지 기다리는 시간(밀리초) | No | 100 |
| useEarliestOffset | Boolean | supervisor가 datasource를 처음 관리할 때 Kafka로부터 시작 offset 집합을 얻음. 이 플래그는 supervisor가 Kafka에서 가장 이른 offset을 가져올지 최신 offset을 가져올지 결정. 정상 상황에서는 이후 태스크들이 이전 세그먼트가 끝난 위치에서 시작하므로 이 플래그는 첫 실행에서만 사용됨 | No | false |
| idleConfig | Object | Kafka supervisor가 언제·어떻게 idle 상태가 될 수 있는지 정의. 자세한 내용은 Idle configuration 참고 | No | null |
여러 토픽에서 수집 (Ingest from multiple topics)
정보
datasource에 대해 멀티-토픽 수집을 활성화하면, 28.0.0보다 이전 버전으로 다운그레이드했을 때 그 datasource의 수집이 실패합니다.
정보
기존 supervisor를 topic 대신 topicPattern을 쓰도록 마이그레이션하는 것은 지원되지 않아요. 기존 supervisor의 topicPattern을 다른 정규식 패턴으로 바꾸는 것도 지원되지 않습니다. 다음과 같이 하면 강제로 마이그레이션할 수 있어요.
- supervisor를 일시 중지(suspend)하세요.
- offset을 리셋하세요.
- 업데이트된 supervisor를 제출하세요.
하나 또는 여러 토픽에서 데이터를 수집할 수 있어요. 여러 토픽에서 수집할 때 Druid는 토픽 이름의 hashcode와 그 토픽 안 파티션의 ID에 기반해 파티션을 할당합니다. 파티션 할당은 모든 태스크에 균일하지 않을 수 있어요. Druid는 개별 토픽의 파티션들이 비슷한 부하를 가진다고 가정합니다. 같은 supervisor에서 높은 부하와 낮은 부하 토픽을 모두 수집하고 싶다면, 높은 부하 토픽에는 파티션 수를 더 많이, 낮은 부하 토픽에는 더 적게 두는 것을 권장해요.
여러 토픽에서 데이터를 수집하려면 topic 대신 topicPattern 속성을 사용하세요. 여러 토픽을 정규식 패턴으로 전달합니다. 예를 들어 clicks와 impressions에서 데이터를 수집하려면 topicPattern을 clicks|impressions로 설정하세요. 마찬가지로 metrics-로 시작하는 모든 토픽에서 수집하려면 topicPattern 값으로 metrics-.*를 사용할 수 있어요. 정규식과 일치하는 새 토픽이 클러스터에 추가되면 Druid는 자동으로 그 새 토픽에서 수집을 시작합니다. my-metrics-12처럼 부분적으로만 일치하는 토픽 이름은 수집에 포함되지 않아요.
Consumer properties
consumer properties는 supervisor가 Kafka 스트림에서 이벤트 메시지를 읽고 처리하는 방식을 제어해요. consumer 구성과 고급 사용 사례에 대한 자세한 내용은 Kafka 문서를 참고하세요.
consumer properties에 <BROKER_1>:<PORT_1>,<BROKER_2>:<PORT_2>,... 형식의 Kafka broker 목록으로 bootstrap.servers를 반드시 포함해야 합니다. 어떤 경우에는 consumer properties를 런타임에 가져와야 할 수도 있어요. 예를 들어 bootstrap.servers가 알려져 있지 않거나 정적이 아닐 때요.
consumerProperties의 isolation.level 속성은 Druid가 트랜잭션적으로 쓰인 메시지를 어떻게 읽을지 결정해요. Druid의 기본값인 read_committed에서는 커밋된 트랜잭션만 읽힙니다. 트랜잭션을 지원하지 않는 이전 Kafka 버전을 쓰거나, 중단된(aborted) 트랜잭션까지 읽고 싶다면 isolation.level을 read_uncommitted로 설정하세요.
Kafka 클러스터가 consumer group ACL을 활성화한 경우, consumerProperties에서 group.id를 설정해서 기본 자동 생성 group ID를 덮어쓸 수 있어요.
SSL 연결을 활성화하려면 keystore, truststore, key의 비밀번호를 안전하게 제공해야 합니다. 이 설정들은 jaas.conf 로그인 구성 파일이나 consumerProperties의 sasl.jaas.config로 지정할 수 있어요. 민감한 정보를 보호하려면 환경 변수 동적 구성 제공자(dynamic config provider)를 사용해서 자격 증명을 평문 대신 시스템 환경 변수에 저장하세요. Kafka 수집용 SSL 구성을 지정하는 데 password provider 인터페이스를 사용할 수도 있지만, 이 기능은 deprecated이므로 동적 구성 제공자를 사용하는 것을 고려하세요.
예를 들어 Kafka에서 SASL과 SSL을 사용할 때, Overlord와 Peon 서비스를 실행하는 머신의 Druid 사용자에 대해 다음 환경 변수를 설정하세요. 값을 자신의 환경 구성에 맞게 바꾸세요.
export KAFKA_JAAS_CONFIG="org.apache.kafka.common.security.plain.PlainLoginModule required username='accesskey' password='secret key';"export SSL_KEY_PASSWORD=mysecretkeypasswordexport SSL_KEYSTORE_PASSWORD=mysecretkeystorepasswordexport SSL_TRUSTSTORE_PASSWORD=mysecrettruststorepassword
supervisor spec에서 consumer properties를 정의할 때 동적 구성 제공자를 사용해서 환경 변수를 참조하세요.
"consumerProperties": { "bootstrap.servers": "localhost:9092", "security.protocol": "SASL_SSL", "sasl.mechanism": "PLAIN", "ssl.keystore.location": "/opt/kafka/config/kafka01.keystore.jks", "ssl.truststore.location": "/opt/kafka/config/kafka.truststore.jks", "druid.dynamic.config.provider": { "type": "environment", "variables": { "sasl.jaas.config": "KAFKA_JAAS_CONFIG", "ssl.key.password": "SSL_KEY_PASSWORD", "ssl.keystore.password": "SSL_KEYSTORE_PASSWORD", "ssl.truststore.password": "SSL_TRUSTSTORE_PASSWORD" } }}
Kafka에 연결할 때 Druid는 환경 변수를 그에 해당하는 값으로 대체합니다.
Idle 구성
정보
idle 상태 전환은 현재 experimental로 지정돼 있어요.
supervisor가 idle 상태에 들어가면, 현재 실행 중인 태스크가 완료된 이후로는 새 태스크가 시작되지 않습니다. 이 전략은 간헐적으로만 데이터가 들어오는 토픽을 쓰는 클러스터 운영자의 비용을 줄여 줄 수 있어요.
idleConfig 구성 옵션은 다음 표와 같아요.
| Property | Description | Required |
| enabled | true면 입력 스트림 또는 토픽에 일정 시간 동안 데이터가 없을 때 supervisor가 idle 상태가 됨 | No |
| inactiveAfterMillis | 입력 토픽의 모든 기존 데이터가 읽히고 inactiveAfterMillis 밀리초 동안 새 데이터가 게시되지 않으면 supervisor가 idle 상태가 됨 | No |
다음 예시는 idle 구성이 활성화된 supervisor spec이에요.
예시 보기
{ "type": "kafka", "spec": { "dataSchema": {...}, "ioConfig": { "topic": "metrics", "inputFormat": { "type": "json" }, "consumerProperties": { "bootstrap.servers": "localhost:9092" }, "autoScalerConfig": { "enableTaskAutoScaler": true, "taskCountMax": 6, "taskCountMin": 2, "minTriggerScaleActionFrequencyMillis": 600000, "autoScalerStrategy": "lagBased", "lagCollectionIntervalMillis": 30000, "lagCollectionRangeMillis": 600000, "scaleOutThreshold": 6000000, "triggerScaleOutFractionThreshold": 0.3, "scaleInThreshold": 1000000, "triggerScaleInFractionThreshold": 0.9, "scaleActionStartDelayMillis": 300000, "scaleActionPeriodMillis": 60000, "scaleInStep": 1, "scaleOutStep": 2 }, "taskCount": 1, "replicas": 1, "taskDuration": "PT1H", "idleConfig": { "enabled": true, "inactiveAfterMillis": 600000 } }, "tuningConfig": {...} }}
데이터 형식
Kafka indexing service는 inputFormat을 지원해요. 자세한 내용은 Source input formats을 참고하세요.
Kafka indexing service는 inputFormat에 대해 다음 값을 지원합니다.
csvtvsjsonkafkaavro_streamprotobufthrift
Kafka input format supervisor spec 예시
kafka input format은 Kafka 페이로드(payload) 값 내용에 더해 Kafka 메타데이터 필드도 파싱할 수 있게 해 줍니다.
kafka input format은 페이로드 파싱 input format을 감싸고, 그 출력에 Kafka 이벤트 타임스탬프, Kafka 토픽 이름, Kafka 이벤트 헤더, 그리고 그 자체도 사용 가능한 어떤 input format으로든 파싱할 수 있는 key 필드를 추가합니다.
예를 들어 개발 환경에서 위키 편집을 나타내는 Kafka 메시지의 다음 구조를 고려해 보세요.
- Kafka 타임스탬프:
1680795276351 - Kafka 토픽:
wiki-edits - Kafka 헤더: -
env=development-zone=z1 - Kafka key:
wiki-edit - Kafka 페이로드 값:
{"channel":"#sv.wikipedia","timestamp":"2016-06-27T00:00:11.080Z","page":"Salo Toraut","delta":31,"namespace":"Main"}
input format으로 { "type": "json" }을 사용하면 페이로드 값만 파싱돼요. 페이로드에 더해 Kafka 메타데이터까지 파싱하려면 kafka input format을 사용하세요.
다음과 같이 구성합니다.
valueFormat: 페이로드 값을 파싱하는 방법을 정의. 페이로드 파싱 input format({ "type": "json" })으로 설정.timestampColumnName: 페이로드의 컬럼과 충돌하지 않도록 Druid 스키마에서 Kafka 타임스탬프에 대한 커스텀 이름을 제공. 기본값은kafka.timestamp.topicColumnName: 페이로드의 컬럼과 충돌하지 않도록 Druid 스키마에서 Kafka 토픽에 대한 커스텀 이름을 제공. 기본값은kafka.topic. 이 필드는 여러 토픽에서 같은 datasource로 데이터를 수집할 때 유용.headerFormat: 기본값string은 Kafka 헤더의 문자열을 UTF-8 인코딩으로 디코딩. 지원되는 다른 인코딩 형식은 다음과 같음: -ISO-8859-1: ISO Latin Alphabet No. 1, 즉 ISO-LATIN-1. -US-ASCII: 일곱 비트 ASCII. ISO646-US라고도 함. Unicode 문자 집합의 Basic Latin 블록. -UTF-16: 십육 비트 UCS Transformation Format, 선택적 byte-order mark로 바이트 순서 식별. -UTF-16BE: 십육 비트 UCS Transformation Format, big-endian 바이트 순서. -UTF-16LE: 십육 비트 UCS Transformation Format, little-endian 바이트 순서.headerColumnPrefix: 페이로드 컬럼과의 충돌을 피하기 위해 Kafka 헤더에 접두사 제공. 기본값은kafka.header.. 예시의 헤더를 고려하면 Druid는 헤더를kafka.header.env,kafka.header.zone컬럼으로 매핑.keyFormat: key를 파싱할 input format 제공. 첫 번째 값만 사용됨. 예시처럼 key 값이 단순 문자열이면tsvformat으로 파싱할 수 있음.{ "type": "tsv", "findColumnsFromHeader": false, "columns": ["x"]}참고로tsv,csv,regex형식에서는 유효한 input format을 만들기 위해columns배열을 제공해야 함. 첫 번째 값만 사용되고, 그 이름은keyColumnName을 대신 사용하므로 무시됨.keyColumnName: 페이로드의 컬럼과 충돌하지 않도록 Kafka key 컬럼의 이름을 제공. 기본값은kafka.key.
다음 input format은 timestampColumnName, topicColumnName, headerColumnPrefix, keyColumnName에 기본값을 사용합니다.
{ "type": "kafka", "valueFormat": { "type": "json" }, "headerFormat": { "type": "string" }, "keyFormat": { "type": "tsv", "findColumnsFromHeader": false, "columns": ["x"] }}
예시 메시지를 다음과 같이 파싱합니다.
{ "channel": "#sv.wikipedia", "timestamp": "2016-06-27T00:00:11.080Z", "page": "Salo Toraut", "delta": 31, "namespace": "Main", "kafka.timestamp": 1680795276351, "kafka.topic": "wiki-edits", "kafka.header.env": "development", "kafka.header.zone": "z1", "kafka.key": "wiki-edit"}
마지막으로 이 Kafka 메타데이터 컬럼들을 dimensionsSpec에 추가하거나, dimensionsSpec이 컬럼을 자동 감지하도록 설정하세요.
다음 supervisor spec은 Kafka 헤더·key·타임스탬프·토픽을 Druid dimension으로 수집하는 방법을 보여줘요.
예시 보기
{ "type": "kafka", "spec": { "ioConfig": { "type": "kafka", "consumerProperties": { "bootstrap.servers": "localhost:9092" }, "topic": "wiki-edits", "inputFormat": { "type": "kafka", "valueFormat": { "type": "json" }, "headerFormat": { "type": "string" }, "keyFormat": { "type": "tsv", "findColumnsFromHeader": false, "columns": ["x"] } }, "useEarliestOffset": true }, "dataSchema": { "dataSource": "wikiticker", "timestampSpec": { "column": "timestamp", "format": "posix" }, "dimensionsSpec": "dimensionsSpec": { "useSchemaDiscovery": true, "includeAllDimensions": true }, "granularitySpec": { "queryGranularity": "none", "rollup": false, "segmentGranularity": "day" } }, "tuningConfig": { "type": "kafka" } }}
Druid가 데이터를 수집한 뒤에는 다음과 같이 Kafka 메타데이터 컬럼을 쿼리할 수 있어요.
SELECT "kafka.header.env", "kafka.key", "kafka.timestamp", "kafka.topic"FROM "wikiticker"
이 쿼리는 다음을 반환합니다.
| kafka.header.env | kafka.key | kafka.timestamp | kafka.topic |
| development | wiki-edit | 1680795276351 | wiki-edits |
Tuning 구성
다음 표는 Kafka에 특화된 tuningConfig 구성 속성을 정리한 것이에요. 모든 스트리밍 수집 방법에 공통인 속성은 Supervisor tuning configuration을 참고하세요.
| Property | Type | Description | Required | Default |
| numPersistThreads | Integer | 디스크에 증분 세그먼트(incremental segment)를 만들고 persist하는 데 사용할 스레드 수. 수집 데이터 처리량이 높을수록 증분 세그먼트 수가 많아져 디스크에 증분 세그먼트를 만드는 데 상당한 CPU 시간이 소요됨. 컬럼 수가 수백~수천 개에 이르는 datasource에서는 증분 세그먼트 생성이 수 초 단위까지 시간이 걸릴 수 있음. 두 시나리오 모두에서 수집이 자주 멈추거나 일시 중지돼 뒤처질 수 있음. 충분한 CPU 자원이 있다면 추가 스레드로 세그먼트 생성을 병렬화해서 수집을 막지 않을 수 있음 | No | 1 |
Kafka 파티션과 Druid 세그먼트에 대한 배포 참고 사항
Druid는 각 Kafka indexing task에 Kafka 파티션을 할당해요. 태스크는 Kafka에서 소비한 이벤트를 세그먼트 granularity interval에 대해 단일 세그먼트로 쓰다가 maxRowsPerSegment, maxTotalRows, intermediateHandoffPeriod 중 하나에 도달합니다. 이 시점에 태스크는 이후 이벤트를 담기 위해 이 세그먼트 granularity에 대해 새 파티션을 만듭니다.
Kafka indexing task는 또한 증분 hand-off(값 hand-off)를 수행해요. 그래서 세그먼트가 준비되는 대로 사용 가능해지고, 태스크 duration이 끝날 때까지 모든 세그먼트를 기다릴 필요가 없습니다. 태스크가 maxRowsPerSegment, maxTotalRows, intermediateHandoffPeriod 중 하나에 도달하면 모든 세그먼트를 hand-off 하고 이후 이벤트를 위해 새 세그먼트 집합을 만듭니다. 이를 통해 태스크가 Middle Manager 서비스에 오래된 세그먼트를 로컬로 쌓아 두지 않고도 더 긴 duration으로 실행될 수 있습니다.
Kafka indexing service는 여전히 작은 세그먼트를 일부 만들 수 있어요. 예를 들어 다음 시나리오를 고려해 보세요.
- 태스크 duration이 4시간
- 세그먼트 granularity가 HOUR로 설정
- supervisor가 9:10에 시작
4시간 뒤인 13:10에 Druid는 새 태스크 집합을 시작합니다. 13:00 ~ 14:00 interval의 이벤트는 기존 태스크와 새 태스크 집합에 나뉘어 들어가면서 작은 세그먼트가 생길 수 있어요. 이들을 이상적인 크기(세그먼트당 약 500~700 MB 범위)의 새 세그먼트로 병합하려면, 선택적으로 다른 세그먼트 granularity로 re-indexing 태스크를 예약할 수 있습니다.
세그먼트 크기 최적화 방법은 Segment size optimization을 참고하세요.
더 알아보기 (Learn more)
관련 주제는 다음을 참고하세요.
- Supervisor API — API로 supervisor를 관리·모니터링하는 방법
- Supervisor — supervisor 상태와 용량 계획
- Loading from Apache Kafka — Apache Kafka에서 데이터 스트리밍하는 튜토리얼
- Kafka input format —
kafkainput format에 대해 배우기