Apache Kafka에서 수집
Apache Kafka에서 수집 (Ingest from Apache Kafka)
Apache Kafka 토픽의 레코드 스트림을 Pinot 테이블로 수집하는 방법을 단계별로 보여 주는 가이드예요.
본문
이 가이드는 Apache Kafka 토픽의 레코드 스트림을 Pinot 테이블로 수집하는 방법을 보여 줘요.
스트림 처리 플랫폼인 Kafka에서 데이터를 수집하는 방법을 배워요. 클러스터 설정 (Set up a cluster) 지침에 따라 로컬 클러스터가 실행 중이어야 해요.
{% hint style="info" %}
이 가이드는 Kafka 3.0 커넥터(kafka30)를 사용합니다. Pinot은 KRaft 모드 Kafka 클러스터용 Kafka 4.0 커넥터도 지원합니다. 올바른 커넥터 선택에 대한 자세한 내용은 Kafka 커넥터 버전 (Kafka Connector Versions) 참고.
{% endhint %}
Kafka 설치와 시작
먼저 로컬 머신에 Kafka를 다운로드해요.
Docker:
docker pull apache/kafka:4.0.0
라우처 스크립트 (Launcher Scripts):
kafka.apache.org/quickstart#quickstart_download에서 Kafka를 다운로드하고 압축을 풀어요:
tar -xzf kafka_2.13-4.0.0.tgz
cd kafka_2.13-4.0.0
다음으로 Kafka 브로커를 띄울 거예요. Kafka 4.0은 기본적으로 KRaft 모드를 사용하며 ZooKeeper가 필요 없어요:
Docker:
docker run --network pinot-demo --name=kafka \
-e KAFKA_NODE_ID=1 \
-e KAFKA_PROCESS_ROLES=broker,controller \
-e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 \
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092 \
-e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
-e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT \
-e KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 \
-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
-e CLUSTER_ID=MkU3OEVBNTcwNTJENDM2Qk \
apache/kafka:4.0.0
참고: --network pinot-demo 플래그는 선택 사항이며, Kafka 컨테이너를 연결하려는 pinot-demo라는 Docker 네트워크가 있다고 가정해요.
라우처 스크립트:
Kafka 4.0은 기본적으로 KRaft 모드를 사용해요. 클러스터 ID를 생성하고 스토리지 디렉터리를 포맷한 다음 브로커를 시작해요:
Kafka 브로커 시작 (KRaft 모드)
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format --standalone -t $KAFKA_CLUSTER_ID -c config/server.properties
bin/kafka-server-start.sh config/server.properties
데이터 소스
터미널에서 다음 스크립트로 JSON 메시지를 생성할 거예요:
import datetime
import uuid
import random
import json
while True:
ts = int(datetime.datetime.now().timestamp()* 1000)
id = str(uuid.uuid4())
count = random.randint(0, 1000)
print(
json.dumps({"ts": ts, "uuid": id, "count": count})
)
datagen.py
이 스크립트(python datagen.py)를 실행하면 다음 출력이 보여요:
{"ts": 1644586485807, "uuid": "93633f7c01d54453a144", "count": 807}
{"ts": 1644586485836, "uuid": "87ebf97feead4e848a2e", "count": 41}
{"ts": 1644586485866, "uuid": "960d4ffa201a4425bb18", "count": 146}
Kafka로 데이터 수집
다음 명령으로 그 메시지 스트림을 Kafka로 파이프해요:
Docker:
python datagen.py | docker exec -i kafka /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic events;
라우처 스크립트:
python datagen.py | bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic events;
다음 명령으로 얼마나 많은 메시지가 수집되었는지 확인할 수 있어요:
Docker:
docker exec -i kafka /opt/kafka/bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic events
라우처 스크립트:
bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic events
출력
events:0:11940
다음 명령으로 메시지 자체를 출력해 볼 수도 있어요:
Docker:
docker exec -i kafka /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic events
라우처 스크립트:
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic events
출력
...
{"ts": 1644586485807, "uuid": "93633f7c01d54453a144", "count": 807}
{"ts": 1644586485836, "uuid": "87ebf97feead4e848a2e", "count": 41}
{"ts": 1644586485866, "uuid": "960d4ffa201a4425bb18", "count": 146}
...
스키마
스키마는 테이블에 어떤 필드가 있는지, JSON 형식으로 그 데이터 타입까지 정의해요.
/tmp/pinot/schema-stream.json이라는 파일을 만들고 다음 내용을 추가하세요.
{
"schemaName": "events",
"dimensionFieldSpecs": [
{
"name": "uuid",
"dataType": "STRING"
}
],
"metricFieldSpecs": [
{
"name": "count",
"dataType": "INT"
}
],
"dateTimeFieldSpecs": [{
"name": "ts",
"dataType": "TIMESTAMP",
"format" : "1:MILLISECONDS:EPOCH",
"granularity": "1:MILLISECONDS"
}]
}
테이블 설정
테이블은 관련 데이터 집합을 나타내는 논리적 추상화예요. 컬럼과 행(Pinot에서는 document)으로 구성돼요. 테이블 설정은 테이블의 속성을 JSON 형식으로 정의해요.
/tmp/pinot/table-config-stream.json이라는 파일을 만들고 다음 내용을 추가하세요.
{
"tableName": "events",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "ts",
"schemaName": "events",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "events",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"realtime.segment.flush.threshold.rows": "0",
"realtime.segment.flush.threshold.time": "24h",
"realtime.segment.flush.threshold.segment.size": "50M",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest"
}
},
"metadata": {
"customConfigs": {}
}
}
JSON 스트림 페이로드 포맷
org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder는 stream.kafka.decoder.prop.jsonFormat이 설정되지 않은 경우 여전히 기본적으로 UTF-8 텍스트 JSON을 사용하므로, 기존 Kafka 테이블 설정은 기록된 동작을 유지해요.
바이너리 JSON 스트림을 디코딩하려면 streamConfigs 블록에 stream.kafka.decoder.prop.jsonFormat을 추가하고 페이로드 인코딩을 고정하세요:
"stream.kafka.decoder.prop.jsonFormat": "SQLITE_JSONB"
지원되는 값은 TEXT, POSTGRES_JSONB, SQLITE_JSONB, SMILE, CBOR, AUTO예요.
AUTO는 토픽이 그 인코딩 중 두 개 이상을 정당하게 포함할 수 있을 때만 사용하세요. AUTO는 opt-in이며, 그 CBOR 감지는 각 메시지가 CBOR self-describe 태그를 포함할 때만 동작해요. 프로듀서가 항상 하나의 알려진 포맷을 내보낸다면 대신 그 포맷을 명시적으로 고정하세요.
이 설정은 스트림 디코딩에만 적용돼요. JSONRecordReader를 사용하는 배치 수집은 여전히 텍스트 JSON 파일을 읽어요.
스키마와 테이블 생성
아래 적절한 명령으로 테이블과 스키마를 생성해요:
Docker:
docker run --rm -ti --network=pinot-demo -v /tmp/pinot:/tmp/pinot apachepinot/pinot:1.0.0 AddTable -schemaFile /tmp/pinot/schema-stream.json -tableConfigFile /tmp/pinot/table-config-stream.json -controllerHost pinot-controller -controllerPort 9000 -exec
라우처 스크립트:
bin/pinot-admin.sh AddTable -schemaFile /tmp/pinot/schema-stream.json -tableConfigFile /tmp/pinot/table-config-stream.json
쿼리
localhost:9000/#/query로 이동해 events 테이블을 클릭해 이 테이블의 처음 10개 행을 보여 주는 쿼리를 실행해요.
_events 테이블 쿼리_
Kafka 수집 가이드라인
Pinot의 Kafka 커넥터 모듈
Pinot은 두 개의 Kafka 커넥터 모듈을 제공해요:
pinot-kafka-3.0-- Kafka 클라이언트 라이브러리 3.x(현재 3.9.2) 사용. Pinot 배포판에 포함된 기본 커넥터예요. 소비자 팩토리 클래스:org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory.pinot-kafka-4.0-- Kafka 클라이언트 라이브러리 4.x(현재 4.1.1) 사용. 이 커넥터는 ZooKeeper 기반 Scala 의존성을 제거하고 순수 Java Kafka 클라이언트를 사용하며, KRaft 모드 Kafka 클러스터에 적합해요. 소비자 팩토리 클래스:org.apache.pinot.plugin.stream.kafka40.KafkaConsumerFactory.
{% hint style="info" %}
레거시 kafka-0.9와 kafka-2.x 커넥터 모듈은 제거되었습니다. org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory를 사용한 이전 Pinot 릴리스에서 업그레이드한다면, 테이블 설정을 위에 나열된 현재 커넥터 클래스 중 하나로 업데이트하세요.
{% endhint %}
{% hint style="info" %} Pinot은 고수준 Kafka 소비자(HLC) 사용을 지원하지 않습니다. Pinot은 정확한 결과 보장, 운영 복잡성 감소, 수평 확장, 스토리지 오버헤드 최소화를 위해 저수준(파티션 수준) 소비자를 사용합니다. {% endhint %}
kafka-2.x 커넥터에서 마이그레이션
기존 테이블 설정이 제거된 kafka-2.x 커넥터를 참조한다면 stream.kafka.consumer.factory.class.name 속성을 업데이트하세요:
- 기존값:
org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory - 바꿀 값(Kafka 3.x):
org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory - 바꿀 값(Kafka 4.x):
org.apache.pinot.plugin.stream.kafka40.KafkaConsumerFactory
다른 스트림 설정 변경은 필요 없어요. Kafka 3.x 커넥터는 Kafka 브로커 2.x 이상과 호환돼요. Kafka 4.x 커넥터는 Kafka 브로커 4.0 이상을 요구해요.
Pinot의 Kafka 설정
ConfigProvider로 Kafka 클라이언트 시크릿 해석
Kafka 3.x와 4.x 커넥터는 실시간 테이블 설정에서 Kafka ConfigProvider 참조를 지원해요. SSL 비밀번호 같은 Kafka 클라이언트 속성을 테이블 설정에 직접 저장하는 대신, Pinot이 Kafka 소비자나 공유 AdminClient를 구성할 때 로드해야 할 때 provider를 사용해요.
provider 별칭과 그 파라미터를 그것을 참조하는 Kafka 속성과 같은 streamConfigs 또는 streamConfigMaps 객체에 선언하세요. 이 예시는 마운트된 properties 파일에서 키스토어 비밀번호를 읽어요:
{
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "events",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"security.protocol": "SSL",
"config.providers": "file",
"config.providers.file.class": "org.apache.kafka.common.config.provider.FileConfigProvider",
"config.providers.file.param.allowed.paths": "/vault/secrets",
"ssl.keystore.password": "${file:/vault/secrets/kafka.properties:keystore.password}"
}
}
참조의 별칭(위 file)은 config.providers에 나열되어야 해요. 그렇지 않으면 Pinot은 그 값을 자체 환경 변수 또는 시스템 속성 표현식으로 처리해요. provider 별칭과 config.providers.<alias>.* 파라미터는 그것이 포함된 스트림 설정 맵에 범위가 한정되며, 이는 다중 스트림 테이블 설정에서 중요해요.
ssl.keystore.password처럼 stream.kafka.consumer.prop. 접두사가 없는 Kafka 클라이언트 속성 이름을 사용하세요. 그 접두사는 임의의 Kafka 클라이언트 속성을 위한 범용 네임스페이스가 아니에요.
Kafka는 클라이언트가 구성될 때 참조를 해석해요. 소비자를 재생성하면 provider 소스를 다시 읽지만, 파일을 바꾼다고 이미 실행 중인 Kafka 클라이언트가 핫 리로드되지는 않아요.
SSL과 함께 Kafka 파티션(저수준) 소비자 사용
Kafka와 schema-registry와 통신하는 데 SSL 기반 인증을 사용하는 예시 설정이에요. 두 세트의 SSL 옵션에 주목하세요: ssl.로 시작하는 것은 Kafka 소비자용이고, stream.kafka.decoder.prop.schema.registry.가 붙은 것은 KafkaConfluentSchemaRegistryAvroMessageDecoder가 사용하는 SchemaRegistryClient용이에요.
{
"tableName": "transcript",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "timestamp",
"timeType": "MILLISECONDS",
"schemaName": "transcript",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "transcript-topic",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.avro.confluent.KafkaConfluentSchemaRegistryAvroMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "localhost:9092",
"schema.registry.url": "",
"security.protocol": "SSL",
"ssl.truststore.location": "",
"ssl.keystore.location": "",
"ssl.truststore.password": "",
"ssl.keystore.password": "",
"ssl.key.password": "",
"stream.kafka.decoder.prop.schema.registry.rest.url": "",
"stream.kafka.decoder.prop.schema.registry.ssl.truststore.location": "",
"stream.kafka.decoder.prop.schema.registry.ssl.keystore.location": "",
"stream.kafka.decoder.prop.schema.registry.ssl.truststore.password": "",
"stream.kafka.decoder.prop.schema.registry.ssl.keystore.password": "",
"stream.kafka.decoder.prop.schema.registry.ssl.keystore.type": "",
"stream.kafka.decoder.prop.schema.registry.ssl.truststore.type": "",
"stream.kafka.decoder.prop.schema.registry.ssl.key.password": "",
"stream.kafka.decoder.prop.schema.registry.ssl.protocol": ""
}
},
"metadata": {
"customConfigs": {}
}
}
Confluent Schema Registry를 JSON 인코딩 메시지와 함께 사용
Kafka 메시지가 JSON으로 인코딩되어 Confluent Schema Registry에 등록되어 있다면 KafkaConfluentSchemaRegistryJsonMessageDecoder를 사용하세요. 이 디코더는 Confluent KafkaJsonSchemaDeserializer를 사용해 JSON 스키마가 레지스트리에서 관리되는 메시지를 디코딩해요.
이 디코더를 쓸 때
- Kafka 프로듀서가 Confluent JSON 스키마 직렬화기로 메시지를 직렬화할 때.
- JSON 스키마가 Confluent Schema Registry에 등록되어 있을 때.
- JSON 메시지에 스키마 검증과 진화(evolution) 지원을 원할 때.
메시지가 Avro로 인코딩되고 Schema Registry에 등록되어 있다면 대신 KafkaConfluentSchemaRegistryAvroMessageDecoder를 사용해요(위 SSL 예시에 나옴). 메시지가 스키마 레지스트리 없는 일반 JSON이라면 JSONMessageDecoder를 사용하세요.
예시 테이블 설정
{
"tableName": "events",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "created_at",
"timeType": "MILLISECONDS",
"schemaName": "events",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "events",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.confluent.KafkaConfluentSchemaRegistryJsonMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "localhost:9092",
"stream.kafka.schema.registry.url": "http://localhost:8081",
"stream.kafka.decoder.prop.schema.registry.rest.url": "http://localhost:8081",
"realtime.segment.flush.threshold.rows": "0",
"realtime.segment.flush.threshold.time": "24h",
"realtime.segment.flush.threshold.segment.size": "50M",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest"
}
},
"metadata": {
"customConfigs": {}
}
}
이 디코더의 핵심 설정 속성:
stream.kafka.decoder.class.name--org.apache.pinot.plugin.inputformat.json.confluent.KafkaConfluentSchemaRegistryJsonMessageDecoder로 설정.stream.kafka.decoder.prop.schema.registry.rest.url-- Confluent Schema Registry의 URL.
인증
이 디코더는 Avro 스키마 레지스트리 디코더와 동일한 인증 옵션을 지원해요. stream.kafka.decoder.prop.schema.registry.* 속성으로 Kafka 소비자와 Schema Registry 클라이언트 모두에 SSL 또는 SASL_SSL 인증을 구성할 수 있어요. 자세한 내용은 위의 SSL 예시와 SASL_SSL 예시 참고.
Schema Registry 기본 인증의 경우 다음 속성을 추가하세요:
"stream.kafka.decoder.prop.basic.auth.credentials.source": "USER_INFO",
"stream.kafka.decoder.prop.schema.registry.basic.auth.user.info": "<username>:<password>"
{% hint style="info" %} 이 디코더는 Pinot 1.4에서 추가되었습니다. Pinot 배포가 1.4 이상 실행 중인지 확인하세요. {% endhint %}
트랜잭션으로 커밋된 메시지 소비
Kafka 3.x와 4.x 커넥터는 Kafka 트랜잭션을 지원해요. 트랜잭션 지원은 Kafka 스트림 설정의 kafka.isolation.level 설정으로 제어되며, read_committed 또는 read_uncommitted(기본)일 수 있어요. read_committed로 설정하면 Kafka 스트림에서 트랜잭션으로 커밋된 메시지만 수집해요.
예를 들어,
{
"tableName": "transcript",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "timestamp",
"timeType": "MILLISECONDS",
"schemaName": "transcript",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "transcript-topic",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.avro.confluent.KafkaConfluentSchemaRegistryAvroMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.isolation.level": "read_committed"
}
},
"metadata": {
"customConfigs": {}
}
}
이 설정의 기본값은 모든 메시지를 읽는 read_uncommitted라는 점에 유의하세요. 또한 이 설정은 저수준 소비자만 지원해요.
만료된 Kafka 오프셋에서 복구
Kafka는 뒤처진 Pinot 소비자가 읽기 전에 레코드를 삭제할 수 있어요. 요청된 오프셋이 만료되면 Pinot은 이제 Kafka 클라이언트의 원시 auto.offset.reset 속성을 기본적으로 earliest로 설정해요. 소비는 Kafka가 여전히 보유한 가장 오래된 레코드에서 재개되며, Pinot의 기존 스트림 데이터 손실 감지가 Kafka가 이미 삭제한 레코드에 대한 불가피한 간격을 보고해요.
이 동작을 재정의하려면 streamConfigs에 원시 Kafka 속성을 추가하세요:
"auto.offset.reset": "latest"
이 원시 속성은 earliest와 latest 같은 Kafka 값을 받아요. stream.kafka.consumer.prop.auto.offset.reset와 혼동하지 마세요: 접두사가 붙은 Pinot 속성은 smallest, largest, 시간 기간, 또는 타임스탬프를 받고 파티션에 저장된 Pinot 오프셋이 없을 때만 시작 위치를 선택해요. 실제로 만료된 저장 오프셋으로부터의 Kafka 클라이언트 복구는 제어하지 않아요.
SASL_SSL과 함께 Kafka 파티션(저수준) 소비자 사용
Kafka와 schema-registry와 통신하는 데 SASL_SSL 기반 인증을 사용하는 예시 설정이에요. 두 세트의 SSL 옵션에 주목하세요: 일부는 Kafka 소비자용이고, stream.kafka.decoder.prop.schema.registry.가 붙은 것은 KafkaConfluentSchemaRegistryAvroMessageDecoder가 사용하는 SchemaRegistryClient용이에요.
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "mytopic",
"stream.kafka.consumer.prop.auto.offset.reset": "largest",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.schema.registry.url": "https://xxx",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.avro.confluent.KafkaConfluentSchemaRegistryAvroMessageDecoder",
"stream.kafka.decoder.prop.schema.registry.rest.url": "https://xxx",
"stream.kafka.decoder.prop.basic.auth.credentials.source": "USER_INFO",
"stream.kafka.decoder.prop.schema.registry.basic.auth.user.info": "schema_registry_username:schema_registry_password",
"sasl.mechanism": "PLAIN" ,
"security.protocol": "SASL_SSL" ,
"sasl.jaas.config":"org.apache.kafka.common.security.scram.ScramLoginModule required username=\"kafkausername\" password=\"kafkapassword\";",
"realtime.segment.flush.threshold.rows": "0",
"realtime.segment.flush.threshold.time": "24h",
"realtime.segment.flush.autotune.initialRows": "3000000",
"realtime.segment.flush.threshold.segment.size": "500M"
},
레코드 헤더를 Pinot 테이블 컬럼으로 추출
Pinot의 Kafka 커넥터는 레코드 헤더와 메타데이터를 Pinot 테이블 컬럼으로 자동 추출하는 것을 지원해요. 레코드 헤더/메타데이터에서 Pinot 테이블 컬럼 이름으로의 매핑:
| Kafka 레코드 (Record) | Pinot 테이블 컬럼 (Column) | 설명 (Description) |
|---|---|---|
| 레코드 키: 모든 타입 | __key : String |
설계 단순화를 위해 레코드 키는 항상 UTF-8 인코딩 String이라고 가정 |
| 레코드 헤더: Map<String, String> | 각 헤더 키는 별도 컬럼으로: __header$HeaderKeyName : String |
설계 단순화를 위해 kafka 레코드의 문자열 헤더를 pinot 테이블 컬럼으로 직접 매핑 |
| 레코드 메타데이터 - offset : long | __metadata$offset : String |
|
| 레코드 메타데이터 - partition : int | __metadata$partition : String |
|
| 레코드 메타데이터 - recordTimestamp : long | __metadata$recordTimestamp : String |
Kafka 테이블에서 메타데이터 추출을 활성화하려면 스트림 설정 metadata.populate를 true로 설정할 수 있어요.
이에 더해, 이 컬럼들 중 어떤 것이든 테이블에서 사용하려면 테이블의 스키마에 명시적으로 나열해야 해요.
예를 들어 Pinot 테이블에 offset과 key만 dimension 컬럼으로 추가하고 싶다면 스키마에 다음과 같이 나열할 수 있어요:
"dimensionFieldSpecs": [
{
"name": "__key",
"dataType": "STRING"
},
{
"name": "__metadata$offset",
"dataType": "STRING"
},
{
"name": "__metadata$partition",
"dataType": "STRING"
},
...
],
스키마가 업데이트되면 이 컬럼들은 다른 pinot 컬럼과 유사해져요. 여기에 수집 변환을 적용하거나 인덱스를 정의할 수 있어요.
{% hint style="info" %} 기존 테이블의 스키마를 업데이트할 때는 스키마 진화 가이드라인을 따르는 것을 잊지 마세요! {% endhint %}
Pinot에 Avro 스키마 위치 알려주기
Avro 파일에서 스키마를 생성하는 독립 실행 유틸리티가 있어요. 자세한 내용은 avro 스키마와 JSON 데이터에서 pinot 스키마 추론 참고.
The Avro schema must be provided 같은 오류를 피하려면 streamConfigs 섹션에서 스키마의 위치를 지정하세요. 예를 들어 현재 섹션이 다음과 같다면:
...
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.avro.SimpleAvroMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "",
"stream.kafka.consumer.prop.auto.offset.reset": "largest"
...
}
그런 다음 "stream.kafka.decoder.prop.schema" 키를 추가하고 그 값으로 스키마 위치를 지정하세요.
부분 파티션 수집 (Subset partition ingestion)
기본적으로 Pinot REALTIME 테이블은 설정된 Kafka 토픽의 모든 파티션에서 소비해요. 일부 시나리오에서는 테이블이 토픽 파티션의 일부만 소비하길 원할 수 있어요. stream.kafka.partition.ids 설정으로 Pinot 테이블이 정확히 어떤 Kafka 파티션을 소비할지 지정할 수 있어요.
부분 파티션 수집을 쓸 때
- 토픽 분할 수집 (Split-topic ingestion) -- 여러 Pinot 테이블이 같은 Kafka 토픽을 공유하고, 각 테이블이 다른 파티션 집합을 담당해요. 같은 토픽이 키로 파티셔닝된 논리적으로 구분된 데이터를 포함하고, 각 파티션 그룹에 대해 별도 테이블(또는 인덱스)을 원할 때 유용해요.
- 멀티 테이블 파티션 배정 -- 처리량이 높은 토픽의 파티션을 여러 Pinot 테이블에 분산해 워크로드 격리, 독립 스케일링, 또는 다른 보존 정책을 적용하고 싶을 때.
- 선택적 소비 -- 토픽의 특정 파티션의 데이터만 필요할 때(예: 특정 지역 또는 테넌트에 해당하는 파티션).
설정
테이블 설정의 streamConfigMaps 항목에 stream.kafka.partition.ids를 추가하세요. 값은 개별 0-기반 파티션 ID, 포함 범위, 또는 둘의 혼합을 쉼표로 구분한 문자열로 포함할 수 있어요:
"stream.kafka.partition.ids": "0-3,6,8-9"
예를 들어 "0,2,5"는 세 개의 명시적 파티션을, "0-3"은 0부터 3까지 포함 파티션을, "0-3,6,8-9"는 한 값에서 두 형식을 혼합해요.
이 설정이 있으면 Pinot은 해석된 파티션 집합에서만 소비해요. 없거나 비어 있으면 Pinot은 토픽의 모든 파티션에서 소비해요(기본 동작).
테이블이 파티션 프루닝에 segmentPartitionConfig도 사용한다면, numPartitions을 소비 하위 집합 크기가 아니라 전체 Kafka 토픽 파티션 수에 맞추세요. Pinot은 stream.kafka.partition.ids가 수집을 하위 집합으로 제한해도 실시간 세그먼트 파티션 메타데이터를 전체 토픽 파티션 수에서 파생해요.
예시: 토픽을 두 테이블에 분할
events라는 파티션 두 개(0과 1)를 가진 Kafka 토픽이 있다고 가정해요. 각각 한 파티션에서 소비하는 Pinot 테이블 두 개를 만들 수 있어요:
테이블 events_part_0:
{
"tableName": "events_part_0",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "ts",
"schemaName": "events",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP"
},
"ingestionConfig": {
"streamIngestionConfig": {
"streamConfigMaps": [
{
"streamType": "kafka",
"stream.kafka.topic.name": "events",
"stream.kafka.partition.ids": "0",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest",
"realtime.segment.flush.threshold.rows": "0",
"realtime.segment.flush.threshold.time": "24h",
"realtime.segment.flush.threshold.segment.size": "50M"
}
]
}
},
"metadata": {
"customConfigs": {}
}
}
테이블 events_part_1:
{
"tableName": "events_part_1",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "ts",
"schemaName": "events",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP"
},
"ingestionConfig": {
"streamIngestionConfig": {
"streamConfigMaps": [
{
"streamType": "kafka",
"stream.kafka.topic.name": "events",
"stream.kafka.partition.ids": "1",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest",
"realtime.segment.flush.threshold.rows": "0",
"realtime.segment.flush.threshold.time": "24h",
"realtime.segment.flush.threshold.segment.size": "50M"
}
]
}
},
"metadata": {
"customConfigs": {}
}
}
검증 규칙과 제한 사항
- 파티션 ID는 음이 아닌 정수여야 해요. 범위 값은 포함이며, 범위의 시작은 끝보다 작거나 같아야 해요.
- 정수가 아닌 값(예:
"abc")과 잘못된 범위(예:"5-")는 검증 오류를 일으켜요. - 중복 ID는 범위 확장 후 조용히 중복 제거돼요. 예를 들어
"0,2,0,5"는"0,2,5"로 처리돼요. - 파티션 ID는 설정에 지정된 순서와 무관하게 안정적 정렬을 위해 내부적으로 정렬돼요.
- Pinot은 소비를 시작하기 전에 해석된 파티션 ID를 Kafka 토픽 메타데이터와 검증해요. 지정된 파티션 ID가 토픽에 없으면 오류가 발생해요.
- 해석된 집합에는 최대 10,000개의 고유 파티션 ID를 포함할 수 있어요.
- 부분 수집 실시간 테이블에
segmentPartitionConfig를 구성한다면numPartitions을 전체 Kafka 토픽 파티션 수로 설정하세요. 예를 들어 테이블이 8-파티션 토픽의"0,3"을 소비한다면2가 아닌8을 사용하세요. - 같은 토픽에서 소비하는 여러 테이블과 부분 파티션 수집을 사용할 때, 각 레코드가 정확히 한 테이블에서 소비되길 원한다면 파티션 배정이 겹치지 않도록 하세요. Pinot은 테이블 간 겹치지 않는 파티션 배정을 강제하지 않아요.
- 파티션 ID, 쉼표, 범위 경계 주변의 공백은 트리밍돼요 (예:
" 0 - 3 , 5 "는 유효).
Protocol Buffers (Protobuf) 포맷 사용
Pinot은 설정에 따라 여러 디코드 옵션으로 Kafka에서 Protocol Buffer 메시지 디코딩을 지원해요.
ProtoBufMessageDecoder (descriptor 파일 기반)
Protobuf 스키마에 대한 사전 컴파일된 .desc(descriptor) 파일이 있을 때 ProtoBufMessageDecoder를 사용해요. 이 디코더는 동적 메시지 파싱을 사용하며 컴파일된 Java 클래스가 필요하지 않아요.
필수 스트림 설정 속성:
| 속성 (Property) | 설명 (Description) |
|---|---|
stream.kafka.decoder.prop.descriptorFile |
.desc descriptor 파일의 경로 또는 URI. 로컬 파일 경로, HDFS, 기타 Pinot 지원 파일 시스템 지원 |
stream.kafka.decoder.prop.protoClassName |
(선택) descriptor 내의 정규화된 Protobuf 메시지 이름. 생략하면 descriptor의 첫 번째 메시지 타입이 사용됨 |
stream.kafka.decoder.prop.descriptorFileFallbackEnabled |
(선택) true 또는 false. 이 테이블에 대한 클러스터 전체 pinot.server.protobuf.descriptor.fallback.enabled 설정을 재정의. 기본값은 true |
예시 streamConfigs:
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "my-protobuf-topic",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.protobuf.ProtoBufMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.decoder.prop.descriptorFile": "/path/to/my_message.desc",
"stream.kafka.decoder.prop.protoClassName": "mypackage.MyMessage"
}
원격 descriptor URI의 경우, 이 디코더가 생성될 때마다 Pinot이 descriptor를 새로 가져오므로 제자리 descriptor 업데이트는 다음 소비 중인 세그먼트 전환에서 반영돼요. 가져오기가 실패하면 Pinot은 같은 서버 JVM에서 같은 URI에 대해 이전에 성공적으로 가져오고 해석한 마지막 descriptor 내용을 재사용할 수 있어요. 이 폴백은 기본적으로 활성화되어 있으며 사용 시 경고를 기록해요. descriptor를 한 번도 로드한 적 없는 새 서버 프로세스에는 도움이 되지 않고, 로컬 descriptor 파일에는 적용되지 않아요.
Pinot은 원격 가져오기가 성공했지만 descriptor가 비어 있거나 손상되었거나 설정된 메시지 타입이 없을 때는 폴백하지 않아요; 잘못된 스키마 배포가 보이도록 디코더 초기화가 실패해요. 원격 가져오기 실패 시 수집을 실패시키고 싶다면 테이블의 streamConfigs에서 stream.kafka.decoder.prop.descriptorFileFallbackEnabled를 false로 설정하거나 클러스터 설정을 비활성화하세요. 테이블 설정이 우선해요. 이 폴백은 ProtoBufMessageDecoder에만 적용되고 컴파일된 JAR 또는 배치 Protobuf 리더에는 적용되지 않아요.
ProtoBufCodeGenMessageDecoder (컴파일된 JAR 기반)
컴파일된 Protobuf Java 클래스를 포함하는 JAR이 있을 때 ProtoBufCodeGenMessageDecoder를 사용해요. 이 디코더는 런타임 코드 생성을 사용해 디코딩 성능을 향상해요.
필수 스트림 설정 속성:
| 속성 (Property) | 설명 (Description) |
|---|---|
stream.kafka.decoder.prop.jarFile |
컴파일된 Protobuf 클래스를 포함하는 JAR 파일의 경로 또는 URI |
stream.kafka.decoder.prop.protoClassName |
Protobuf 메시지의 정규화된 Java 클래스 이름 (필수) |
예시 streamConfigs:
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "my-protobuf-topic",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.protobuf.ProtoBufCodeGenMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.decoder.prop.jarFile": "/path/to/my-protobuf-classes.jar",
"stream.kafka.decoder.prop.protoClassName": "com.example.proto.MyMessage"
}
KafkaConfluentSchemaRegistryProtoBufMessageDecoder (Confluent Schema Registry)
Protobuf 스키마가 Confluent Schema Registry에서 관리될 때 KafkaConfluentSchemaRegistryProtoBufMessageDecoder를 사용해요. 이 디코더는 런타임에 레지스트리에서 스키마를 자동으로 해석해요.
필수 스트림 설정 속성:
| 속성 (Property) | 설명 (Description) |
|---|---|
stream.kafka.decoder.prop.schema.registry.rest.url |
Confluent Schema Registry의 URL |
선택 속성:
| 속성 (Property) | 설명 (Description) |
|---|---|
stream.kafka.decoder.prop.cached.schema.map.capacity |
캐시된 스키마의 최대 수. 기본값: 1000 |
stream.kafka.decoder.prop.schema.registry.* |
Schema Registry 연결을 위한 SSL 및 인증 옵션 (Avro Confluent 디코더와 같은 패턴) |
예시 streamConfigs:
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "my-protobuf-topic",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.protobuf.KafkaConfluentSchemaRegistryProtoBufMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.decoder.prop.schema.registry.rest.url": "http://schema-registry:8081"
}
BSON 포맷 사용
Pinot은 BSONMessageDecoder로 Kafka에서 BSON 메시지 디코딩을 지원해요. 각 Kafka 메시지가 단일 바이너리 인코딩 BSON 문서(예: MongoDB change-data-capture 파이프라인으로 생성된 레코드)를 포함할 때 사용해요.
예시 streamConfigs:
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "mongo-events",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.bson.BSONMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092"
}
BSONMessageDecoder는 배치 BSON 수집과 같은 추출기를 사용하므로, ObjectId 값은 hex 문자열로, Date와 BsonTimestamp 값은 java.sql.Timestamp로, Decimal128 값은 NaN과 Infinity가 null로 표출되는 BigDecimal로, 바이너리 값은 byte[]로, 내장 문서는 Map<String, Object>로, 배열은 Object[]로 디코딩돼요.
Apache Arrow 포맷 사용
Pinot은 ArrowMessageDecoder로 Kafka에서 Apache Arrow IPC 스트리밍 포맷 메시지 디코딩을 지원해요. 업스트림 시스템이 Arrow 포맷으로 직렬화된 데이터를 생성할 때 유용해요.
선택적 스트림 설정 속성:
| 속성 (Property) | 설명 (Description) |
|---|---|
stream.kafka.decoder.prop.arrow.allocator.limit |
Arrow 할당자 최대 메모리(바이트). 기본값: 268435456 (256 MB) |
stream.kafka.decoder.prop.extractRawTimeValues |
Arrow Date, Time, Timestamp 값을 Pinot의 기본 변환 값 대신 원시 정수로 유지. 기본값: false |
예시 streamConfigs:
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "my-arrow-topic",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.arrow.ArrowMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.broker.list": "kafka:9092",
"stream.kafka.decoder.prop.arrow.allocator.limit": "536870912",
"stream.kafka.decoder.prop.extractRawTimeValues": "true"
}
{% hint style="info" %} Arrow 디코더는 각 Kafka 메시지가 완전한 Arrow IPC 스트림(스키마 + 레코드 배치)을 포함할 것이라고 기대합니다. 프로듀서가 IPC 스트리밍 포맷으로 Arrow 데이터를 직렬화하는지 확인하세요. {% endhint %}
stream.kafka.decoder.prop.extractRawTimeValues가 false(기본)일 때 Pinot은 추출 중 Arrow Date, Time, Timestamp 값을 변환해요. true로 설정하면 원시 정수를 유지해요: Date는 epoch 이후 일 수로 유지되는 반면, Time과 Timestamp는 스키마에 선언된 Arrow 단위로 유지돼요.
각 Arrow Kafka 메시지는 0개, 1개 또는 여러 Pinot 행을 산출할 수 있어요. 빈 배치는 무시되고, 단일 행 배치는 하나의 Pinot 행으로 수집되며, 여러 행 배치는 여러 Pinot 행으로 분기돼요.
Kafka 파티션 부분 소비
기본적으로 Pinot 실시간 테이블은 Kafka 토픽의 모든 파티션을 소비해요. stream.kafka.partition.ids 속성으로 수집을 특정 파티션 하위 집합으로 제한할 수 있어요. 이것은 다음과 같은 경우 유용해요:
- 독립적인 스케일링을 위해 단일 Kafka 토픽을 여러 Pinot 테이블에 분할
- 각기 다른 테이블이 다른 파티션 범위를 소유하는 멀티 테넌트 시나리오
설정
streamConfigs에 파티션 ID, 포함 범위, 또는 둘 다의 쉼표로 구분된 문자열로 stream.kafka.partition.ids를 추가하세요:
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.topic.name": "myTopic",
"stream.kafka.broker.list": "localhost:9092",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.partition.ids": "0-3,6,8-9"
}
예를 들어 "0,2,5"는 개별 파티션을, "0-3"은 포함 범위를, "0-3,6,8-9"는 한 값에서 두 형식을 혼합해요.
테이블이 segmentPartitionConfig도 사용한다면 numPartitions을 소비 하위 집합 크기 대신 전체 Kafka 토픽 파티션 수로 설정하세요. Pinot은 수집이 특정 Kafka 파티션으로 제한되어도 전체 토픽 파티션 수에서 실시간 세그먼트 파티션 메타데이터를 계산해요.
참고
- 파티션 ID는 음이 아닌 정수여야 해요. 범위는 포함이며
start <= end여야 해요. - Pinot은 소비를 시작하기 전에 해석된 파티션 ID를 Kafka 토픽 메타데이터와 검증해요.
- 목록의 중복 ID는 자동으로 중복 제거돼요.
- 해석된 집합에는 최대 10,000개의 고유 파티션 ID를 포함할 수 있어요.
- 브로커에 보고되는 총 파티션 수는 전체 Kafka 토픽 크기를 반영해, 같은 토픽을 공유하는 테이블 간 올바른 쿼리 라우팅을 보장해요.
segmentPartitionConfig를 구성한다면numPartitions을 전체 Kafka 토픽 파티션 수로 설정하세요. 예를 들어 테이블이 6-파티션 토픽의"1,4"를 소비한다면2가 아닌6을 사용하세요.- 토픽을 두 테이블에 분할할 때, 하나는 짝수 ID, 다른 하나는 홀수 ID로 구성하세요 (예: 4-파티션 토픽의 경우
"0,2"와"1,3").