실제 적용에서의 배치 데이터 수집
실제 적용에서의 배치 데이터 수집 (Batch Data Ingestion in Practice)
실제로는 Pinot 데이터 수집을 파이프라인이나 예약된 작업으로 실행해야 해요.
pinot-distribution이 이미 빌드됐다고 가정하면, examples 디렉터리 안에서 여러 샘플 테이블 레이아웃을 찾을 수 있어요.
출처: 문서
본문
테이블 레이아웃 (Table Layout)
보통 각 테이블은 airlineStats처럼 자신만의 디렉터리를 가져요.
테이블 디렉터리 안에서 모든 입력 데이터를 두기 위해 rawdata가 만들어져요.
일반적으로 타임스탬프가 있는 데이터 이벤트의 경우 데이터를 파티셔닝해 일별 폴더에 저장해요. 예를 들어 일반적인 레이아웃은 rawdata/%yyyy%/%mm%/%dd%/[daily_input_files] 패턴을 따르게 돼요.
/var/pinot/airlineStats/rawdata/2014/01/01/airlineStats_data_2014-01-01.avro
/var/pinot/airlineStats/rawdata/2014/01/02/airlineStats_data_2014-01-02.avro
/var/pinot/airlineStats/rawdata/2014/01/03/airlineStats_data_2014-01-03.avro
...
/var/pinot/airlineStats/rawdata/2014/01/31/airlineStats_data_2014-01-31.avro
배치 수집 작업 구성 (Configuring batch ingestion job)
데이터를 어떻게 수집할지 설명하는 배치 수집 job spec 파일을 만드세요.
아래는 예시예요(examples/batch/airlineStats/ingestionJobSpec.yaml에도 있음).
# executionFrameworkSpec: 실행할 수집 작업 정의
executionFrameworkSpec:
# name: 실행 프레임워크 이름
name: 'standalone'
# segmentGenerationJobRunnerClassName: org.apache.pinot.spi.batch.ingestion.runner.SegmentGenerationJobRunner 인터페이스 구현 클래스 이름
segmentGenerationJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.standalone.SegmentGenerationJobRunner'
# segmentTarPushJobRunnerClassName: org.apache.pinot.spi.batch.ingestion.runner.SegmentTarPushJobRunner 인터페이스 구현 클래스 이름
segmentTarPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.standalone.SegmentTarPushJobRunner'
# segmentUriPushJobRunnerClassName: org.apache.pinot.spi.batch.ingestion.runner.SegmentUriPushJobRunner 인터페이스 구현 클래스 이름
segmentUriPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.standalone.SegmentUriPushJobRunner'
# jobType: Pinot 수집 작업 유형
# 지원되는 job 유형:
# 'SegmentCreation'
# 'SegmentTarPush'
# 'SegmentUriPush'
# 'SegmentCreationAndTarPush'
# 'SegmentCreationAndUriPush'
jobType: SegmentCreationAndTarPush
# inputDirURI: 입력 데이터의 루트 디렉터리, PinotFS에 구성된 scheme 있어야 함
inputDirURI: 'examples/batch/airlineStats/rawdata'
# includeFileNamePattern / excludeFileNamePattern: Java NIO PathMatcher 패턴 (glob: 또는 regex:)
# https://docs.pinot.apache.org/configuration-reference/job-specification#file-name-patterns 참고
# 샘플 사용법:
# 'glob:**/*.avro' 또는 'regex:.*[.]avro'는 inputDirURI 아래의 Avro 경로(전체 경로)를 매칭
includeFileNamePattern: 'glob:**/*.avro'
# excludeFileNamePattern: 'glob:**/*.tmp' 또는 'regex:.*[.]tmp'
# excludeFileNamePattern: ''
# outputDirURI: 출력 세그먼트의 루트 디렉터리, PinotFS에 구성된 scheme 있어야 함
outputDirURI: 'examples/batch/airlineStats/segments'
# overwriteOutput: 기존 출력 세그먼트 덮어쓰기 여부
overwriteOutput: true
# 세그먼트 메타데이터 push 작업의 데이터 전송을 줄이기 위해 분리된 metadata 전용 tar gz 파일 생성
createMetadataTarGz: true
# 세그먼트 생성 작업 병렬도
segmentCreationJobParallelism: 4
# pinotFSSpecs: 모든 관련 Pinot 파일 시스템 정의
pinotFSSpecs:
- # scheme: PinotFS 식별에 사용
# 예: local, hdfs, dbfs 등
scheme: file
# className: PinotFS 인스턴스 생성에 사용되는 클래스 이름
# 예:
# org.apache.pinot.spi.filesystem.LocalPinotFS는 로컬 파일시스템용
# org.apache.pinot.plugin.filesystem.AzurePinotFS는 Azure Data Lake용
# org.apache.pinot.plugin.filesystem.HadoopPinotFS는 HDFS용
className: org.apache.pinot.spi.filesystem.LocalPinotFS
# recordReaderSpec: 모든 record reader 정의
recordReaderSpec:
# dataFormat: 레코드 데이터 형식, 예: 'avro', 'parquet', 'orc', 'csv', 'json', 'thrift'
dataFormat: 'avro'
# className: 해당 RecordReader 클래스 이름
# 예:
# org.apache.pinot.plugin.inputformat.avro.AvroRecordReader
# org.apache.pinot.plugin.inputformat.csv.CSVRecordReader
# org.apache.pinot.plugin.inputformat.parquet.ParquetRecordReader
# org.apache.pinot.plugin.inputformat.json.JSONRecordReader
# org.apache.pinot.plugin.inputformat.orc.ORCRecordReader
# org.apache.pinot.plugin.inputformat.thrift.ThriftRecordReader
className: 'org.apache.pinot.plugin.inputformat.avro.AvroRecordReader'
# tableSpec: 테이블 이름과 해당 테이블 config 및 schema를 가져올 위치 정의
tableSpec:
# tableName: 테이블 이름
tableName: 'airlineStats'
# schemaURI: 테이블 스키마를 읽을 위치 정의, PinotFS 또는 HTTP 지원
# 예:
# hdfs://path/to/table_schema.json
# http://localhost:9000/tables/myTable/schema
schemaURI: 'http://localhost:9000/tables/airlineStats/schema'
# tableConfigURI: 테이블 config를 읽을 위치 정의, PinotFS 또는 HTTP 지원
# 예:
# hdfs://path/to/table_config.json
# http://localhost:9000/tables/myTable
# 참고: pinot controller에서 Pinot 테이블 config를 직접 읽는 API는 JSON wrapper를 포함
# 실제 테이블 config는 'OFFLINE' 필드 아래의 객체
tableConfigURI: 'http://localhost:9000/tables/airlineStats'
# segmentNameGeneratorSpec: SegmentNameGenerator를 초기화하는 방법 정의
segmentNameGeneratorSpec:
# type: 현재 지원되는 유형은 'simple'과 'normalizedDate'
type: normalizedDate
# configs: SegmentNameGenerator 초기화 config
configs:
segment.name.prefix: 'airlineStats_batch'
exclude.sequence.id: true
# pinotClusterSpecs: Pinot 클러스터 접근 지점 정의
pinotClusterSpecs:
- # controllerURI: 테이블/스키마 정보와 데이터 push에 사용
# 예: http://localhost:9000
controllerURI: 'http://localhost:9000'
# pushJobSpec: 세그먼트 push 작업 관련 구성 정의
pushJobSpec:
# 세그먼트 push 작업 병렬도
pushParallelism: 4
# pushAttempts: push 작업 시도 횟수, 기본값 1(재시도 없음)
pushAttempts: 2
# pushRetryIntervalMillis: 재시도 대기 ms, 기본값 1초
pushRetryIntervalMillis: 1000
# URI 및 METADATA push 유형에 적용
# true이고 세그먼트가 아직 딥 스토어에 없으면 딥 스토어로 이동
copyToDeepStoreForMetadataPush: false
# 세그먼트 metadata tar gz 파일이 있으면 push에 사용 선호
preferMetadataTarGz: true
작업 실행 (Executing the job)
아래 명령이 Pinot 클러스터에 예시 테이블을 생성해요.
bin/pinot-admin.sh AddTable -schemaFile examples/batch/airlineStats/airlineStats_schema.json -tableConfigFile examples/batch/airlineStats/airlineStats_offline_table_config.json -exec
아래 명령이 수집 작업을 시작해 Pinot 세그먼트를 생성하고 클러스터에 push해요.
bin/pinot-admin.sh LaunchDataIngestionJob -jobSpecFile examples/batch/airlineStats/ingestionJobSpec.yaml
작업이 끝나면 세그먼트는 입력 디렉터리 레이아웃과 동일한 레이아웃으로 examples/batch/airlineStats/segments에 저장돼요.
/var/pinot/airlineStats/segments/2014/01/01/airlineStats_batch_2014-01-01_2014-01-01.tar.gz
...
/var/pinot/airlineStats/segments/2014/01/31/airlineStats_batch_2014-01-31_2014-01-31.tar.gz
Spark로 작업 실행 (Executing the job using Spark)
pinot-all만으로는 Spark에 부족해요. Spark(및 Hadoop) 배치 러너는plugins-external/아래에 있으며pinot-all-*-jar-with-dependencies.jar에 패키징되지 않아요. 클래스패스에pinot-all만 두면SparkSegmentGenerationJobRunner같은 클래스에 대해ClassNotFoundException이 발생해요. 항상plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-spark-3/의 shaded jar를 추가하고plugins.dir에plugins-external을 포함하세요.
아래 예시는 로컬 파일시스템(S3/HDFS 자격증명 없음)의 샘플 airlineStats 데이터에 대해 로컬 모드의 Spark를 실행해요. Spark 3.x 배포판을 다운로드하고 SPARK_HOME을 설정하거나 사전 설치된 Spark를 사용하세요.
로컬 설치 가이드를 따라 Pinot 바이너리 배포판을 빌드하거나 다운로드하세요.
Local-FS Spark job spec
로컬 경로에는 LocalPinotFS를 사용하세요. Spark 3 러너 클래스는 spark3 패키지에 있어요. examples/batch/airlineStats/sparkIngestionJobSpec.yaml에서 시작해 아래와 같이 클래스 이름/FS를 맞춰 조정할 수 있어요.
# executionFrameworkSpec: 실행할 수집 작업 정의
executionFrameworkSpec:
# name: 실행 프레임워크 이름
name: 'spark'
# Spark 3 패키지: org.apache.pinot.plugin.ingestion.batch.spark3.*
segmentGenerationJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentGenerationJobRunner'
segmentTarPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentTarPushJobRunner'
segmentUriPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentUriPushJobRunner'
segmentMetadataPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark3.SparkSegmentMetadataPushJobRunner'
# extraConfigs: 실행 프레임워크용 추가 config
extraConfigs:
# stagingDir는 분산 파일시스템에서 모든 세그먼트를 호스팅한 다음 이 디렉터리 전체를 출력 디렉터리로 이동하는 데 사용
stagingDir: examples/batch/airlineStats/staging
# jobType: Pinot 수집 작업 유형
# 지원되는 job 유형:
# 'SegmentCreation'
# 'SegmentTarPush'
# 'SegmentUriPush'
# 'SegmentCreationAndTarPush'
# 'SegmentCreationAndUriPush'
# 'SegmentCreationAndMetadataPush'
jobType: SegmentCreationAndTarPush
# inputDirURI: 입력 데이터의 루트 디렉터리, PinotFS에 구성된 scheme 있어야 함
inputDirURI: 'examples/batch/airlineStats/rawdata'
# includeFileNamePattern / excludeFileNamePattern: Java NIO PathMatcher (glob: 또는 regex:)
# configuration-reference/job-specification.md#file-name-patterns 참고
includeFileNamePattern: 'glob:**/*.avro'
# excludeFileNamePattern: 'glob:**/*.tmp'
# outputDirURI: 출력 세그먼트의 루트 디렉터리, PinotFS에 구성된 scheme 있어야 함
outputDirURI: 'examples/batch/airlineStats/segments'
# overwriteOutput: 기존 출력 세그먼트 덮어쓰기 여부
overwriteOutput: true
# pinotFSSpecs: 모든 관련 Pinot 파일 시스템 정의
pinotFSSpecs:
- # 로컬 파일시스템 — 이 레시피에는 클라우드 또는 HDFS 자격증명 불필요
scheme: file
className: org.apache.pinot.spi.filesystem.LocalPinotFS
# recordReaderSpec: 모든 record reader 정의
recordReaderSpec:
dataFormat: 'avro'
className: 'org.apache.pinot.plugin.inputformat.avro.AvroRecordReader'
# tableSpec: 테이블 이름과 해당 테이블 config 및 schema를 가져올 위치 정의
tableSpec:
tableName: 'airlineStats'
schemaURI: 'http://localhost:9000/tables/airlineStats/schema'
# 참고: pinot controller에서 Pinot 테이블 config를 직접 읽는 API는 JSON wrapper를 포함
# 실제 테이블 config는 'OFFLINE' 필드 아래의 객체
tableConfigURI: 'http://localhost:9000/tables/airlineStats'
# segmentNameGeneratorSpec: SegmentNameGenerator를 초기화하는 방법 정의
segmentNameGeneratorSpec:
type: normalizedDate
configs:
segment.name.prefix: 'airlineStats_batch'
exclude.sequence.id: true
# pinotClusterSpecs: Pinot 클러스터 접근 지점 정의
pinotClusterSpecs:
- controllerURI: 'http://localhost:9000'
# pushJobSpec: 세그먼트 push 작업 관련 구성 정의
pushJobSpec:
pushParallelism: 2
pushAttempts: 2
pushRetryIntervalMillis: 1000
로컬 모드 spark-submit
PINOT_DISTRIBUTION_DIR이 압축 풀린 바이너리 배포판(lib/, plugins/, plugins-external/, examples/ 포함)을 가리키는지 확인하세요.
필수 요소:
plugins.dir—plugins(Avro 같은 record reader)와plugins-external(Spark 배치 러너) 둘 다 포함하는 세미콜론 구분 목록spark.driver.extraClassPath/spark.executor.extraClassPath—pinot-batch-ingestion-spark-3-*-shaded.jar및pinot-all-*-jar-with-dependencies.jar
export PINOT_VERSION=1.5.1 #설치한 Pinot 버전으로 설정
export PINOT_DISTRIBUTION_DIR=${PINOT_ROOT_DIR}/build/
cd ${PINOT_DISTRIBUTION_DIR}
SPARK_BATCH_PLUGIN=${PINOT_DISTRIBUTION_DIR}/plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-spark-3/pinot-batch-ingestion-spark-3-${PINOT_VERSION}-shaded.jar
PINOT_ALL_JAR=${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar
${SPARK_HOME}/bin/spark-submit \
--class org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand \
--master "local[2]" \
--deploy-mode client \
--conf "spark.driver.extraJavaOptions=-Dplugins.dir=${PINOT_DISTRIBUTION_DIR}/plugins;${PINOT_DISTRIBUTION_DIR}/plugins-external -Dlog4j2.configurationFile=${PINOT_DISTRIBUTION_DIR}/conf/pinot-ingestion-job-log4j2.xml" \
--conf "spark.driver.extraClassPath=${SPARK_BATCH_PLUGIN}:${PINOT_ALL_JAR}" \
--conf "spark.executor.extraClassPath=${SPARK_BATCH_PLUGIN}:${PINOT_ALL_JAR}" \
local://${PINOT_ALL_JAR} \
-jobSpecFile ${PINOT_DISTRIBUTION_DIR}/examples/batch/airlineStats/sparkIngestionJobSpec.yaml
ClassNotFoundException문제 해결:plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-spark-3/아래의 shaded jar가${PINOT_VERSION}용으로 존재하고, 드라이버와 실행기 클래스패스 양쪽에 있으며,plugins.dir이plugins-external을 나열하는지 확인하세요. Spark 배치 수집 페이지 FAQ도 참고하세요.
Hadoop으로 작업 실행 (Executing the job using Hadoop)
pinot-all만으로는 MapReduce에 부족해요. Hadoop 배치 플러그인은plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-hadoop/아래에 있으며HADOOP_CLASSPATH에 있어야 해요.plugins.dir에plugins-external을 포함하세요.
샘플 Hadoop 수집 job spec(examples/batch/airlineStats/hadoopIngestionJobSpec.yaml에도 있음). 로컬 파일시스템 드라이 런에는 아래처럼 LocalPinotFS를 선호하세요(체크인된 샘플은 HDFS 지향 파이프라인의 경우 HadoopPinotFS를 사용할 수 있음).
# executionFrameworkSpec: 실행할 수집 작업 정의
executionFrameworkSpec:
name: 'hadoop'
segmentGenerationJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.hadoop.HadoopSegmentGenerationJobRunner'
segmentTarPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.hadoop.HadoopSegmentTarPushJobRunner'
segmentUriPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.hadoop.HadoopSegmentUriPushJobRunner'
segmentMetadataPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.hadoop.HadoopSegmentMetadataPushJobRunner'
extraConfigs:
stagingDir: examples/batch/airlineStats/staging
jobType: SegmentCreationAndTarPush
inputDirURI: 'examples/batch/airlineStats/rawdata'
includeFileNamePattern: 'glob:**/*.avro'
# excludeFileNamePattern: 'glob:**/*.tmp'
outputDirURI: 'examples/batch/airlineStats/segments'
overwriteOutput: true
pinotFSSpecs:
- scheme: file
className: org.apache.pinot.spi.filesystem.LocalPinotFS
recordReaderSpec:
dataFormat: 'avro'
className: 'org.apache.pinot.plugin.inputformat.avro.AvroRecordReader'
tableSpec:
tableName: 'airlineStats'
schemaURI: 'http://localhost:9000/tables/airlineStats/schema'
tableConfigURI: 'http://localhost:9000/tables/airlineStats'
segmentNameGeneratorSpec:
type: normalizedDate
configs:
segment.name.prefix: 'airlineStats_batch'
exclude.sequence.id: true
pinotClusterSpecs:
- controllerURI: 'http://localhost:9000'
pushJobSpec:
pushParallelism: 2
pushAttempts: 2
pushRetryIntervalMillis: 1000
PINOT_ROOT_DIR과 PINOT_VERSION이 올바르게 설정됐는지 확인하세요.
export PINOT_VERSION=1.5.1 #설치한 Pinot 버전으로 설정
export PINOT_DISTRIBUTION_DIR=${PINOT_ROOT_DIR}/build/
HADOOP_BATCH_PLUGIN=${PINOT_DISTRIBUTION_DIR}/plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-hadoop/pinot-batch-ingestion-hadoop-${PINOT_VERSION}-shaded.jar
PINOT_ALL_JAR=${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar
export HADOOP_CLIENT_OPTS="-Dplugins.dir=${PINOT_DISTRIBUTION_DIR}/plugins;${PINOT_DISTRIBUTION_DIR}/plugins-external -Dlog4j2.configurationFile=${PINOT_DISTRIBUTION_DIR}/conf/pinot-ingestion-job-log4j2.xml"
export HADOOP_CLASSPATH="${HADOOP_BATCH_PLUGIN}:${PINOT_ALL_JAR}:${HADOOP_CLASSPATH}"
hadoop jar \
${PINOT_ALL_JAR} \
org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand \
-jobSpecFile ${PINOT_DISTRIBUTION_DIR}/examples/batch/airlineStats/hadoopIngestionJobSpec.yaml
ClassNotFoundException문제 해결:plugins.dir에plugins-external을 추가하고pinot-batch-ingestion-hadoop-*-shaded.jar를HADOOP_CLASSPATH에 넣으세요. Hadoop 배치 수집을 참고하세요.
백필 수집 작업 실행 (Executing a backfill ingestion job)
특정 날짜의 백필 입력 디렉터리에 원래 수집보다 적은 파일이 포함될 수 있는 경우, LaunchDataIngestionJob은 기존 세그먼트를 완전히 대체할 수 없어요 — 백필 문서의 Edge case example를 참고하세요. 이 경우 동일한 ingestionJobSpec.yaml을 재사용하고 Pinot의 segment-lineage 메커니즘을 통해 날짜 범위의 기존 세그먼트를 대체하는 LaunchBackfillIngestionJob을 사용하세요.
bin/pinot-admin.sh LaunchBackfillIngestionJob \
-jobSpecFile examples/batch/airlineStats/ingestionJobSpec.yaml \
-startDate 2014-01-01 \
-endDate 2014-01-02
-startDate(포함)와 -endDate(제외) 인자는 yyyy-MM-dd 형식을 사용하며 UTC 일 경계로 파싱돼요.
전체 단계별 워크플로, 지원 옵션(파티션 범위 백필용
-partitionColumn/-partitionColumnValue포함), OFFLINE 전용 제약은 백필 데이터 (Backfill Data)를 참고하세요.
튜닝 (Tuning)
환경 변수 JAVA_OPTS를 설정해 수정할 수 있어요:
-Dlog4j2.configurationFile로 Log4j2 파일 위치-Dplugins.dir=/opt/pinot/plugins로 플러그인 디렉터리 위치(standalone). Spark/Hadoop 배치 작업의 경우plugins-external도 포함하는 세미콜론 구분 목록 사용:-Dplugins.dir=/opt/pinot/plugins;/opt/pinot/plugins-external-Xmx8g -Xms4G같은 JVM 속성
위 세 가지를 JAVA_OPTS에 모두 함께 구성해야 한다는 점을 참고하세요. JAVA_OPTS="-Xmx4g"만 구성하면 plugins.dir이 보통 비어 있어 작업 실패를 유발할 수 있어요.
예: standalone Docker 작업:
docker run --rm -ti -e JAVA_OPTS="-Xms8G -Dlog4j2.configurationFile=/opt/pinot/conf/pinot-admin-log4j2.xml -Dplugins.dir=/opt/pinot/plugins" --name pinot-data-ingestion-job apachepinot/pinot:latest LaunchDataIngestionJob -jobSpecFile /path/to/ingestion_job_spec.yaml
필요하면 맞춤 JAVA_OPTS도 추가할 수 있어요.