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의 이후 버전에서 변경될 것)