수집 구성

수집 구성 (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를 사용하세요.

더 알아보기 (Learn more)