Pinot의 딥 스토리지로 S3 사용하기

Pinot의 딥 스토리지로 S3 사용하기

Pinot을 설정해 S3를 딥 스토어로 사용하는 방법을 알려드려요.

출처: 문서

본문

아래 명령은 pinot distribution 바이너리를 기반으로 해요.

Pinot 클러스터 설정

S3를 딥 스토어로 사용하도록 Pinot을 설정하려면 Controller와 Server에 추가 config를 넣어야 해요.

Controller 시작

아래는 샘플 controller.conf 파일이에요.

controller.data.dir을 s3 버킷으로 구성하세요. 업로드된 모든 세그먼트가 거기에 저장돼요.

그리고 s3를 pinot 저장소로 추가하세요:

pinot.controller.storage.factory.class.s3=org.apache.pinot.plugin.filesystem.S3PinotFS
pinot.controller.storage.factory.s3.region=us-west-2

AWS 자격증명에 관해서는 DefaultAWSCredentialsProviderChain의 관례를 따르세요.

AccessKey와 Secret은 다음을 사용해 지정할 수 있어요:

  • 환경 변수 - AWS_ACCESS_KEY_ID와 AWS_SECRET_ACCESS_KEY(.NET을 제외한 모든 AWS SDK와 CLI가 인식하므로 권장), 또는 AWS_ACCESS_KEY와 AWS_SECRET_KEY(Java SDK만 인식)
  • Java 시스템 속성 - aws.accessKeyId와 aws.secretKey
  • 모든 AWS SDK와 AWS CLI가 공유하는 기본 위치(~/.aws/credentials)의 자격증명 프로필 파일
  • pinot config 파일에 AWS 자격증명 구성(예: config 파일에 pinot.controller.storage.factory.s3.accessKey와 pinot.controller.storage.factory.s3.secretKey 설정) (권장하지 않음)
pinot.controller.storage.factory.s3.accessKey=****************LFVX
pinot.controller.storage.factory.s3.secretKey=****************gfhz

pinot.controller.segment.fetcher.protocols에 s3를 추가하고 pinot.controller.segment.fetcher.s3.class를 org.apache.pinot.common.utils.fetcher.PinotFSSegmentFetcher로 설정하세요.

controller.data.dir=s3://my.bucket/pinot-data/pinot-s3-example/controller-data
controller.local.temp.dir=/tmp/pinot-tmp-data/
controller.zk.str=localhost:2181
controller.host=127.0.0.1
controller.port=9000
controller.helix.cluster.name=pinot-s3-example
pinot.controller.storage.factory.class.s3=org.apache.pinot.plugin.filesystem.S3PinotFS
pinot.controller.storage.factory.s3.region=us-west-2

pinot.controller.segment.fetcher.protocols=file,http,s3
pinot.controller.segment.fetcher.s3.class=org.apache.pinot.common.utils.fetcher.PinotFSSegmentFetcher

버킷 소유자에게 전체 제어권을 부여하려면 config에 다음을 추가하세요:

pinot.controller.storage.factory.s3.disableAcl=false

그런 다음 다음으로 pinot controller를 시작하세요:

bin/pinot-admin.sh StartController -configFileName conf/controller.conf

Broker 시작

Broker는 간단하므로 기본값으로 시작하면 돼요:

bin/pinot-admin.sh StartBroker -zkAddress localhost:2181 -clusterName pinot-s3-example

Server 시작

아래는 샘플 server.conf 파일이에요.

controller config와 유사하게 pinot server에도 s3 config를 설정하세요.

pinot.server.netty.port=8098
pinot.server.adminapi.port=8097
pinot.server.instance.dataDir=/tmp/pinot-tmp/server/index
pinot.server.instance.segmentTarDir=/tmp/pinot-tmp/server/segmentTars

pinot.server.storage.factory.class.s3=org.apache.pinot.plugin.filesystem.S3PinotFS
pinot.server.storage.factory.s3.region=us-east-1
pinot.server.storage.factory.s3.accessKey=myAccessKeyChangeMe
pinot.server.storage.factory.s3.secretKey=mySecretKeyChangeMe
pinot.server.storage.factory.s3.disableAcl=false
pinot.server.storage.factory.s3.endpoint=http://minio:9000
pinot.server.segment.fetcher.protocols=file,http,s3
pinot.server.segment.fetcher.s3.class=org.apache.pinot.common.utils.fetcher.PinotFSSegmentFetcher

버킷 소유자에게 전체 제어권을 부여하려면 config에 다음을 추가하세요:

pinot.controller.storage.factory.s3.disableAcl=false

그런 다음 다음으로 pinot server를 시작하세요:

bin/pinot-admin.sh StartServer -configFileName conf/server.conf -zkAddress localhost:2181 -clusterName pinot-s3-example

테이블 설정

이 데모에서는 airlineStats 테이블을 예시로 사용해요.

다음 명령으로 테이블을 만드세요:

bin/pinot-admin.sh AddTable  -schemaFile examples/batch/airlineStats/airlineStats_schema.json -tableConfigFile examples/batch/airlineStats/airlineStats_offline_table_config.json -exec

수집 작업 설정 (Ingestion Jobs)

Standalone 작업

특정 변경 사항이 있는 샘플 standalone 수집 job spec은 다음과 같아요:

  • jobType은 SegmentCreationAndUriPush

  • inputDirURI는 s3 위치 **s3://my.bucket/batch/airlineStats/rawdata/**로 설정

  • outputDirURI는 s3 위치 s3://my.bucket/output/airlineStats/segments로 설정

  • pinotFSSpecs 아래에 새 PinotFs 추가

    - scheme: s3
      className: org.apache.pinot.plugin.filesystem.S3PinotFS
      configs:
        region: 'us-west-2'
    
  • 라이브러리 버전 < 0.6.0인 경우 segmentUriPrefix를 [scheme]://[bucket.name](예: s3://my.bucket)으로 설정하고, 버전 0.6.0부터는 빈 문자열을 넣거나 segmentUriPrefix를 무시하면 돼요.

샘플 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: SegmentCreationAndUriPush
# inputDirURI: 입력 데이터의 루트 디렉터리, PinotFS에 구성된 scheme 있어야 함
inputDirURI: 's3://my.bucket/batch/airlineStats/rawdata/'

# includeFileNamePattern: 포함할 파일 이름 패턴, 전역 glob 패턴 지원
# 샘플 사용법:
#   'glob:*.avro'는 inputDirURI 바로 아래의 avro 파일만 포함, 하위 디렉터리는 제외
#   'glob:**/*.avro'는 inputDirURI 아래의 모든 avro 파일을 재귀적으로 포함
includeFileNamePattern: 'glob:**/*.avro'

# excludeFileNamePattern: 제외할 파일 이름 패턴, 전역 glob 패턴 지원
# 샘플 사용법:
#   'glob:*.avro'는 inputDirURI 바로 아래의 avro 파일만 제외, 하위 디렉터리는 제외
#   'glob:**/*.avro'는 inputDirURI 아래의 모든 avro 파일을 재귀적으로 제외
# _excludeFileNamePattern: ''

# outputDirURI: 출력 세그먼트의 루트 디렉터리, PinotFS에 구성된 scheme 있어야 함
outputDirURI: 's3://my.bucket/examples/output/airlineStats/segments'

# overwriteOutput: 기존 출력 세그먼트 덮어쓰기 여부
overwriteOutput: true

# 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


  - scheme: s3
    className: org.apache.pinot.plugin.filesystem.S3PinotFS
    configs:
      region: 'us-west-2'

# 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'

# pinotClusterSpecs: Pinot 클러스터 접근 지점 정의
pinotClusterSpecs:
  - # controllerURI: 테이블/스키마 정보와 데이터 push에 사용
    # 예: http://localhost:9000
    controllerURI: 'http://localhost:9000'

# pushJobSpec: 세그먼트 push 작업 관련 구성 정의
pushJobSpec:

  # pushAttempts: push 작업 시도 횟수, 기본값 1(재시도 없음)
  pushAttempts: 2

  # pushRetryIntervalMillis: 재시도 대기 ms, 기본값 1초
  pushRetryIntervalMillis: 1000

  # Pinot 버전 < 0.6.0의 경우 [scheme]://[bucket.name]을 접두사로 사용
  # 예: s3://my.bucket
  segmentUriPrefix: 's3://my.bucket'
  segmentUriSuffix: ''

샘플 작업 출력:

bin/pinot-admin.sh LaunchDataIngestionJob -jobSpecFile  ~/temp/pinot/pinot-s3-test/ingestionJobSpec.yaml

출력 로그는 SegmentGenerationJobSpec 인쇄, file 및 s3 scheme용 PinotFS 초기화, 그리고 각 세그먼트 URI(airlineStats_OFFLINE_...tar.gz)를 controller의 /v2/segments로 전송하는 과정을 보여줘요. 각 요청은 controller에 성공적으로 업로드되어 200 응답과 {"status":"Successfully uploaded segment: airlineStats_OFFLINE_..."} 메시지를 반환해요.

Spark 작업

Spark 클러스터 설정 (이미 있으면 건너뛰기)

이 페이지를 따라 로컬 spark 클러스터를 설정하세요.

Spark 작업 제출

아래는 샘플 Spark 수집 작업이에요.

# executionFrameworkSpec: 실행할 수집 작업 정의
executionFrameworkSpec:

  # name: 실행 프레임워크 이름
  name: 'spark'

  # segmentGenerationJobRunnerClassName: org.apache.pinot.spi.batch.ingestion.runner.SegmentGenerationJobRunner 인터페이스 구현 클래스 이름
  segmentGenerationJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentGenerationJobRunner'

  # segmentTarPushJobRunnerClassName: org.apache.pinot.spi.batch.ingestion.runner.SegmentTarPushJobRunner 인터페이스 구현 클래스 이름
  segmentTarPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentTarPushJobRunner'

  # segmentUriPushJobRunnerClassName: org.apache.pinot.spi.batch.ingestion.runner.SegmentUriPushJobRunner 인터페이스 구현 클래스 이름
  segmentUriPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentUriPushJobRunner'

# jobType: Pinot 수집 작업 유형
# 지원되는 job 유형:
#   'SegmentCreation'
#   'SegmentTarPush'
#   'SegmentUriPush'
#   'SegmentCreationAndTarPush'
#   'SegmentCreationAndUriPush'
jobType: SegmentCreationAndUriPush

# inputDirURI: 입력 데이터의 루트 디렉터리, PinotFS에 구성된 scheme 있어야 함
inputDirURI: 's3://my.bucket/batch/airlineStats/rawdata/'

# includeFileNamePattern: 포함할 파일 이름 패턴, 전역 glob 패턴 지원
includeFileNamePattern: 'glob:**/*.avro'

# outputDirURI: 출력 세그먼트의 루트 디렉터리, PinotFS에 구성된 scheme 있어야 함
outputDirURI: 's3://my.bucket/examples/output/airlineStats/segments'

# overwriteOutput: 기존 출력 세그먼트 덮어쓰기 여부
overwriteOutput: true

# pinotFSSpecs: 모든 관련 Pinot 파일 시스템 정의
pinotFSSpecs:

  - scheme: file
    className: org.apache.pinot.plugin.filesystem.HadoopPinotFS
    configs:
      'hadoop.conf.path': ''

  - scheme: s3
    className: org.apache.pinot.plugin.filesystem.S3PinotFS
    configs:
      region: 'us-west-2'

# 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'
  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: push 작업 병렬도, 기본값 1
  pushParallelism: 2

  # pushAttempts: push 작업 시도 횟수, 기본값 1(재시도 없음)
  pushAttempts: 2

  # pushRetryIntervalMillis: 재시도 대기 ms, 기본값 1초
  pushRetryIntervalMillis: 1000

수집 작업으로 spark 작업을 제출하세요:

${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 -Dlog4j2.configurationFile=${PINOT_DISTRIBUTION_DIR}/conf/pinot-ingestion-job-log4j2.xml" \
  --conf "spark.driver.extraClassPath=${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar" \
  local://${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar \
  -jobSpecFile ${PINOT_DISTRIBUTION_DIR}/examples/sparkIngestionJobSpec.yaml

샘플 결과/스냅샷

controller용 s3 위치의 샘플 스냅샷은 아래와 같아요.

Sample S3 Controller Storage

PropertyStore의 샘플 다운로드 URI가 아래에 있어요. 세그먼트 다운로드 URI가 s3://로 시작할 것으로 기대해요.

Sample segment download URI in PropertyStore

더 알아보기 (Learn more)