스트림 수집 커넥터
스트림 수집 커넥터 (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 |