Spark Pinot 커넥터 읽기 모델

Spark Pinot 커넥터 읽기 모델 (Read Model)

커넥터는 offline, hybrid, realtime 테이블을 스캔할 수 있어요. 기본 두 옵션인 <table 및 tableType> 파라미터는 아래처럼 주어져야 해요.

  • offline 테이블: table: tbl, tableType: OFFLINE or offline
  • realtime 테이블: table: tbl, tableType: REALTIME or realtime
  • hybrid 테이블: table: tbl, tableType: HYBRID or hybrid

출처: 문서

본문

스캔 예시:

val df = spark.read
      .format("pinot")
      .option("table", "airlineStats")
      .option("tableType", "offline")
      .load()

사용자 지정 스키마를 직접 지정할 수 있어요. 스키마를 지정하지 않으면 커넥터는 Pinot controller에서 테이블 스키마를 읽고 Spark 스키마로 변환해요.

아키텍처

커넥터는 Pinot Servers에서 직접 데이터를 읽어요. 이 작업을 위해 커넥터는 먼저 주어진 필터(필터 푸시다운이 활성화된 경우)와 컬럼으로 쿼리를 만든 다음 생성된 쿼리에 대한 라우팅 테이블을 찾아요. 라우팅 테이블과 segmentsPerSplit(자세한 설명은 아래 정의)에 기반해 Spark 파티션당 PINOT SERVER 1개와 SEGMENT 1개 이상을 포함하는 pinot split을 만들어요. 마지막으로 각 파티션은 지정된 pinot server에서 병렬로 데이터를 읽어요.

각 Spark 파티션은 Pinot server와 연결을 열고 데이터를 읽어요. 예를 들어 지정된 쿼리에 대한 라우팅 테이블 정보가 다음과 같다고 가정해요.

- realtime ->
   - realtimeServer1 -> (segment1, segment2, segment3)
   - realtimeServer2 -> (segment4)
- offline ->
   - offlineServer10 -> (segment10, segment20)

segmentsPerSplit이 3이면 아래처럼 3개의 Spark 파티션이 생성돼요.

Spark Partition Queried Pinot Server/Segments
partition1 realtimeServer1 / segment1, segment2, segment3
partition2 realtimeServer2 / segment4
partition3 offlineServer10 / segment10, segment20

segmentsPerSplit이 1이면 6개의 Spark 파티션이 생성돼요.

Spark Partition Queried Pinot Server/Segments
partition1 realtimeServer1 / segment1
partition2 realtimeServer1 / segment2
partition3 realtimeServer1 / segment3
partition4 realtimeServer2 / segment4
partition5 offlineServer10 / segment10
partition6 offlineServer10 / segment20

segmentsPerSplit 값이 너무 낮으면 더 많은 병렬성이라는 뜻이에요. 하지만 이는 Pinot server와 많은 연결이 열리고 Pinot server의 QPS가 증가한다는 뜻이기도 해요.

segmentsPerSplit 값이 너무 높으면 병렬성이 낮다는 뜻이에요. 각 Pinot server는 요청당 더 많은 세그먼트를 스캔해요.

참고: 쿼리가 오면 Pinot server는 세그먼트 메타데이터를 기반으로 세그먼트를 가지치기(prune)해요. 어떤 경우(예: 일부 컬럼 기반 필터링)에는 일부 server가 데이터를 반환하지 않을 수 있어요. 따라서 일부 Spark 파티션은 비어 있을 수 있어요. 이 경우 데이터를 Spark에 로드한 후 효율적인 데이터 분석을 위해 repartition()을 적용할 수 있어요.

필터 및 컬럼 푸시다운

커넥터는 필터와 컬럼 푸시다운을 지원해요. 필터와 컬럼은 pinot server로 푸시돼요. 필터 및 컬럼 푸시다운은 Pinot와 Spark 사이의 데이터 전송을 최소화하므로 데이터를 읽는 동안 성능을 개선해요. 기본적으로 필터 푸시다운이 활성화돼요. 필터를 Spark에서 적용하고 싶다면 usePushDownFilters를 false로 설정해야 해요.

커넥터는 SQL을 사용하므로 모든 SQL 필터가 지원돼요.

세그먼트 가지치기 (Segment Pruning)

커넥터는 주어진 쿼리의 라우팅 테이블을 받아 어떤 Pinot server를 쿼리할지, 어떤 세그먼트를 스캔할지 정보를 얻어요. 주어진 Pinot 테이블에 파티셔닝이 활성화되고 Spark에서 생성된 쿼리가 특정 파티션을 스캔한다면 필요한 Pinot server 및 세그먼트 정보만 얻어져요(Pinot broker처럼 데이터 읽기 전에 세그먼트 가지치기 작업이 적용된다는 뜻). 자세한 내용: Optimizing routing

테이블 쿼리

커넥터는 SQL을 사용해 Pinot 테이블을 쿼리해요.

커넥터는 필터와 필요한 컬럼을 기반으로 realtime 및 offline 쿼리를 만들어요.

  • 쿼리된 테이블 유형이 OFFLINE 또는 REALTIME이면 특정 테이블 유형에 대한 라우팅 테이블 정보를 얻어요.
  • 쿼리된 테이블 유형이 HYBRID이면 realtime 및 offline 라우팅 테이블 정보를 모두 얻어요. 또한 커넥터는 주어진 테이블의 TimeBoundary 정보를 받아 realtime과 offline 세그먼트 데이터 간 겹침이 정확히 한 번 쿼리되도록 쿼리에 사용해요. 자세한 내용: Pinot Broker

쿼리 생성

주어진 사용에 대한 생성된 쿼리 예시(airlineStats 테이블이 hybrid이고 TimeBoundary 정보가 DaysSinceEpoch, 16084라고 가정):

val df = spark.read
      .format("pinot")
      .option("table", "airlineStats")
      .option("tableType", "hybrid")
      .load()

위 사용에 대해 realtime 및 offline SQL 쿼리가 생성돼요.

  • Offline 쿼리: select * from airlineStats_OFFLINE where DaysSinceEpoch < 16084 LIMIT {Int.MaxValue}
  • Realtime 쿼리: select * from airlineStats_REALTIME where DaysSinceEpoch >= 16084 LIMIT {Int.MaxValue}
val df = spark.read
      .format("pinot")
      .option("table", "airlineStats")
      .option("tableType", "offline")
      .load()
      .filter($"DestStateName" === "Florida")
      .filter($"Origin" === "ORD")
      .select($"DestStateName", $"Origin", $"Carrier")
  • Offline 쿼리: select DestStateName, Origin, Carrier from airlineStats_OFFLINE where DestStateName = 'Florida and Origin = 'ORD' LIMIT {Int.MaxValue}

참고: 모든 쿼리에 Limit이 추가돼요. 생성된 쿼리가 Pinot BrokerRequest 클래스로 변환되기 때문이에요. 이 작업에서 pinot는 limit을 자동으로 10으로 설정해요. 따라서 이 문제를 방지하기 위해 LIMIT을 Int.MaxValue로 설정했어요.

커넥터 읽기 파라미터

Configuration Description Required Default Value
table table type이 없는 Pinot 테이블 이름 Yes -
tableType Pinot 테이블 유형(realtime, offline 또는 hybrid) Yes -
controller Pinot controller host:port (useHttps/secureMode에서 스키마 추론) No localhost:9000
broker Pinot broker host:port (useHttps/secureMode에서 스키마 추론) No Pinot Controller에서 테이블의 broker 인스턴스 가져오기
usePushDownFilters 필터를 pinot server로 푸시할지 여부. true면 pinot server와 spark 사이의 데이터 교환이 최소화됨. No true
segmentsPerSplit 한 연결에서 pinot server가 스캔할 최대 세그먼트 수 No 3
pinotServerTimeoutMs pinot server에서 데이터를 얻는 최대 타임아웃(ms) No 10000 (10초)
useGrpcServer gRPC를 통한 읽기 활성화의 불리언 값. 스트리밍을 활용하므로 Pinot server와 Spark executor 양쪽에서 더 메모리 효율적. Pinot server에서 gRPC가 활성화되어 있어야 함. No false
queryOptions Pinot 쿼리 옵션의 쉼표 구분 목록(예: "enableNullHandling=true,skipUpsert=true") No ""
failOnInvalidSegments 응답 메타데이터가 잘못된 세그먼트를 나타내면 읽기 작업 실패 No false
secureMode HTTPS와 gRPC TLS를 활성화하는 통합 스위치 (명시적 useHttps/grpc.use-plain-text가 우선) No false
useHttps REST 호출용 HTTPS 활성화 (REST에 대해 secureMode 재정의) No false
grpc.use-plain-text gRPC에 plaintext 사용 (gRPC에 대해 secureMode 재정의). 설정하지 않으면 secureMode=true는 TLS(grpc.use-plain-text=false)를 의미. No true
grpc.port useGrpcServer=true일 때 사용할 Pinot server gRPC 포트. No 8090
grpc.max-inbound-message-size 스트리밍 응답을 받는 Spark executor의 최대 인바운드 gRPC 메시지 크기(바이트). No 134217728 (128 MB)
grpc.tls.keystore-type gRPC TLS keystore 유형. No JKS
grpc.tls.keystore-path gRPC TLS keystore 파일 경로. 상호 TLS 또는 클라이언트 인증서 인증으로 gRPC TLS를 사용할 때 필요. No None
grpc.tls.keystore-password gRPC 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 gRPC TLS truststore 비밀번호. No None
grpc.tls.ssl-provider gRPC TLS에 사용되는 SSL provider(예: JDK). No JDK
proxy.enabled controller 및 broker 발견을 위한 Pinot proxy 지원 활성화. No false
grpc.proxy-uri gRPC 요청용 proxy URI (proxy.enabled=true일 때만 사용). No None

더 알아보기 (Learn more)