배치 수집 가이드

배치 수집 가이드 (Batch Ingestion Guide)

이미 파일 시스템(S3 등)에 존재하는 데이터로 테이블을 만들고 싶을 때 배치 수집을 써요. 거대한 데이터를 짧은 지연으로 쿼리하거나, 단순한 데이터 파일로 새 기능을 테스트할 때 특히 유용해요.

출처: Batch Ingestion Guide

본문

배치 수집(Batch Ingestion)으로 Apache Pinot에 데이터를 넣을 수 있어요.

배치 수집 모드 선택하기

Pinot은 여러 배치 수집 모드를 제공해요. 아래 표로 환경과 데이터 규모에 맞는 모드를 골라 보세요.

결정 가이드 (Decision Guide)

모드 (Mode) 가장 적합한 경우 (Best For) 인프라 (Infrastructure) 데이터 규모 (Data Scale) 상태 (Status)
Standalone 개발/테스트, 소규모 잡, 스크립트 파이프라인 없음 (단일 JVM) 최대 수 GB 개발 권장
Spark 3 프로덕션 배치 수집 Spark 3.x 클러스터 GB~TB+ 프로덕션 권장
Hadoop 기존 MapReduce 파이프라인 Hadoop 클러스터 GB~TB+ 레거시
Flink 스트림-우선 조직, 백필, 업서트 부트스트랩 Flink 클러스터 GB~TB+ 활성
LaunchDataIngestionJob Standalone용 CLI 래퍼 없음 최대 수 GB 편의 도구

각 모드를 언제 쓸까

Standalone은 가장 단순한 옵션으로 분산 컴퓨팅 프레임워크가 필요 없어요. 단일 JVM 프로세스에서 세그먼트 생성을 수행하므로 개발·테스트, 그리고 수 GB 정도의 소규모 프로덕션 잡에 이상적이에요. 스크립트 기반 CI/CD 파이프라인에서도 잘 맞아요.

Spark 3는 대규모 프로덕션 배치 수집에 권장되는 선택이에요. Spark 3.x 클러스터에 세그먼트 생성을 분산해서 기가바이트에서 테라바이트 이상의 데이터를 처리할 수 있어요. 새 Spark 기반 파이프라인을 만들고 있다면 이 모드를 쓰세요.

Hadoop은 MapReduce를 이용해 Hadoop 클러스터에서 세그먼트를 생성해요. 레거시로 간주되며, 기존 MapReduce 인프라와 파이프라인을 이전할 수 없을 때 주로 유용해요.

Flink는 이미 Apache Flink를 운영하는 조직에 잘 맞아요. 배치·스트림 모드를 모두 지원하며, 특히 오프라인 테이블 백필이나 업서트 테이블 부트스트랩에 유용해요. Flink 커넥터가 업서트 시맨틱에 올바르게 참여하는 파티셔닝된 세그먼트를 쓸 수 있기 때문이에요.

LaunchDataIngestionJob은 내부적으로 Standalone 러너를 호출하는 CLI 편의 래퍼예요. 셸 명령이나 cron 잡에서 커스텀 코드 없이 수집을 트리거하고 싶을 때 사용해요.

Maven 아티팩트 좌표

모든 아티팩트는 그룹 ID org.apache.pinot을 사용해요. ${pinot.version}을 사용 중인 Pinot 릴리스 버전으로 바꾸세요.

모드 (Mode) 아티팩트 ID (Artifact ID) 참고 (Notes)
Standalone pinot-batch-ingestion-standalone Pinot 바이너리 배포판에 포함
Spark 3 pinot-batch-ingestion-spark-3 plugins-external/pinot-batch-ingestion/에 위치
Hadoop pinot-batch-ingestion-hadoop plugins-external/pinot-batch-ingestion/에 위치
Flink pinot-flink-connector pinot-connectors/에 위치
Common pinot-batch-ingestion-common 모든 모드가 공유하는 라이브러리

Spark 3용 Maven 의존성 예시:

<dependency>
  <groupId>org.apache.pinot</groupId>
  <artifactId>pinot-batch-ingestion-spark-3</artifactId>
  <version>${pinot.version}</version>
</dependency>

시작하기 (Getting Started)

파일시스템에서 데이터를 수집하려면 다음 단계를 수행해요 (이 페이지에서 자세히 설명):

  1. 스키마 설정 생성
  2. 테이블 설정 생성
  3. 스키마와 테이블 설정 업로드
  4. 데이터 업로드

배치 수집은 현재 다음 메커니즘으로 데이터를 업로드할 수 있어요:

여기 standalone 로컬 처리를 사용하는 예시예요.

먼저 다음 CSV 데이터로 테이블을 만들어 볼게요.

studentID,firstName,lastName,gender,subject,score,timestampInEpoch
200,Lucy,Smith,Female,Maths,3.8,1570863600000
200,Lucy,Smith,Female,English,3.5,1571036400000
201,Bob,King,Male,Maths,3.2,1571900400000
202,Nick,Young,Male,Physics,3.6,1572418800000

스키마 설정 생성

우리 데이터에서 집계를 수행할 수 있는 유일한 컬럼은 score예요. 그리고 timestampInEpoch이 유일한 타임스탬프 컬럼이에요. 그래서 스키마에서 score는 metric으로, timestampInEpoch은 timestamp 컬럼으로 둘게요.

{
  "schemaName": "transcript",
  "dimensionFieldSpecs": [
    {
      "name": "studentID",
      "dataType": "INT"
    },
    {
      "name": "firstName",
      "dataType": "STRING"
    },
    {
      "name": "lastName",
      "dataType": "STRING"
    },
    {
      "name": "gender",
      "dataType": "STRING"
    },
    {
      "name": "subject",
      "dataType": "STRING"
    }
  ],
  "metricFieldSpecs": [
    {
      "name": "score",
      "dataType": "FLOAT"
    }
  ],
  "dateTimeFieldSpecs": [{
    "name": "timestampInEpoch",
    "dataType": "LONG",
    "format" : "1:MILLISECONDS:EPOCH",
    "granularity": "1:MILLISECONDS"
  }]
}

여기서 format과 granularity라는 두 필드도 정의했어요. format은 데이터 소스에서 타임스탬프 컬럼이 어떻게 포맷되었는지를 지정해요. 지금은 밀리초 단위이므로 1:MILLISECONDS:EPOCH로 지정했어요.

테이블 설정 생성

transcript 테이블을 정의하고 이전 단계에서 만든 스키마를 연결해요. 배치 데이터이므로 tableType을 OFFLINE으로 두어요.

{
  "tableName": "transcript",
  "tableType": "OFFLINE",
  "segmentsConfig": {
    "replication": 1,
    "timeColumnName": "timestampInEpoch",
    "timeType": "MILLISECONDS",
    "retentionTimeUnit": "DAYS",
    "retentionTimeValue": 365
  },
  "tenants": {
    "broker":"DefaultTenant",
    "server":"DefaultTenant"
  },
  "tableIndexConfig": {
    "loadMode": "MMAP"
  },
  "ingestionConfig": {
    "batchIngestionConfig": {
      "segmentIngestionType": "APPEND",
      "segmentIngestionFrequency": "DAILY"
    },
    "continueOnError": true,
    "rowTimeValueCheck": true,
    "segmentTimeValueCheck": false

  },
  "metadata": {}
}

스키마와 테이블 설정 업로드

이제 두 설정이 모두 있으니 업로드 후 다음 명령으로 테이블을 만들어요:

bin/pinot-admin.sh AddTable \
  -tableConfigFile /path/to/table-config.json \
  -schemaFile /path/to/table-schema.json -exec

[Rest API]에서 테이블 설정과 스키마를 확인해 정상 업로드되었는지 살펴보세요.

데이터 업로드

이제 Pinot에 빈 테이블이 생겼어요. 다음으로 이 빈 테이블에 CSV 파일을 업로드해요.

테이블은 여러 세그먼트로 구성돼요. 세그먼트는 다음 세 가지 방법으로 만들 수 있어요:

  • Minion 기반 수집
  • Upload API
  • 수집 잡 (Ingestion jobs)

Minion 기반 수집

SegmentGenerationAndPushTask를 참고하세요.

Upload API

소규모 파일로 빠른 수집 테스트에 쓸 수 있는 컨트롤러 API가 2개 있어요.

{% hint style="danger" %} 이 API들이 호출되면 컨트롤러가 파일을 다운로드하고 로컬에서 세그먼트를 빌드해야 합니다.

따라서 이 API는 프로덕션 환경과 대용량 입력 파일에는 적합하지 않습니다. {% endhint %}

/ingestFromFile

이 API는 주어진 파일로 세그먼트를 만들고 Pinot으로 푸시해요. 모든 단계가 컨트롤러에서 수행돼요.

사용 예시:

JSON 파일 data.json을 foo_OFFLINE 테이블에 업로드하려면 아래 명령을 사용해요.

쿼리 파라미터는 URL 인코딩이 필요해요. 예를 들어 아래 명령의 {"inputFormat":"json"} 은 %7B%22inputFormat%22%3A%22json%22%7D 로 변환해야 해요.

curl -X POST -F [email protected] \
  -H "Content-Type: multipart/form-data" \
  "http://localhost:9000/ingestFromFile?tableNameWithType=foo_OFFLINE&
  batchConfigMapStr={"inputFormat":"json"}"

batchConfigMapStr은 파일 디코딩에 필요한 추가 속성을 전달하는 데 쓸 수 있어요. 예를 들어 csv의 경우 구분자(delimiter)를 제공해야 할 수 있어요.

curl -X POST -F [email protected] \
  -H "Content-Type: multipart/form-data" \
  "http://localhost:9000/ingestFromFile?tableNameWithType=foo_OFFLINE&
batchConfigMapStr={
  "inputFormat":"csv",
  "recordReader.prop.delimiter":"|"
}"
/ingestFromURI

이 API는 주어진 URI의 파일로 세그먼트를 만들고 Pinot으로 푸시해요. 파일시스템 접근에 필요한 속성은 batchConfigMap에 제공돼요. 모든 단계가 컨트롤러에서 수행돼요.

사용 예시:

{% hint style="warning" %} /ingestFromURI는 기본적으로 원격 파일만 허용합니다. 컨트롤러 로컬의 file:///... URI와 LocalPinotFS 기반 읽기는 controller.conf에서 controller.ingestFromURI.allowLocalFileSystem=true를 명시적으로 설정하고 컨트롤러를 재시작하지 않는 한 거부됩니다. {% endhint %}

curl -X POST "http://localhost:9000/ingestFromURI?tableNameWithType=foo_OFFLINE
&batchConfigMapStr={
  "inputFormat":"json",
  "input.fs.className":"org.apache.pinot.plugin.filesystem.S3PinotFS",
  "input.fs.prop.region":"us-central",
  "input.fs.prop.accessKey":"foo",
  "input.fs.prop.secretKey":"bar"
}
&sourceURIStr=s3://test.bucket/path/to/json/data/data.json"

신뢰할 수 있는 환경에서 컨트롤러 호스트 로컬의 파일을 의도적으로 수집한다면, 위 컨트롤러 설정을 활성화하고 file:///... URI를 사용해요:

curl -X POST "http://localhost:9000/ingestFromURI?tableNameWithType=foo_OFFLINE&batchConfigMapStr={\"inputFormat\":\"json\"}&sourceURIStr=file:///var/tmp/pinot/input.json"

수집 잡 (Ingestion jobs)

세그먼트는 DataIngestionJobs라 불리는 태스크로 생성·업로드할 수 있어요. 잡에는 자체 설정도 필요하고, 이를 JobSpec이라고 불러요.

우리 CSV 파일과 테이블의 JobSpec은 다음과 같아요:

executionFrameworkSpec:
  name: 'standalone'
  segmentGenerationJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.standalone.SegmentGenerationJobRunner'
  segmentTarPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.standalone.SegmentTarPushJobRunner'
  segmentUriPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.standalone.SegmentUriPushJobRunner'
  segmentMetadataPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.standalone.SegmentMetadataPushJobRunner'

# Recommended to set jobType to SegmentCreationAndMetadataPush for production environment where Pinot Deep Store is configured  
jobType: SegmentCreationAndTarPush

inputDirURI: '/tmp/pinot-quick-start/rawdata/'
includeFileNamePattern: 'glob:**/*.csv'
outputDirURI: '/tmp/pinot-quick-start/segments/'
overwriteOutput: true
pinotFSSpecs:
  - scheme: file
    className: org.apache.pinot.spi.filesystem.LocalPinotFS
recordReaderSpec:
  dataFormat: 'csv'
  className: 'org.apache.pinot.plugin.inputformat.csv.CSVRecordReader'
  configClassName: 'org.apache.pinot.plugin.inputformat.csv.CSVRecordReaderConfig'
tableSpec:
  tableName: 'transcript'
pinotClusterSpecs:
  - controllerURI: 'http://localhost:9000'
pushJobSpec:
  pushAttempts: 2
  pushRetryIntervalMillis: 1000

자세한 내용은 Ingestion job spec을 참고하세요.

이제 transcript 테이블의 잡 스펙이 있으니 다음 명령으로 잡을 실행할 수 있어요:

bin/pinot-admin.sh LaunchDataIngestionJob \
    -jobSpecFile /tmp/pinot-quick-start/batch-job-spec.yaml

잡이 성공적으로 끝나면 [query console]로 가서 데이터를 가지고 놀아 보세요.

세그먼트 푸시 잡 유형 (Segment push job type)

Pinot 세그먼트를 업로드하는 방법은 3가지예요:

  • Segment tar push
  • Segment URI push
  • Segment metadata push

Segment tar push

이것이 원래의 기본 푸시 메커니즘이에요.

Tar push는 세그먼트가 로컬에 저장되어 있거나 PinotFS에서 InputStream으로 열 수 있어야 해요. 그래서 전체 세그먼트 tar 파일을 컨트롤러로 스트리밍할 수 있어요.

푸시 잡이 수행할 일:

  1. 전체 세그먼트 tar 파일을 Pinot 컨트롤러에 업로드.

Pinot 컨트롤러가 수행할 일:

  1. 세그먼트를 컨트롤러 세그먼트 디렉터리(로컬 또는 모든 PinotFS)에 저장.
  2. 세그먼트 메타데이터 추출.
  3. 테이블에 세그먼트 추가.

Segment URI push

이 푸시 방식은 딥스토어에 세그먼트 tar 파일이 저장되어 있고, 전역에서 접근 가능한 세그먼트 tar URI가 있어야 해요.

URI push는 클라이언트 쪽이 가볍고, 컨트롤러 쪽은 tar push와 동일한 작업이 필요해요.

푸시 잡이 수행할 일:

  1. 이 세그먼트 tar URI를 Pinot 컨트롤러에 POST.

Pinot 컨트롤러가 수행할 일:

  1. URI에서 세그먼트를 다운로드해 컨트롤러 세그먼트 디렉터리(로컬 또는 모든 PinotFS)에 저장.
  2. 세그먼트 메타데이터 추출.
  3. 테이블에 세그먼트 추가.

Segment metadata push

이 푸시 방식도 딥스토어에 세그먼트 tar 파일이 저장되어 있고, 전역에서 접근 가능한 세그먼트 tar URI가 있어야 해요.

Metadata push는 컨트롤러 쪽이 가볍고, 컨트롤러 쪽에서 딥스토어 다운로드가 발생하지 않아요.

푸시 잡이 수행할 일:

  1. URI를 기반으로 세그먼트 다운로드.
  2. 메타데이터 추출.
  3. 메타데이터를 Pinot 컨트롤러에 업로드.

Pinot 컨트롤러가 수행할 일:

  1. 메타데이터를 기반으로 테이블에 세그먼트 추가.

4. Segment Metadata Push with copyToDeepStore

이 방식은 딥스토어로 쓰지 않는 위치에 세그먼트를 푸시하는 경우를 위해 원래 Segment Metadata Push를 확장해요. 수집 잡은 여전히 메타데이터 푸시를 하되, Pinot 컨트롤러에 세그먼트를 딥스토어로 복사해 달라고 요청할 수 있어요. 이 용례는 보통 수집 잡이 딥스토어에 직접 접근할 수 없으면서 그 효율성 때문에 메타데이터 푸시를 쓰길 원할 때, 스테이징 위치에 세그먼트를 임시로 두는 방식으로 발생해요.

참고: 스테이징 위치와 딥스토어는 같은 스토리지 스킴을 사용해야 해요 (예: 둘 다 s3). 복사가 PinotFS.copyDir 인터페이스를 통해 이루어지므로 그걸 가정하기 때문이에요. 또한 이 복사는 스토리지 시스템 쪽에서 수행되므로 세그먼트가 Pinot 컨트롤러를 통과할 필요가 전혀 없어요.

이를 동작시키려면 Pinot 컨트롤러에 스테이징 위치 접근 권한을 부여하세요. 예를 들어 AWS에서는 컨트롤러 EC2 인스턴스에 다음과 같은 접근 정책을 추가해야 할 수 있어요:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Action": "s3:ListAllMyBuckets",
            "Resource": "*"
        },
        {
            "Effect": "Allow",
            "Action": "s3:*",
            "Resource": [
                "arn:aws:s3:::metadata-push-staging",
                "arn:aws:s3:::metadata-push-staging/*"
            ]
        }
    ]
}

그다음 메타데이터 푸시에 다음처럼 설정 하나를 추가해 사용해요:

...
jobType: SegmentCreationAndMetadataPush
...
outputDirURI: 's3://metadata-push-staging/stagingDir/'
...
pushJobSpec:
  copyToDeepStoreForMetadataPush: true
...

일관된 데이터 푸시와 롤백 (Consistent data push and rollback)

Pinot은 세그먼트 단위로 원자적 업데이트를 지원해요. 즉 여러 세그먼트로 구성된 데이터를 테이블에 푸시할 때 세그먼트가 하나씩 교체되므로, 이 업로드 단계에서 브로커에 보내지는 쿼리는 이전 데이터와 새 데이터가 섞여 일관적이지 않은 결과를 만들어낼 수 있어요.

이 기능을 활성화하는 방법은 consistent-push-and-rollback.md를 참고하세요.

세그먼트 페처 (Segment fetchers)

외부 시스템(Hadoop/spark 등)에서 Pinot 세그먼트 파일을 만들 때, 그 데이터를 Pinot 컨트롤러와 서버에 푸시하는 방법은 여러 가지예요:

  1. 공유 NFS에 세그먼트를 푸시하고 Pinot이 그 NFS 위치에서 세그먼트 파일을 가져오게 함. Segment URI Push 참고.
  2. 웹 서버에 세그먼트를 푸시하고 Pinot이 HTTP/HTTPS 링크로 웹 서버에서 세그먼트 파일을 가져오게 함. Segment URI Push 참고.
  3. PinotFS(HDFS/S3/GCS/ADLS)에 세그먼트를 푸시하고 Pinot이 PinotFS URI에서 세그먼트 파일을 가져오게 함. Segment URI Push와 Segment Metadata Push 참고.
  4. 다른 시스템에 세그먼트를 푸시하고 자체 세그먼트 페처를 구현해 그 시스템에서 데이터를 가져오게 함.

처음 세 가지 옵션은 Pinot 패키지에 기본으로 제공돼요. 원격 잡이 Pinot 컨트롤러에 파일 URI를 보내기만 하면, 컨트롤러가 파일을 받아 적절한 Pinot 서버와 브로커에 배정해요. PinotFS 지원을 활성화하려면 PinotFS 설정과 적절한 Hadoop 의존성을 제공해야 해요.

영속성 (Persistence)

기본적으로 Pinot에는 스토리지 계층이 없으므로, 시스템 크래시 시 전송된 데이터가 저장되지 않아요. 생성된 세그먼트를 영속적으로 저장하려면 컨트롤러와 서버 설정을 변경해 딥스토리지를 추가해야 해요. 관련 정보와 설정은 파일 시스템 (File systems)을 참고하세요.

튜닝 (Tuning)

Standalone

Pinot은 Java로 작성되었으므로 세그먼트 러너 잡을 튜닝하기 위해 다음과 같은 기본 Java 설정을 지정할 수 있어요.

  • Log4j2 파일 위치는 -Dlog4j2.configurationFile
  • 플러그인 디렉터리 위치는 -Dplugins.dir=/opt/pinot/plugins
  • JVM 속성, 예: -Xmx8g -Xms4G

Docker를 사용한다면 JAVA_OPTS 변수 아래에 위 설정을 지정할 수 있어요.

Hadoop

-D mapreduce.map.memory.mb=8192를 설정해 Hadoop 잡 제출 시 매퍼 메모리 크기를 지정할 수 있어요.

Hadoop 배치 플러그인은 plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-hadoop/에 있어요. plugins.dir을 plugins와 plugins-external 둘 다(세미콜론 구분) 가리키고, 셰이딩된 Hadoop 플러그인 jar를 HADOOP_CLASSPATH에 넣으세요. Hadoop 참고.

Spark

Spark 잡 제출 시 세그먼트 생성 메모리를 튜닝하기 위해 spark.executor.memory 설정을 추가할 수 있어요.

Spark 3 배치 플러그인은 plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-spark-3/에 있어요. pinot-all에는 포함되지 않아요. plugins.dir을 plugins와 plugins-external 둘 다 가리키고, 셰이딩된 Spark 플러그인 jar를 Spark 클래스패스에 넣으세요. Spark 참고.

더 알아보기 (Learn more)