수집 구성
수집 구성 (Ingestion Configuration)
Apache Pinot 수집(ingestion) 구성 레퍼런스 문서예요.
출처: 문서
본문
ingestionConfig는 테이블 구성의 최상위 필드 중 하나로, 데이터 수집 관련 속성을 정의해요. 여기에는 배치 수집, stream 수집, 필터링, 변환, 복합 타입 처리 설정이 포함돼요.
ingestionConfig
| 속성 | 설명 |
|---|---|
| batchIngestionConfig | 배치 수집 관련 구성. 아래 batchIngestionConfig 참고. |
| filterConfig | 수집 중 레코드 필터링 구성. 아래 filterConfig 참고. |
| transformConfigs | 컬럼 변환 함수 목록. 아래 transformConfigs 참고. |
| streamIngestionConfig | 스트림 수집 관련 구성. 아래 streamIngestionConfig 참고. |
| complexTypeConfig | 복합 타입(MAP, OPEN_STRUCT) 처리 구성. |
| schemaConformingMode | 수집 중 스키마에 맞춰 입력을 변형할지 제어(NONE이 기본). |
| ignoreNullsForCompositeKeys | upsert/dedup 테이블의 복합 기본 키에서 null 값을 무시할지 여부. |
batchIngestionConfig
배치 수집을 앞단에서 실행하고 세그먼트를 Pinot으로 push하는 데 사용돼요.
| 속성 | 설명 |
|---|---|
| segmentIngestionType | 세그먼트 수집 유형: APPEND 또는 REFRESH. |
| segmentIngestionFrequency | 세그먼트 수집 빈도: DAILY, HOURLY 등. |
| batchConfigMaps | 배치 수집 작업 구성 목록. 각 맵은 다음 키를 포함할 수 있어요: |
| inputDirURI | 입력 데이터의 루트 디렉터리 URI. |
| includeFileNamePattern | 포함할 파일 이름 패턴(예: glob:**/*.csv). |
| excludeFileNamePattern | 제외할 파일 이름 패턴(예: glob:**/*.tmp). |
| inputFormat | 입력 형식(예: csv, avro, parquet, json, orc). |
| outputDirURI | 출력 세그먼트 디렉터리 URI. |
| input.fs.className | 입력 파일시스템 클래스(예: org.apache.pinot.plugin.filesystem.S3PinotFS). |
| input.fs.prop.* | 파일시스템 속성(예: input.fs.prop.region, input.fs.prop.accessKey, input.fs.prop.secretKey). |
| push.mode | push 모드: tar(기본), uri, metadata. |
| segmentNameSpec | 세그먼트 이름 생성 설정. |
| pushSpec | push 작업 설정(pushParallelism, pushAttempts, pushRetryIntervalMillis). |
예시:
"batchIngestionConfig": {
"segmentIngestionType": "APPEND",
"segmentIngestionFrequency": "DAILY",
"batchConfigMaps": [
{
"inputDirURI": "s3://my-bucket/baseballStats/rawdata",
"includeFileNamePattern": "glob:**/*.csv",
"excludeFileNamePattern": "glob:**/*.tmp",
"inputFormat": "csv",
"outputDirURI": "s3://my-bucket/baseballStats/segments",
"input.fs.className": "org.apache.pinot.plugin.filesystem.S3PinotFS",
"input.fs.prop.region": "us-west-2",
"input.fs.prop.accessKey": "${AWS_ACCESS_KEY}",
"input.fs.prop.secretKey": "${AWS_SECRET_KEY}",
"push.mode": "tar"
}
],
"segmentNameSpec": {},
"pushSpec": {}
}
filterConfig
수집 중 레코드를 필터링해 특정 조건에 맞는 레코드만 세그먼트에 포함시켜요.
| 속성 | 설명 |
|---|---|
| filterFunction | 레코드가 필터를 통과할지 결정하는 함수 표현식. false를 반환하면 레코드를 건너뛰어요. Groovy 스크립트 사용 가능. |
예시:
"filterConfig": {
"filterFunction": "Groovy({foo == \"VALUE1\"}, foo)"
}
transformConfigs
수집 중 컬럼을 변환하는 함수 목록이에요. columnName과 transformFunction 쌍으로 지정해요.
| 속성 | 설명 |
|---|---|
| columnName | 변환 결과를 저장할 대상 컬럼 이름. |
| transformFunction | 적용할 변환 함수. Pinot 내장 변환 함수를 사용해요. |
예시:
"transformConfigs": [
{
"columnName": "bar",
"transformFunction": "lower(moo)"
},
{
"columnName": "hoursSinceEpoch",
"transformFunction": "toEpochHours(millis)"
}
]
변환 함수 목록은 변환 (Transformations)을 참고하세요.
streamIngestionConfig
스트림 수집(실시간 테이블) 관련 구성이에요. Kafka, Kinesis, Pulsar 등 stream type별 설정을 streamConfigMaps로 정의해요.
| 속성 | 설명 |
|---|---|
| streamConfigMaps | 스트림별 구성 맵 목록. 각 맵은 다음 키를 포함할 수 있어요: |
| streamType | 스트림 유형: kafka, kinesis, pulsar 등. |
| stream.kafka.broker.list | Kafka broker 목록. |
| stream.kafka.topic.name | Kafka 토픽 이름. |
| stream.kafka.consumer.factory.class.name | 소비자 팩토리 클래스(예: org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory). |
| stream.kafka.consumer.prop.auto.offset.reset | 오프셋 리셋 전략: largest, smallest. |
| stream.kafka.decoder.class.name | 디코더 클래스(예: org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder). |
| stream.kafka.decoder.prop.schema.registry.rest.url | Confluent 스키마 레지스트리 URL. |
| stream.kafka.decoder.prop.schema.registry.schema.name | 스키마 레지스트리의 스키마 이름. |
| realtime.segment.flush.threshold.rows | 소비 세그먼트를 flush 할 행 수 임계값. |
| realtime.segment.flush.threshold.time | 소비 세그먼트를 flush 할 시간 임계값(예: 24h). |
| realtime.segment.flush.threshold.segment.size | 소비 세그먼트 대상 크기(예: 100M). |
예시:
"streamIngestionConfig": {
"streamConfigMaps": [
{
"realtime.segment.flush.threshold.rows": "0",
"realtime.segment.flush.threshold.time": "24h",
"realtime.segment.flush.threshold.segment.size": "150M",
"stream.kafka.broker.list": "XXXX",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.consumer.prop.auto.offset.reset": "largest",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"stream.kafka.decoder.prop.schema.registry.rest.url": "XXXX",
"stream.kafka.decoder.prop.schema.registry.schema.name": "XXXX",
"stream.kafka.topic.name": "XXXX",
"streamType": "kafka"
}
]
}
전체 stream config 키 목록은 스트림 수집 (Stream ingestion) 페이지를 참고하세요.
complexTypeConfig
복합 타입(MAP, OPEN_STRUCT, LIST) 컬럼 처리 구성을 정의해요.
| 속성 | 설명 |
|---|---|
| complexTypeHandlingStrategy | 복합 타입 처리 전략. 기본은 NONE, PASS_RAW(JSON 문자열로 전달), CONVERT_TO_JSON, CREATE_MAP 등. |
| fieldsNotToIndex | 인덱싱에서 제외할 필드 목록. |
| fieldToIndexConfigs | 복합 필드 인덱싱의 세부 구성. |
streamConfigs (deprecated)
streamConfigs 섹션은 릴리스 0.7.0부터 deprecated됐으며, streamIngestionConfig.streamConfigMaps를 사용하세요.