스트림 수집 커넥터

스트림 수집 커넥터 (Stream Ingestion Connectors)

지원되는 각 스트림 수집 커넥터의 모든 구성(configuration)을 정리한 문서예요.

출처: 문서

본문

이 문서는 지원되는 각 스트림 수집 커넥터의 모든 구성을 나열해요.

모든 스트림 커넥터에 공통 (Applicable to all Stream Connectors)

구성 설명
stream.<stream_type>.consumer.factory.class.name 스트림 컨슈머에 사용할 팩토리 클래스
stream.<stream_type>.consumer.prop.auto.offset.reset 데이터 소비를 시작할 소스 스트림의 오프셋·위치. 유효 값: smallest - 스트림의 가장 이른 데이터부터 소비 시작 largest - 스트림의 가장 최신 데이터부터 소비 시작 timestamp - yyyy-MM-dd'T'HH:mm:ss.SSSZ 형식으로 지정된 타임스탬프 이후의 오프셋부터 소비 시작 datetime - 현재 시각부터 지정된 기간·지속 시간 이후의 오프셋부터 소비 시작. 예: 2d 기본값: largest
stream.<stream_type>.topic.name 소비할 소스 스트림 이름
stream.<stream_type>.fetch.timeout.millis 컨슈머에 대한 각 fetch 호출에 사용할 타임아웃(밀리초). 데이터가 제공되기 전에 타임아웃이 만료되면 컨슈머는 빈 배치를 반환해요. 기본값: 5_000
stream.<stream_type>.connection.timeout.millis 업스트림에 대한 연결 생성에 사용하는 타임아웃(밀리초) (업스트림에 대한 초기 연결 타임아웃) 기본값: 30_000
stream.<stream_type>.idle.timeout.millis 스트림이 지정된 시간 동안 유휴(데이터 없음) 상태로 남으면 클라이언트 연결이 리셋되고 새 컨슈머 인스턴스가 생성돼요. 기본값: 180_000
stream.<stream_type>.decoder.class.name 스트림 페이로드를 디코딩하는 데 사용할 디코더 클래스 이름
stream.<stream_type>.decoder.prop 디코더 특정 속성에 사용하는 접두사
topic.consumption.rate.limit 전체 토픽의 메시지 속도 상한. 이 구성을 무시하려면 -1을 사용. 기본값: -1 자세한 내용은 여기 참고
stream.<stream_type>.metadata.populate true로 설정하면 지원되는 컨슈머가 들어오는 페이로드에서 키, 사용자 헤더, 레코드 메타데이터를 추출할 수 있어요. 현재 Kafka 커넥터에서만 지원돼요.
realtime.segment.flush.threshold.time 실시간 세그먼트의 시간 기반 플러시 임계값. 실시간 세그먼트가 커밋·닫기·디스크로 플러시될 준비가 됐는지 결정하는 데 사용돼요. ⚠ 이 시간은 해당 토픽에 구성된 보존 기간보다 작아야 해요.
realtime.segment.flush.threshold.size 완성된 실시간 세그먼트의 크기. ℹ 이 구성은 realtime.segment.flush.threshold.rows가 0으로 설정된 경우에만 사용돼요.
realtime.segment.flush.threshold.rows 실시간 세그먼트의 행 수 기반 플러시 임계값. 이 값이 0으로 설정되면 컨슈머가 파티션에서 소비하는 행 수를 조정해 완성된 세그먼트가 올바른 크기가 되도록 해요 (threshold.time가 먼저 도달하지 않는 한).
realtime.segment.flush.autotune.initialRows SegmentSizeBasedFlushThresholdUpdater에 사용할 초기 행 수. 이 임계값 업데이터는 컨트롤러가 이전 세그먼트의 크기를 기반으로 새 세그먼트의 플러시 임계값을 계산하는 데 사용돼요. ⚠ 이 플러시 임계값 업데이터는 realtime.segment.flush.threshold.rows가 <=0으로 설정된 경우에만 사용돼요. 그렇지 않으면 DefaultFlushThresholdUpdater가 사용돼요.
realtime.segment.commit.timeoutSeconds 컨트롤러가 서버에 의해 세그먼트가 빌드되기를 기다릴 시간 임계값

Kafka 파티션 수준 커넥터 (Kafka Partition-Level Connector)

Kafka 3.x / 4.x

Pinot는 두 개의 Kafka 커넥터 모듈을 제공해요: pinot-kafka-3.0(Kafka 클라이언트 3.9.2, 기본)과 pinot-kafka-4.0(Kafka 클라이언트 4.1.1, KRaft 모드 클러스터용). 레거시 kafka-0.9와 kafka-2.x 모듈은 제거됐어요.

구성 설명
stream.kafka.consumer.factory.class.name 허용 값: - org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory (Kafka 3.x, 기본) - org.apache.pinot.plugin.stream.kafka40.KafkaConsumerFactory (Kafka 4.x)
stream.kafka.topic.name (필수) 수집할 kafka 토픽 이름
stream.kafka.broker.list (필수) kafka 브로커 연결 문자열
stream.kafka.partition.ids 소비할 Kafka 파티션 ID 및/또는 포함 범위의 선택적 쉼표 구분 문자열 (예: "0,2,5", "0-3", 또는 "0-3,6,8-9"). 설정하면 해석된 파티션만 이 테이블이 소비해요. 없거나 비어 있으면 모든 토픽 파티션이 소비돼요(기본 동작). 파티션 ID는 음수가 아닌 정수여야 하며, 범위는 포함적이고, 중복은 조용히 제거되며, 해석된 집합은 10,000개의 고유 파티션 ID로 제한돼요. Pinot는 소비 시작 전에 해석된 ID를 토픽 메타데이터와 대조해 검증해요. 자세한 내용과 예시는 하위 집합 파티션 수집을 보세요.
stream.kafka.buffer.size 기본값: 512000
stream.kafka.socket.timeout 기본값: 10000
stream.kafka.fetcher.size 기본값: 100000
stream.kafka.isolation.level 허용 값: read_committed, read_uncommitted 기본값: read_uncommitted 참고: Kafka에서 트랜잭션을 사용할 때는 read_committed로 설정해야 해요.

지원되는 디코더 클래스 (Supported Decoder Classes)

디코더 클래스 설명
org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder 스키마 레지스트리 없이 일반 JSON 메시지를 디코딩해요.
org.apache.pinot.plugin.inputformat.avro.SimpleAvroMessageDecoder stream.kafka.decoder.prop.schema를 통해 제공되는 스키마를 사용해 Avro 메시지를 디코딩해요.
org.apache.pinot.plugin.inputformat.avro.confluent.KafkaConfluentSchemaRegistryAvroMessageDecoder Confluent Schema Registry에 등록된 스키마를 가진 Avro 메시지를 디코딩해요. stream.kafka.decoder.prop.schema.registry.rest.url 필요.
org.apache.pinot.plugin.inputformat.json.confluent.KafkaConfluentSchemaRegistryJsonMessageDecoder Confluent Schema Registry에 등록된 스키마를 가진 JSON 메시지를 디코딩해요. stream.kafka.decoder.prop.schema.registry.rest.url 필요. Pinot 1.4에 추가됨.
org.apache.pinot.plugin.inputformat.protobuf.ProtoBufMessageDecoder Protocol Buffer 메시지를 디코딩해요.

Kinesis 파티션 수준 커넥터 (Kinesis Partition-level Connector)

구성 설명
stream.kinesis.consumer.factory.class.name 허용 값: org.apache.pinot.plugin.stream.kinesis.KinesisConsumerFactory
stream.kinesis.topic.name (필수) 소비할 Kinesis 데이터 스트림 이름
region (필수) 구성된 Kinesis 데이터 스트림이 있는 AWS 리전
maxRecordsToFetch 단일 GetRecords 요청 중 가져올 최대 레코드 수 기본값: 10000
requests_per_second_limit Pinot가 샤드당 시도할 초당 최대 Kinesis 읽기 요청. 여러 컨슈머가 하나의 샤드 예산을 공유할 수 있도록 0.25 같은 분수 값도 수용해요. Pinot는 이 제한을 GetRecords와 GetShardIterator 읽기 모두에 적용하고, 제한될 때 fetch 타임아웃까지 백오프·재시도해요. 기본값: 1.0
shardIteratorType Pinot가 샤드 반복자를 열 때 사용하는 AWS 샤드 반복자 유형. Pinot는 이 값을 그대로 Kinesis 클라이언트에 전달해요. 기본값: LATEST

키 기반 인증 속성 (Key-based Authentication Properties)

구성 설명
accessKey (필수) AWS Kinesis 데이터 스트림 접근에 사용하는 AWS Access key
secretKey (필수) AWS Kinesis 데이터 스트림 접근에 사용하는 AWS Secret key

IAM 역할 기반 인증 속성 (IAM Role-based Authentication Properties)

구성 설명
iamRoleBasedAccessEnabled AWS Kinesis 데이터 스트림 연결에 IAM 역할 기반 인증을 사용할 때 true로 설정 기본값: false
roleArn (필수) 교차 계정 IAM 역할의 ARN
roleSessionName 클라이언트가 IAM 역할을 가정할 때의 세션 고유 식별자 기본값: pinot-kinesis-<UUID>
externalId AWS 계정 간 신뢰를 관리하고 confused deputy 문제를 방지하는 데 사용하는 고유 식별자. 자세한 내용 여기
sessionDurationSeconds 역할 세션의 지속 시간(초) 기본값: 900
asyncSessionUpdateEnabled 세션 업데이트 활성화 여부를 결정하는 플래그 기본값: true

더 알아보기 (Learn more)