Java
Java
Pinot는 broker 라우팅 SQL 쿼리를 위한 네이티브 Java 쿼리 클라이언트를 제공해요. 클라이언트는 tenant 인지이며 blocking 및 async 실행을 지원하고 ZooKeeper 또는 controller를 통해 broker를 발견할 수 있어요.
출처: 문서
본문
테이블, 스키마, 세그먼트, tenant, 인스턴스, 작업 관리 같은 controller REST 작업에는 Java admin client를 사용하세요.
설치
클라이언트는 다음 의존성을 포함해 사용할 수 있어요.
<dependency>
<groupId>org.apache.pinot</groupId>
<artifactId>pinot-java-client</artifactId>
<version>1.4.0</version>
</dependency>
include 'org.apache.pinot:pinot-java-client:1.4.0'
Java 클라이언트 코드를 로컬에서 빌드해 사용할 수도 있어요.
기본 인증은 Pinot 0.10.0부터 JDBC와 Java 클라이언트에서 지원돼요. 그 릴리스 이상의 클라이언트 JAR을 사용하고 있는지 확인하세요.
사용 방법
pinot-java-client로 Pinot를 쿼리하는 예시예요.
import org.apache.pinot.client.Connection;
import org.apache.pinot.client.ConnectionFactory;
import org.apache.pinot.client.ResultSetGroup;
import org.apache.pinot.client.ResultSet;
/**
* pinot-java-client를 사용해 Java에서 Pinot를 쿼리하는 방법을 보여줌.
*/
public class PinotClientExample {
public static void main(String[] args) {
// Pinot 연결
String zkUrl = "localhost:2181";
String pinotClusterName = "PinotCluster";
Connection pinotConnection = ConnectionFactory.fromZookeeper(zkUrl + "/" + pinotClusterName);
String query = "SELECT COUNT(*) FROM myTable GROUP BY foo";
ResultSetGroup pinotResultSetGroup = pinotConnection.execute(query);
ResultSet resultTableResultSet = pinotResultSetGroup.getResultSet(0);
int numRows = resultTableResultSet.getRowCount();
int numColumns = resultTableResultSet.getColumnCount();
String columnValue = resultTableResultSet.getString(0, 1);
String columnName = resultTableResultSet.getColumnName(1);
System.out.println("ColumnName: " + columnName + ", ColumnValue: " + columnValue);
}
}
ConnectionFactory
클라이언트는 Pinot 클러스터에 연결을 생성하는 ConnectionFactory 클래스를 제공해요. 현재 소스는 다음 쿼리 연결 패턴을 지원해요.
- ZooKeeper (권장):
ConnectionFactory.fromZookeeper(...)가 Helix external view에서 broker를 동적으로 해석하고 테이블별로 쿼리를 라우팅. - Broker list:
ConnectionFactory.fromHostList(...)가 고정 broker 목록을 사용. 단순 배포, 로드 밸런싱된 broker, 로컬 테스트에 주로 유용. - Controller address:
ConnectionFactory.fromController(...)가 ZooKeeper를 직접 관찰하는 대신 controller에서 broker 매핑을 주기적으로 새로 고침. - Properties object:
ConnectionFactory.fromProperties(Properties)가brokerList속성을 읽고 고정 broker 목록 연결을 구축.
fromController(...)와 fromControllerGrpc(...)의 경우 controller를 host:port로 전달하세요. 클라이언트는 scheme 속성에서 http 또는 https를 적용해요.
Pinot 클러스터가 Kubernetes 안에서 실행 중이고 Kubernetes 밖에서 연결하려 한다면, ZooKeeper 방식은 해석할 수 없는 내부 호스트 이름을 반환할 가능성이 높아요.
따라서 Kubernetes 배포에서는 broker 앞에 있는 로드 밸런서의 호스트 이름을 전달하는 것이 좋아요.
예시:
Properties properties = new Properties();
properties.setProperty("brokerList", "broker-1:8099,broker-2:8099");
Connection zkConnection =
ConnectionFactory.fromZookeeper("some-zookeeper-server:2181/zookeeperPath/pinot-cluster");
Connection controllerConnection =
ConnectionFactory.fromController("controller.example.com:9000");
Connection brokerConnection =
ConnectionFactory.fromHostList("broker-1:8099", "broker-2:8099");
Connection propertiesConnection = ConnectionFactory.fromProperties(properties);
gRPC 연결
Java 클라이언트는 ConnectionFactory.fromControllerGrpc(...), fromZookeeperGrpc(...), fromHostListGrpc(...)로 gRPC broker 연결도 노출해요.
Properties properties = new Properties();
properties.setProperty("usePlainText", "false");
properties.setProperty("maxInboundMessageSizeBytes", "134217728");
properties.setProperty("channelKeepAliveTimeSeconds", "60");
properties.setProperty("channelKeepAliveTimeoutSeconds", "20");
properties.setProperty("channelKeepAliveWithoutCalls", "true");
properties.setProperty("channelShutdownTimeoutSeconds", "10");
properties.setProperty("tls.truststore.path", "/path/to/grpc-truststore.jks");
properties.setProperty("tls.truststore.password", "changeit");
GrpcConnection connection = ConnectionFactory.fromControllerGrpc(properties, "localhost:9000");
ResultSetGroup resultSetGroup = connection.execute(
"SELECT COUNT(*) FROM myTable",
Map.of("blockRowSize", "10000", "encoding", "JSON", "compression", "ZSTD"));
아래 gRPC 전송 속성은 GrpcConfig를 통해 연결 Properties에서 직접 읽혀요. TLS 설정은 pinot.*.tls.*가 아니라 tls.* 네임스페이스를 사용해요.
| Property | Default | Notes |
|---|---|---|
usePlainText |
true |
일반 텍스트 gRPC 전송 사용. TLS를 활성화하려면 false로 설정. |
maxInboundMessageSizeBytes |
134217728 (128 MB) |
클라이언트가 수락하는 최대 인바운드 gRPC 메시지 크기. |
channelKeepAliveTimeSeconds |
-1 (disabled) |
Keepalive ping 간격. keepalive를 활성화하려면 양수 값 설정. |
channelKeepAliveTimeoutSeconds |
20 |
keepalive 승인을 기다리는 타임아웃. |
channelKeepAliveWithoutCalls |
true |
활성 RPC가 없어도 keepalive ping 허용. |
channelShutdownTimeoutSeconds |
10 |
종료 시 클라이언트가 gRPC 채널 종료를 기다리는 시간. |
tls.keystore.type |
JVM 기본 keystore 유형 (KeyStore.getDefaultType()) |
상호 TLS용 클라이언트 keystore 유형. |
tls.keystore.path |
None | 상호 TLS용 클라이언트 keystore 경로. |
tls.keystore.password |
None | 클라이언트 keystore 비밀번호. |
tls.truststore.type |
JVM 기본 keystore 유형 (KeyStore.getDefaultType()) |
broker 인증서를 검증하는 데 사용되는 truststore 유형. |
tls.truststore.path |
None | broker 인증서를 검증하는 데 사용되는 truststore 경로. |
tls.truststore.password |
None | Truststore 비밀번호. |
tls.ssl.provider |
JDK |
gRPC 클라이언트 SSL 컨텍스트를 구축할 때 사용되는 SSL 공급자. |
tls.insecure |
false |
broker 인증서 검증 건너뛰기. 비프로덕션 테스트에만 적합. |
tls.protocols |
JVM TLS 기본값 | TLSv1.2,TLSv1.3 같은 쉼표 구분 TLS 프로토콜 허용 목록. |
인증 헤더 같은 기본 메타데이터에는 headers.<name> 속성 접두사를 사용하세요. blockRowSize, compression, encoding 같은 쿼리 특정 gRPC 메타데이터는 execute(..., metadataMap) 또는 executeGrpc(..., metadataMap)의 메타데이터 맵에 전달해야 해요.
쿼리 메서드
쿼리를 blocking 및 async 방식으로 실행할 수 있어요.
Connection.execute(String)— blocking 쿼리Connection.executeAsync(String)— future 객체를 반환하는 비동기 쿼리
ResultSetGroup resultSetGroup =
connection.execute("select * from foo...");
// OR
Future<ResultSetGroup> futureResultSetGroup =
connection.executeAsync("select * from foo...");
PreparedStatement로 쿼리 파라미터를 이스케이프할 수도 있어요. Prepared Statement를 데이터베이스에 저장하지 않으므로 이후 쿼리 성능은 향상되지 않아요.
PreparedStatement statement =
connection.prepareStatement("select * from foo where a = ?");
statement.setString(1, "bar");
ResultSetGroup resultSetGroup = statement.execute();
// OR
Future<ResultSetGroup> futureResultSetGroup = statement.executeAsync();
Connection.execute(...)에는 명시적 테이블 이름 또는 테이블 이름 iterable을 받는 오버로드도 있어요. 그 오버로드는 SQL을 다시 파싱하지 않고 클라이언트가 broker를 선택할 수 있게 해요.
커서 페이지네이션
Connection 뒤의 HTTP 전송은 커서 페이지네이션을 구현해요. 쿼리 결과 집합이 클 수 있고 모든 것을 하나의 응답으로 로드하는 대신 페이지 단위 탐색을 원할 때 openCursor(query, pageSize)를 사용하세요.
try (ResultCursor cursor = connection.openCursor(
"SELECT playerName, yearID FROM baseballStats ORDER BY yearID", 1000)) {
CursorResultSetGroup firstPage = cursor.getCurrentPage();
while (cursor.hasNext()) {
CursorResultSetGroup nextPage = cursor.next();
// nextPage.getResultSet(...) 처리
}
}
커서 API는 다음을 지원해요.
getCurrentPage()— 현재 로드된 페이지 조회next()/nextAsync()및previous()/previousAsync()— 탐색seekToPage()/seekToPageAsync()— 직접 페이지 이동getCursorId(),getCurrentPageNumber(),getTotalRows(),isExpired()— 커서 메타데이터close()— 서버 측 커서를 삭제하고 리소스 해제
커서 페이지네이션은 기본 HTTP 전송이 구현하는 CursorCapable을 하부 전송이 구현할 때만 사용할 수 있어요.
결과 집합
getResultSet(int)로 얻은 첫 번째 ResultSet에 다양한 getter 메서드로 결과를 얻을 수 있어요.
String query = "select foo, bar from baz where quux = 'quuux'";
ResultSetGroup resultSetGroup = connection.execute(query);
ResultSet resultSet = resultSetGroup.getResultSet(0);
for (int i = 0; i < resultSet.getRowCount(); ++i) {
System.out.println("foo: " + resultSet.getString(i, 0));
System.out.println("bar: " + resultSet.getInt(i, 1));
}
인증
Pinot는 기본 HTTP 인증을 지원하며 구성으로 클러스터에 활성화할 수 있어요. 클라이언트 측 Java 애플리케이션에서 기본 HTTP 인증을 지원하려면 Pinot Java Client 0.10.0 이상을 사용하고 있는지 확인하세요. 다음 코드 스니펫은 기본 HTTP 인증이 활성화된 Pinot 클러스터에 Java 클라이언트를 사용해 연결하고 쿼리하는 방법을 보여줘요.
final String username = "admin";
final String password = "verysecret";
// 사용자 이름과 비밀번호를 연결하고 base64로 인코딩
String plainCredentials = username + ":" + password;
String base64Credentials = new String(
Base64.getEncoder().encode(plainCredentials.getBytes()));
String authorizationHeader = "Basic " + base64Credentials;
Map<String, String> headers = new HashMap<>();
headers.put("Authorization", authorizationHeader);
JsonAsyncHttpPinotClientTransportFactory factory =
new JsonAsyncHttpPinotClientTransportFactory();
factory.setHeaders(headers);
PinotClientTransport clientTransport = factory
.buildTransport();
Properties properties = new Properties();
properties.put("brokerList", "localhost:8000,localhost:8001");
Connection connection = ConnectionFactory.fromProperties(properties, clientTransport);
String query = "select count(*) FROM baseballStats limit 1";
ResultSetGroup rs = connection.execute(query);
System.out.println(rs);
connection.close();
연결 속성
Java 쿼리 클라이언트는 다음 연결 속성을 Properties에서 직접 읽어요.
| Property | Default | Used by | Notes |
|---|---|---|---|
brokerConnectTimeoutMs |
2000 |
HTTP broker 전송 | broker 연결 타임아웃(밀리초) |
brokerReadTimeoutMs |
60000 |
HTTP broker 전송 및 커서 fetch | broker 읽기 타임아웃(밀리초) |
brokerHandshakeTimeoutMs |
2000 |
HTTP broker 전송 | TLS 핸드셰이크 타임아웃(밀리초) |
controllerConnectTimeoutMs |
2000 |
Controller 기반 broker 캐시 | fromController(...)의 controller 연결 타임아웃 |
controllerReadTimeoutMs |
60000 |
Controller 기반 broker 캐시 | broker-map 새로 고침의 controller 읽기 타임아웃 |
controllerHandshakeTimeoutMs |
2000 |
Controller 기반 broker 캐시 | Controller TLS 핸드셰이크 타임아웃 |
headers.<name> |
None | HTTP broker 전송 및 controller broker 캐시 | headers.Authorization 같은 기본 HTTP 헤더 추가 |
scheme |
http |
HTTP broker 전송 및 controller broker 캐시 | TLS 활성화 broker 및 controller용으로 https로 설정 |
queryOptions |
Empty string | HTTP broker 전송 | JSON 요청 본문에 Pinot 쿼리 옵션으로 주입 |
useMultistageEngine |
false |
HTTP broker 전송 | HTTP 요청을 /query/sql에서 /query로 전환 |
appId |
None | HTTP broker 전송 및 controller broker 캐시 | 생성된 user agent 문자열 접두사 |
failOnExceptions |
true |
Connection |
broker 응답에 쿼리 처리 예외가 포함되면 PinotClientException throw |
preferTLS |
false |
ZooKeeper 기반 broker 발견 | Helix에서 broker 발견 시 broker TLS 포트 선호 |
useGrpcPort |
false |
ZooKeeper/controller broker 발견 | broker 목록 구축 시 broker gRPC 포트 선호 |
pinot.java_client.tls.* |
None | HTTP broker 전송 및 controller broker 캐시 | TlsUtils.extractTlsConfig(...)가 소비하는 TLS 구성 네임스페이스 |
brokerTlsV10Enabled |
false |
HTTP broker 전송 | broker 요청에 TLSv1.0 재활성화 |
controllerTlsV10Enabled |
false |
Controller broker 캐시 | controller 요청에 TLSv1.0 재활성화 |
예시:
Properties properties = new Properties();
properties.setProperty("scheme", "https");
properties.setProperty("headers.Authorization", "Basic <base64-credentials>");
properties.setProperty("queryOptions", "timeoutMs=10000");
properties.setProperty("brokerReadTimeoutMs", "10000");
properties.setProperty("controllerReadTimeoutMs", "10000");
properties.setProperty("pinot.java_client.tls.truststore.path", "/path/to/truststore.jks");
properties.setProperty("pinot.java_client.tls.truststore.password", "changeit");
Connection connection = ConnectionFactory.fromController(properties, "controller.example.com:9000");
JsonAsyncHttpPinotClientTransportFactory를 직접 사용자 지정할 때 scheme 속성은 Properties에 scheme을 설정할 때만 팩토리를 재정의해요. scheme이 없으면 withConnectionProperties(...)는 팩토리에 이미 설정된 scheme을 유지하고, 둘 다 지정하지 않을 때만 http로 폴백해요.
요청 추적
Java 클라이언트(JsonAsyncHttpPinotClientTransport)는 모든 HTTP 쿼리에 X-Correlation-Id 헤더(요청별 고유 UUID)를 자동으로 첨부해요. 이 ID는:
- 클라이언트가
org.apache.pinot.client아래DEBUG레벨로 기록 - broker 액세스 로그에 나타나 프록시 및 로드 밸런서 전반의 종단 간 추적을 가능하게 함
구성이 필요 없어요. 클라이언트 로그에서 상관 ID를 보려면 로그 레벨을 설정하세요.
log4j.logger.org.apache.pinot.client=DEBUG