Spark-Pinot 커넥터

Spark-Pinot 커넥터 (Spark-Pinot Connector)

Spark-Pinot 커넥터를 사용해 Pinot에서 데이터를 읽고 쓰세요.

Spark-pinot 커넥터로 Pinot에서 데이터를 읽어요.

상세 읽기 모델 문서는 여기: spark-pinot-connector-read-model

쓰기 모델은 실험적이며 문서는 여기: spark-pinot-connector-write-model

출처: 문서

본문

기능

  • realtime, offline 또는 hybrid 테이블 쿼리
  • 분산, 병렬 스캔
  • gRPC를 사용한 스트리밍 읽기(선택)
  • PQL 대신 SQL 지원
  • 성능 최적화를 위한 컬럼 및 필터 푸시다운
  • hybrid 테이블의 경우 realtime과 offline 세그먼트 간 겹침이 정확히 한 번 쿼리됨
  • 스키마 발견
    • 동적 추론
    • case class의 정적 분석
  • 쿼리 옵션 지원
  • 보안 연결을 위한 HTTPS/TLS 지원

Quick Start

import org.apache.spark.sql.SparkSession

val spark: SparkSession = SparkSession
      .builder()
      .appName("spark-pinot-connector-test")
      .master("local")
      .getOrCreate()

import spark.implicits._

val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .load()
  .filter($"DestStateName" === "Florida")

data.show(100)

보안 구성

Pinot 1.5.0부터 사용 가능.

통합 스위치 또는 명시적 플래그로 HTTP와 gRPC를 모두 보호할 수 있어요.

  • 통합: secureMode=true를 설정해 HTTPS와 gRPC TLS를 함께 활성화(권장)
  • 명시적: REST에 useHttps, gRPC에 grpc.use-plain-text=false

빠른 예시

// 통합 보안 모드 (HTTPS + gRPC TLS 기본 활성화)
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("secureMode", "true")
  .load()

// 명시적 HTTPS 전용 (gRPC는 기본 plaintext 유지)
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("useHttps", "true")
  .load()

// 명시적 gRPC TLS 전용 (REST는 기본 HTTP 유지)
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("grpc.use-plain-text", "false")
  .load()

HTTPS 구성

HTTPS가 활성화되면(secureMode=true 또는 useHttps=true로) 필요에 따라 keystore/truststore를 구성할 수 있어요.

val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("useHttps", "true")
  .option("keystorePath", "/path/to/keystore.jks")
  .option("keystorePassword", "keystorePassword")
  .option("truststorePath", "/path/to/truststore.jks")
  .option("truststorePassword", "truststorePassword")
  .load()

HTTPS 구성 옵션

Option Description Required Default
secureMode HTTPS와 gRPC TLS를 활성화하는 통합 스위치 No false
useHttps HTTPS 연결 활성화 (REST에 대해 secureMode 재정의) No false
keystorePath 클라이언트 keystore 파일 경로 (JKS 형식) No None
keystorePassword keystore 비밀번호 No None
truststorePath truststore 파일 경로 (JKS 형식) No None
truststorePassword truststore 비밀번호 No None

참고: HTTPS가 활성화될 때 truststore가 제공되지 않으면 커넥터는 모든 인증서를 신뢰해요(프로덕션 사용에 권장되지 않음).

인증 지원

Pinot 1.5.0부터 사용 가능.

커넥터는 Pinot 클러스터에 대한 보안 접근을 위해 사용자 지정 인증 헤더를 지원해요.

// Bearer token 인증 사용
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("authToken", "my-jwt-token")  // 자동으로 "Authorization: Bearer ***" 추가
  .load()

// 사용자 지정 인증 헤더 사용
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("authHeader", "Authorization")
  .option("authToken", "Bearer my-custom-token")
  .load()

// API key 인증 사용
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("authHeader", "X-API-Key")
  .option("authToken", "my-api-key")
  .load()

인증 구성 옵션

Option Description Required Default
authHeader 사용자 지정 인증 헤더 이름 No Authorization (authToken 제공 시)
authToken 인증 토큰/값 No None

참고: authHeader 없이 authToken만 제공되면 커넥터는 자동으로 Authorization: Bearer <token>을 사용해요.

Pinot Proxy 지원

Pinot 1.5.0부터 사용 가능.

커넥터는 proxy가 노출된 유일한 엔드포인트인 보안 클러스터 접근을 위한 Pinot Proxy를 지원해요. proxy가 활성화되면 controller/broker에 대한 모든 HTTP 요청과 server에 대한 gRPC 요청이 proxy를 통해 라우팅돼요.

Proxy 구성 예시

// 기본 proxy 구성
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("controller", "pinot-proxy:8080")  // Proxy endpoint
  .option("proxy.enabled", "true")
  .load()

// 인증이 있는 Proxy
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("controller", "pinot-proxy:8080")
  .option("proxy.enabled", "true")
  .option("authToken", "my-proxy-token")
  .load()

// gRPC 구성이 있는 Proxy
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("controller", "pinot-proxy:8080")
  .option("proxy.enabled", "true")
  .option("grpc.proxy-uri", "pinot-proxy:8094")  // gRPC proxy endpoint
  .load()

Proxy 구성 옵션

Option Description Required Default
proxy.enabled controller 및 broker 요청에 Pinot Proxy 사용 No false

참고: proxy가 활성화되면 커넥터는 요청을 실제 Pinot 서비스로 라우팅하기 위해 FORWARD_HOST와 FORWARD_PORT 헤더를 추가해요.

gRPC 구성

Pinot 1.5.0부터 사용 가능.

커넥터는 Pinot 서버와의 보안 및 최적화된 통신을 위한 포괄적인 gRPC 구성을 지원해요.

gRPC 구성 예시

// 기본 gRPC 구성
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("grpc.port", "8091")
  .option("grpc.max-inbound-message-size", "256000000")  // 256MB
  .load()

// gRPC with TLS (명시적)
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("grpc.use-plain-text", "false")
  .option("grpc.tls.keystore-path", "/path/to/grpc-keystore.jks")
  .option("grpc.tls.keystore-password", "keystore-password")
  .option("grpc.tls.truststore-path", "/path/to/grpc-truststore.jks")
  .option("grpc.tls.truststore-password", "truststore-password")
  .load()

// gRPC with proxy
val data = spark.read
  .format("pinot")
  .option("table", "airlineStats")
  .option("tableType", "offline")
  .option("proxy.enabled", "true")
  .option("grpc.proxy-uri", "pinot-proxy:8094")
  .load()

gRPC 구성 옵션

Option Description Required Default
grpc.port Pinot gRPC 포트 No 8090
grpc.max-inbound-message-size gRPC 클라이언트 초기화 시 최대 인바운드 메시지 바이트 No 128MB
grpc.use-plain-text gRPC 통신에 plain text 사용 (gRPC에 대해 secureMode 재정의) No true
grpc.tls.keystore-type gRPC 연결용 TLS keystore 유형 No JKS
grpc.tls.keystore-path gRPC 연결용 TLS keystore 파일 위치 No None
grpc.tls.keystore-password TLS keystore 비밀번호 No None
grpc.tls.truststore-type gRPC 연결용 TLS truststore 유형 No JKS
grpc.tls.truststore-path gRPC 연결용 TLS truststore 파일 위치 No None
grpc.tls.truststore-password TLS truststore 비밀번호 No None
grpc.tls.ssl-provider SSL provider No JDK
grpc.proxy-uri Pinot Rest Proxy gRPC 엔드포인트 URI No None

참고: gRPC를 proxy와 함께 사용할 때 커넥터는 올바른 요청 라우팅을 위해 FORWARD_HOST와 FORWARD_PORT 메타데이터 헤더를 자동으로 추가해요.

spark-shell로 예시 실행

https://github.com/apache/pinot/tree/master/pinot-connectors/pinot-spark-3-connector/examples 아래에 예시가 있어요.

사전 요구 사항

  • Apache Spark 3.x 설치 및 PATH에 spark-shell 사용 가능.

  • PINOT_HOME 환경 변수 설정:

    export PINOT_HOME=/path/to/pinot
    
  • Pinot Spark 3 Connector shaded JAR 빌드 및 다음 위치에 준비:

    $PINOT_HOME/pinot-connectors/pinot-spark-3-connector/target/pinot-spark-3-connector-*-shaded.jar
    
  • 예시 Scala 스크립트 위치:

    $PINOT_HOME/pinot-connectors/pinot-spark-3-connector/examples/read_pinot_from_proxy_with_auth_token.scala
    

Pinot Proxy에서 읽는 Scala 스크립트

import org.apache.spark.sql.SparkSession

val spark = SparkSession.builder().appName("read-pinot-airlineStats").master("local[*]").getOrCreate()

val df = spark.read.
  format("org.apache.pinot.connector.spark.v3.datasource.PinotDataSource").
  option("table", "myTable").
  option("tableType", "offline").
  option("controller", "pinot-proxy:8080").
  option("secureMode", "true").
  option("authToken", "st-xxxxxxx").
  option("proxy.enabled", "true").
  option("grpc.proxy-uri", "pinot-proxy:8094").
  option("useGrpcServer", "true").
  load()

println("Schema:")
df.printSchema()

println("Sample rows:")
df.show(10, truncate = false)

println(s"Total rows: ${df.count()}")

spark.stop()

spark-shell로 실행

다음 명령으로 spark-shell에서 예시를 실행하세요.

spark-shell 
    --master 'local[*]' \
    --name read-pinot \
    --jars "$PINOT_HOME/pinot-connectors/pinot-spark-3-connector/target/pinot-spark-3-connector-*-shaded.jar" < "$PINOT_HOME/pinot-connectors/pinot-spark-3-connector/examples/read_pinot_from_proxy_with_auth_token.scala"

샘플 출력

spark-shell --master 'local[*]' --name read-pinot --jars "$PINOT_HOME/pinot-connectors/pinot-spark-3-connector/target/pinot-spark-3-connector-*-shaded.jar" < "$PINOT_HOME/pinot-connectors/pinot-spark-3-connector/examples/read_pinot_from_proxy_with_auth_token.scala"

25/09/04 07:59:29 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
Spark context Web UI available at http://xiang-mac-home.wyvern-sun.ts.net:4040
Spark context available as 'sc' (master = local[*], app id = local-1756997971428).
Spark session available as 'spark'.
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 3.5.1
      /_/

Using Scala version 2.13 (default) (OpenJDK 64-Bit Server VM, Java 17.0.15)
Type in expressions to have them evaluated.
Type :help for more information.

scala> import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.SparkSession

scala>

scala> val spark = SparkSession.builder().appName("read-pinot-table").master("local[*]").getOrCreate()
25/09/04 07:59:35 WARN SparkSession: Using an existing Spark session; only runtime SQL configurations will take effect.
spark: org.apache.spark.sql.SparkSession = org.apache.spark.sql.SparkSession@c377641

scala>

scala> val df = spark.read.
     |   format("org.apache.pinot.connector.spark.v3.datasource.PinotDataSource").
     |   option("table", "api_gateway_agg_monthly").
     |   option("tableType", "REALTIME").
     |   option("controller", "pinot.xxx.yyy.startree.cloud").
     |   option("broker", "broker.pinot.xxx.yyy.startree.cloud").
     |   option("secureMode", "true").
     |   option("authToken", "st-xxx-yyy").
     |   option("proxy.enabled", "true").
     |   option("grpc.proxy-uri", "proxy-grpc.pinot.xxx.yyy.startree.cloud").
     |   option("useGrpcServer", "true").
     |   load()
25/09/04 07:59:35 WARN HttpUtils: No truststore configured, trusting all certificates (not recommended for production)
df: org.apache.spark.sql.DataFrame = [api_calls_count: bigint, developer_account_id: string ... 1 more field]

scala>

scala> println("Schema:")
Schema:

scala> df.printSchema()
root
 |-- api_calls_count: long (nullable = true)
 |-- developer_account_id: string (nullable = true)
 |-- monthsSinceEpoch: long (nullable = true)

scala>

scala> println("Sample rows:")
Sample rows:

scala> df.show(10, truncate = false)
25/09/04 07:59:39 WARN HttpUtils: No truststore configured, trusting all certificates (not recommended for production)
+---------------+------------------------------------+----------------+
|api_calls_count|developer_account_id                |monthsSinceEpoch|
+---------------+------------------------------------+----------------+
|276            |000e2e63-12ef-e353-af76-6fe98d2e8747|1748736000000   |
...
|277            |002f40b0-409c-b3e8-bb69-049b3e321589|1748736000000   |
+---------------+------------------------------------+----------------+
only showing top 10 rows

scala>

scala> println(s"Total rows: ${df.count()}")
25/09/04 08:00:38 WARN HttpUtils: No truststore configured, trusting all certificates (not recommended for production)
Total rows: 60000

scala>

scala> spark.stop()

scala> :quit

spark-submit으로 예시 실행

예시를 로컬 Pinot 클러스터를 시작해 독립형 모드로(예: IDE 사용) 로컬 실행할 수 있어요. docs 참조.

다음 명령으로 클러스터 모드에서도 테스트를 실행할 수 있어요.

export SPARK_CLUSTER=<YOUR_YARN_OR_SPARK_CLUSTER>

# ExampleSparkPinotConnectorTest를 편집해 `.master("local")`를 제거하고 이 명령 실행 전에 jar 재빌드
spark-submit \
    --class org.apache.pinot.connector.spark.v3.datasource.ExampleSparkPinotConnectorTest \
    --jars ./target/pinot-spark-3-connector-1.3.0-shaded.jar \
    --master $SPARK_CLUSTER \
    --deploy-mode cluster \
  ./target/pinot-spark-3-connector-1.3.0-tests.jar

이 예시는 Pinot Spark 3 Connector를 사용해 인증 토큰 지원과 함께 proxy를 통해 Pinot 클러스터에서 데이터를 읽는 방법을 보여줘요.

보안 모범 사례

Pinot 1.5.0부터 사용 가능.

프로덕션 HTTPS 구성

  • 프로덕션 환경에서 항상 HTTPS 사용
  • 적절한 파일 권한으로 보안 위치에 인증서 저장
  • 유효한 truststore로 올바른 인증서 검증 사용
  • 인증서 정기적으로 회전

프로덕션 인증

  • 최소 필요 권한으로 서비스 계정 사용
  • 인증 토큰을 보안하게 저장(환경 변수, 비밀 관리 시스템)
  • 토큰 회전 정책 구현
  • 인증 실패 모니터링

프로덕션 gRPC 구성

  • 프로덕션에서 gRPC 통신에 TLS 활성화
  • 가능하면 인증서 기반 인증 사용
  • 데이터에 따라 적절한 메시지 크기 제한 구성
  • 고처리량 시나리오에 연결 풀링 사용

향후 작업

  • 읽기 작업에 대한 통합 테스트 추가
  • 쓰기 지원 추가(pinot 세그먼트 쓰기 로직은 pinot의 이후 버전에서 변경될 것)

더 알아보기 (Learn more)