Spark를 사용해 S3에서 Parquet 파일 수집하기

Spark를 사용해 S3에서 Parquet 파일 수집하기

Pinot 사용의 주요 장점 중 하나는 플러그형 아키텍처예요. 플러그인 덕분에 실행 프레임워크, 파일시스템, 입력 형식 등 어떤 서드파티 시스템에 대한 지원도 쉽게 추가할 수 있어요.

이 튜토리얼에서는 세 가지 플러그인을 사용해 데이터를 쉽게 수집하고 pinot 클러스터에 push할 거예요. 사용할 플러그인은 다음과 같아요:

  • pinot-batch-ingestion-spark
  • pinot-s3
  • pinot-parquet

사용 가능한 모든 플러그인은 배치 수집 (Batch Ingestion), 파일 시스템 (File systems), 입력 형식 (Input formats)에서 확인할 수 있어요.

설정 (Setup)

이 튜토리얼에서는 다음 도구와 프레임워크를 사용해요:

입력 데이터 (Input Data)

먼저 수집할 입력 데이터를 얻어야 해요. 데모를 위해 작은 parquet 파일 몇 개를 만들어 S3 버킷에 업로드할 거예요. 가장 쉬운 방법은 CSV 파일을 만든 다음 parquet으로 변환하는 것이에요. CSV는 사람이 읽을 수 있어서 데모에서 실패가 발생해도 입력을 수정하기 쉬워요. 이 파일을 students.csv라고 부를게요.

timestampInEpoch,id,name,age,score
1597044264380,1,david,15,98
1597044264381,2,henry,16,97
1597044264382,3,katie,14,99
1597044264383,4,catelyn,15,96
1597044264384,5,emma,13,93
1597044264390,6,john,15,100
1597044264396,7,isabella,13,89
1597044264399,8,linda,17,91
1597044264502,9,mark,16,67
1597044264670,10,tom,14,78

이제 Spark를 사용해 위 CSV 파일에서 parquet 파일을 만들 거예요. 작은 프로그램이므로 본격적인 Spark 코드를 작성하는 대신 Spark 셸을 사용할 거예요.

scala> val df = spark.read.format("csv").option("header", true).load("path/to/students.csv")
scala> df.write.option("compression","none").mode("overwrite").parquet("/path/to/batch_input/")

이제 .parquet 파일이 /path/to/batch_input 디렉터리에서 찾을 수 있어요. 이 디렉터리를 S3에 UI를 사용하거나 다음 명령으로 업로드할 수 있어요.

aws s3 cp /path/to/batch_input s3://my-bucket/batch-input/ --recursive

스키마와 테이블 만들기 (Create Schema and Table)

수집될 데이터를 쿼리하기 위한 테이블을 만들어야 해요. pinot의 모든 테이블은 스키마와 연결돼 있어요. 구성 생성에 대한 자세한 내용은 테이블 구성 (Table configuration)과 스키마 구성 (Schema configuration)에서 확인할 수 있어요.

데모를 위해 다음 스키마와 테이블 구성을 사용할 거예요.

{
    "schemaName": "students",
    "dimensionFieldSpecs": [
        {
            "name": "id",
            "dataType": "INT"
        },
        {
            "name": "name",
            "dataType": "STRING"
        },
        {
            "name": "age",
            "dataType": "INT"
        }
    ],
    "metricFieldSpecs": [
        {
            "name": "score",
            "dataType": "INT"
        }
    ],
    "dateTimeFieldSpecs": [
        {
            "name": "timestampInEpoch",
            "dataType": "LONG",
            "format": "1:MILLISECONDS:EPOCH",
            "granularity": "1:MILLISECONDS"
        }
    ]
}
{
    "tableName": "students",
    "segmentsConfig": {
        "timeColumnName": "timestampInEpoch",
        "timeType": "MILLISECONDS",
        "replication": "1",
        "schemaName": "students"
    },
    "tableIndexConfig": {
        "invertedIndexColumns": [],
        "loadMode": "MMAP"
    },
    "tenants": {
        "broker": "DefaultTenant",
        "server": "DefaultTenant"
    },
    "tableType": "OFFLINE",
    "metadata": {}
}

이제 이 구성을 pinot에 업로드하고 빈 테이블을 만들 수 있어요. 이 목적으로 pinot-admin.sh CLI를 사용할 거예요.

pinot-admin.sh AddTable -tableConfigFile /path/to/student_table.json -schemaFile /path/to/student_schema.json -controllerHost localhost -controllerPort 9000 -exec

사용 가능한 모든 명령은 Command-Line Interface (CLI)에서 확인할 수 있어요.

이제 우리 테이블은 Pinot 데이터 탐색기에서 사용할 수 있어요.

데이터 수집 (Ingest Data)

데이터가 S3에 있고 Pinot에 테이블도 있으므로 데이터 수집 과정을 시작할 수 있어요. Pinot에서 데이터 수집은 다음 단계를 포함해요:

  • 입력에서 데이터를 읽고 압축된 세그먼트 파일 생성
  • 압축된 세그먼트 파일을 출력 위치에 업로드
  • 세그먼트 파일의 위치를 controller에 push

위치가 controller에 전달되면 controller는 서버에 세그먼트 파일을 다운로드하고 테이블을 채우도록 알릴 수 있어요.

위 단계는 Hadoop, Spark, Flink 등 원하는 분산 실행기로 수행할 수 있어요. 이 데모에서는 Apache Spark를 사용해 단계를 실행할 거예요.

Pinot은 Spark용 러너를 기본 제공해요. 따라서 사용자는 한 줄의 코드도 작성할 필요가 없어요. 제공된 인터페이스를 사용해 다른 실행기용 러너를 작성할 수도 있어요.

먼저 데이터 수집 과정에 대한 job spec 구성 파일을 만들 거예요.

executionFrameworkSpec:
  name: 'spark'
  segmentGenerationJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentGenerationJobRunner'
  segmentTarPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentTarPushJobRunner'
  segmentUriPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentUriPushJobRunner'
  segmentMetadataPushJobRunnerClassName: 'org.apache.pinot.plugin.ingestion.batch.spark.SparkSegmentMetadataPushJobRunner'
  extraConfigs:
      stagingDir: s3://my-bucket/spark/staging/
# jobType: Pinot 수집 작업 유형
# 지원되는 job 유형:
#   'SegmentCreation'
#   'SegmentTarPush'
#   'SegmentUriPush'
#   'SegmentCreationAndTarPush'
#   'SegmentCreationAndUriPush'
#   'SegmentCreationAndMetadataPush'
jobType: SegmentCreationAndMetadataPush
inputDirURI: 's3://my-bucket/path/to/batch-input/'
outputDirURI: 's3:///my-bucket/path/to/batch-output/'
overwriteOutput: true
pinotFSSpecs:
  - scheme: s3
    className: org.apache.pinot.plugin.filesystem.S3PinotFS
    configs:    
      region: 'us-west-2'
recordReaderSpec:
  dataFormat: 'parquet'
  className: 'org.apache.pinot.plugin.inputformat.parquet.ParquetRecordReader'
tableSpec:
  tableName: 'students'
pinotClusterSpecs:
  - controllerURI: 'http://localhost:9000'
pushJobSpec:
  pushParallelism: 2
  pushAttempts: 2
  pushRetryIntervalMillis: 1000

job spec에서 실행 프레임워크를 spark로 유지하고 각 단계에 적절한 러너를 구성했어요. 또한 spark 작업용 임시 stagingDir이 필요해요. 이 디렉터리는 작업이 실행된 후 정리돼요.

또한 구성에서 사용할 S3 파일시스템과 Parquet reader 구현도 제공해요. 전체 구성 목록은 수집 Job Spec (Ingestion Job Spec)을 참고할 수 있어요.

이제 spark 작업을 실행해 모든 단계를 실행하고 pinot에 데이터를 채울 수 있어요.

export PINOT_VERSION=1.4.0 #설치한 Pinot 버전으로 설정
export PINOT_DISTRIBUTION_DIR=/path/to/apache-pinot-${PINOT_VERSION}-bin

spark-submit //
--class org.apache.pinot.tools.admin.command.LaunchDataIngestionJobCommand //
--master local --deploy-mode client //
--conf "spark.driver.extraJavaOptions=-Dplugins.dir=${PINOT_DISTRIBUTION_DIR}/plugins -Dplugins.include=pinot-s3,pinot-parquet -Dlog4j2.configurationFile=${PINOT_DISTRIBUTION_DIR}/conf/pinot-ingestion-job-log4j2.xml" //
--conf "spark.driver.extraClassPath=${PINOT_DISTRIBUTION_DIR}/plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-spark/pinot-batch-ingestion-spark-${PINOT_VERSION}-shaded.jar:${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar:${PINOT_DISTRIBUTION_DIR}/plugins/pinot-file-system/pinot-s3/pinot-s3-${PINOT_VERSION}-shaded.jar:${PINOT_DISTRIBUTION_DIR}/plugins/pinot-input-format/pinot-parquet/pinot-parquet-${PINOT_VERSION}-shaded.jar" //
--conf "spark.executor.extraClassPath=${PINOT_DISTRIBUTION_DIR}/plugins-external/pinot-batch-ingestion/pinot-batch-ingestion-spark/pinot-batch-ingestion-spark-${PINOT_VERSION}-shaded.jar:${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar:${PINOT_DISTRIBUTION_DIR}/plugins/pinot-file-system/pinot-s3/pinot-s3-${PINOT_VERSION}-shaded.jar:${PINOT_DISTRIBUTION_DIR}/plugins/pinot-input-format/pinot-parquet/pinot-parquet-${PINOT_VERSION}-shaded.jar" //
local://${PINOT_DISTRIBUTION_DIR}/lib/pinot-all-${PINOT_VERSION}-jar-with-dependencies.jar -jobSpecFile /path/to/spark_job_spec.yaml

오류가 발생하면 Spark 수집 가이드의 FAQ 섹션을 살펴볼 수 있어요.

짜잔! 이제 데이터가 성공적으로 수집됐어요. Pinot의 broker에서 데이터를 쿼리해보세요.

bin/pinot-admin.sh PostQuery -brokerHost localhost -brokerPort 8000 -queryType sql -query "SELECT * FROM students LIMIT 10"

모든 것이 올바르게 수행됐다면 다음 출력을 받아야 해요.

{
  "resultTable": {
    "dataSchema": {
      "columnNames": [
        "age",
        "id",
        "name",
        "score",
        "timestampInEpoch"
      ],
      "columnDataTypes": [
        "INT",
        "INT",
        "STRING",
        "INT",
        "LONG"
      ]
    },
    "rows": [
      [ 15, 1, "david", 98, 1597044264380 ],
      [ 16, 2, "henry", 97, 1597044264381 ],
      [ 14, 3, "katie", 99, 1597044264382 ],
      [ 15, 4, "catelyn", 96, 1597044264383 ],
      [ 13, 5, "emma", 93, 1597044264384 ],
      [ 15, 6, "john", 100, 1597044264390 ],
      [ 13, 7, "isabella", 89, 1597044264396 ],
      [ 17, 8, "linda", 91, 1597044264399 ],
      [ 16, 9, "mark", 67, 1597044264502 ],
      [ 14, 10, "tom", 78, 1597044264670 ]
    ]
  },
  "exceptions": [],
  "numServersQueried": 1,
  "numServersResponded": 1,
  "numSegmentsQueried": 1,
  "numSegmentsProcessed": 1,
  "numSegmentsMatched": 1,
  "numConsumingSegmentsQueried": 0,
  "numDocsScanned": 10,
  "numEntriesScannedInFilter": 0,
  "numEntriesScannedPostFilter": 50,
  "numGroupsLimitReached": false,
  "totalDocs": 10,
  "timeUsedMs": 11,
  "segmentStatistics": [],
  "traceInfo": {},
  "minConsumingFreshnessTimeMs": 0
}

더 알아보기 (Learn more)