Flink
Flink
Apache Flink를 이용해 Apache Pinot에 배치 수집을 수행하는 방법을 다루는 페이지예요.
출처: Flink
본문
Apache Pinot은 Apache Flink를 처리 프레임워크로 사용해 세그먼트를 생성하고 업로드하는 것을 지원해요. Pinot 배포판에는 Flink 애플리케이션(스트리밍 또는 배치)에 통합해 데이터를 바로 세그먼트로 만들어 Pinot 테이블에 쓸 수 있는 PinotSink이 포함되어 있어요.
PinotSink은 오프라인 테이블, 실시간 테이블, 그리고 업서트 테이블(전체 업서트만)을 지원해요. 데이터는 메모리에 버퍼링되고 설정된 임계값에 도달하면 세그먼트로 플러시되어 Pinot 클러스터에 업로드돼요.
요구 사항 (Requirements)
- Flink 2.2.0 이상 – 새 Flink 2.x
SinkAPI를 사용해요. Java 21 지원이 포함돼요. - Java 11 이상 – Flink 2.x는 Java 11 이상을 요구해요.
Maven 의존성
Flink 잡에서 Pinot Flink 커넥터를 사용하려면 pom.xml에 다음 의존성을 추가하세요.:
<dependency>
<groupId>org.apache.pinot</groupId>
<artifactId>pinot-flink-connector</artifactId>
<version>${pinot.version}</version>
</dependency>
${pinot.version}을 사용 중인 Pinot 버전으로 바꾸세요. 최신 안정 버전은 Apache Pinot releases에서 확인하세요.
참고: 커넥터는 다음 의존성을 전이적으로 포함해요:
pinot-controller- 컨트롤러 클라이언트 API용pinot-segment-writer-file-based- 세그먼트 생성용flink-streaming-java- Flink 2.x 핵심 의존성
오프라인 테이블 수집 (Offline Table Ingestion)
퀵스타트 예시 (Quick Start Example)
// Set up Flink environment and data source
StreamExecutionEnvironment execEnv = StreamExecutionEnvironment.getExecutionEnvironment();
execEnv.setParallelism(2);
// Configure row type
RowTypeInfo typeInfo = new RowTypeInfo(
new TypeInformation[]{Types.FLOAT, Types.FLOAT, Types.STRING, Types.STRING},
new String[]{"lon", "lat", "address", "name"});
DataStream<Row> srcRows = execEnv.fromData(...);
// Create a PinotAdminClient to fetch Pinot schema and table config
String controllerUrl = "http://localhost:9000";
URI controllerUri = URI.create(controllerUrl);
String controllerAddress = controllerUri.getAuthority();
String controllerPath = controllerUri.getPath();
if (controllerPath != null && !controllerPath.isEmpty() && !"/".equals(controllerPath)) {
controllerAddress += controllerPath.endsWith("/") ? controllerPath.substring(0, controllerPath.length() - 1)
: controllerPath;
}
Properties properties = new Properties();
properties.setProperty(PinotAdminTransport.ADMIN_TRANSPORT_SCHEME, controllerUri.getScheme());
try (PinotAdminClient client = new PinotAdminClient(controllerAddress, properties)) {
// Fetch Pinot schema
Schema schema = client.getSchemaClient().getSchemaObject("starbucksStores");
// Fetch Pinot table config
TableConfig tableConfig =
client.getTableClient().getTableConfigObjectForType("starbucksStores", TableType.OFFLINE);
// Create Flink Pinot Sink (Flink 2.x API)
srcRows.sinkTo(new PinotSink<>(
new FlinkRowGenericRowConverter(typeInfo),
tableConfig,
schema,
controllerUrl));
}
execEnv.execute();
PinotSink에 controllerUrl을 전달하면 sink가 세그먼트 업로드에 필요한 push.controllerUri와 기본 outputDirURI 값을 주입할 수 있어요.
테이블 설정 (Table Configuration)
PinotSink은 TableConfig를 사용해 세그먼트 생성과 업로드의 배치 수집 설정을 결정해요. 테이블 설정 예시:
{
"tableName": "starbucksStores_OFFLINE",
"tableType": "OFFLINE",
"segmentsConfig": {
// ...
},
"tenants": {
// ...
},
"tableIndexConfig": {
// ...
},
"ingestionConfig": {
"batchIngestionConfig": {
"segmentIngestionType": "APPEND",
"segmentIngestionFrequency": "HOURLY",
"batchConfigMaps": [
{
"outputDirURI": "file:///tmp/pinotoutput",
"overwriteOutput": "false",
"push.controllerUri": "http://localhost:9000"
}
]
}
}
}
필수 설정:
outputDirURI- 업로드 전에 세그먼트가 쓰여지는 디렉터리push.controllerUri- 세그먼트 업로드를 위한 Pinot 컨트롤러 URL
완전한 실행 가능한 예제는 FlinkQuickStart.java를 참고하세요.
실시간 테이블 수집 (Realtime Table Ingestion)
Non-업서트 실시간 테이블
업서트가 없는 표준 실시간 테이블은 오프라인 테이블과 같은 방식을 사용하되 테이블 타입으로 REALTIME을 지정해요:
// Reuse the PinotAdminClient setup from the offline example above...
// Fetch table config for realtime table
Schema schema = client.getSchemaClient().getSchemaObject("myTable");
TableConfig tableConfig =
client.getTableClient().getTableConfigObjectForType("myTable", TableType.REALTIME);
// Same sink configuration
srcRows.sinkTo(new PinotSink<>(
new FlinkRowGenericRowConverter(typeInfo),
tableConfig,
schema,
controllerUrl));
execEnv.execute();
업서트 테이블 (Upsert Tables)
전체 업서트 테이블 (Full Upsert Tables)
Flink 커넥터는 각 레코드가 모든 컬럼을 포함하는 전체 업서트 테이블의 백필을 지원해요. 업로드된 세그먼트는 비교 컬럼(comparison column) 값에 따라 업서트 시맨틱에 올바르게 참여해요.
요구 사항:
- 파티셔닝: 데이터는 업스트림 스트림(예: Kafka)과 같은 전략으로 파티셔닝되어야 해요
- 병렬도: Flink 잡 병렬도는 업스트림 스트림/테이블 파티션 수와 일치해야 해요
- 비교 컬럼: 비교 컬럼의 값 순서는 업스트림 스트림과 일관되어야 해요. 그래야 Pinot이 주어진 키에 대해 어느 레코드가 최신인지 올바르게 해석할 수 있어요. 중요한 고려 사항은 Pinot upsert comparison column docs 참고.
예시:
// Set up Flink environment
StreamExecutionEnvironment execEnv = StreamExecutionEnvironment.getExecutionEnvironment();
execEnv.setParallelism(2); // MUST match number of partitions in stream/table
// Configure row type matching your upsert table schema
RowTypeInfo typeInfo = new RowTypeInfo(
new TypeInformation[]{Types.INT, Types.STRING, Types.STRING, Types.FLOAT, Types.LONG, Types.BOOLEAN},
new String[]{"playerId", "name", "game", "score", "timestampInEpoch", "deleted"});
DataStream<Row> srcRows = execEnv.fromData(...);
String controllerUrl = "http://localhost:9000";
URI controllerUri = URI.create(controllerUrl);
String controllerAddress = controllerUri.getAuthority();
String controllerPath = controllerUri.getPath();
if (controllerPath != null && !controllerPath.isEmpty() && !"/".equals(controllerPath)) {
controllerAddress += controllerPath.endsWith("/") ? controllerPath.substring(0, controllerPath.length() - 1)
: controllerPath;
}
Properties properties = new Properties();
properties.setProperty(PinotAdminTransport.ADMIN_TRANSPORT_SCHEME, controllerUri.getScheme());
try (PinotAdminClient client = new PinotAdminClient(controllerAddress, properties)) {
Schema schema = client.getSchemaClient().getSchemaObject("myUpsertTable");
TableConfig tableConfig =
client.getTableClient().getTableConfigObjectForType("myUpsertTable", TableType.REALTIME);
// IMPORTANT: Partition data by primary key using the SAME logic as the stream
srcRows.partitionCustom(
(Partitioner<Integer>) (key, partitions) -> key % partitions,
r -> (Integer) r.getField("playerId")) // Primary key field
.sinkTo(new PinotSink<>(
new FlinkRowGenericRowConverter(typeInfo),
tableConfig,
schema,
controllerUrl));
}
execEnv.execute();
파티셔닝 동작 방식:
업서트 테이블에 세그먼트를 업로드할 때 Pinot은 파티션 ID를 인코딩하는 특수 세그먼트 명명 규약 UploadedRealtimeSegmentName을 사용해요. 형식은:
{prefix}__{tableName}__{partitionId}__{uploadTimeMs}__{sequenceId}
예시: flink__myTable__0__1724045187__1
각 Flink 서브태스크는 서브태스크 인덱스를 기반으로 특정 파티션의 세그먼트를 생성해요. 그 세그먼트들은 스트림으로 소비된 세그먼트의 그 파티션을 처리하는 동일한 서버 인스턴스에 배정되어, 모든 세그먼트에서 올바른 업서트 동작을 보장해요.
설정 옵션:
추가 생성자 파라미터로 세그먼트 생성을 커스터마이징할 수 있어요:
new PinotSink<>(
recordConverter,
tableConfig,
schema,
segmentFlushMaxNumRecords, // Default: 500,000, number of rows per segment
executorPoolSize, // Default: 5, number of threads to use to upload segment
segmentNamePrefix, // Default: "flink"
segmentUploadTimeMs // Default: current time, upload time value to encode in segment name
)
부분 업서트 테이블 (Partial Upsert Tables)
경고(WARNING): 부분 업서트 테이블에는 Flink 기반 업로드를 권장하지 않아요.
부분 업서트 테이블에서 업로드된 세그먼트는 컬럼의 일부분 또는 주 키의 중간 로우만 포함해요. 업로드된 로우가 최종 상태가 아니고 이후 업데이트가 스트림으로 도착하면, 부분 업서트 머저(merger)가 복제본 간 일관성 없는 결과를 만들어낼 수 있어요. 이로 인해 탐지·해결이 어려운 데이터 불일치가 발생할 수 있어요.
부분 업서트 테이블은 스트림 기반 수집만 사용하거나, 업로드된 데이터가 각 주 키의 최종 상태를 나타내도록 해야 해요.
고급 설정 (Advanced Configuration)
세그먼트 플러시 제어
세그먼트가 플러시·업로드되는 시점을 제어:
// Same setup as previous examples...
long segmentFlushMaxNumRecords = 1000000; // Flush after 1M records
int executorPoolSize = 10; // Thread pool size for async uploads
srcRows.sinkTo(new PinotSink<>(
new FlinkRowGenericRowConverter(typeInfo),
tableConfig,
schema,
segmentFlushMaxNumRecords,
executorPoolSize
));
세그먼트 명명
더 나은 정리를 위해 세그먼트 명명과 업로드 시간을 커스터마이징:
// Same setup as previous examples...
String segmentNamePrefix = "flink_job_daily";
Long segmentUploadTimeMs = 1724045185000L; // Group segments by upload run time
srcRows.sinkTo(new PinotSink<>(
new FlinkRowGenericRowConverter(typeInfo),
tableConfig,
schema,
DEFAULT_SEGMENT_FLUSH_MAX_NUM_RECORDS,
DEFAULT_EXECUTOR_POOL_SIZE,
segmentNamePrefix,
segmentUploadTimeMs
));
Flink 1.x에서의 마이그레이션
중요: 커넥터는 이제 Flink 2.2.0 이상과 Java 11+을 요구해요. 옛
PinotSinkFunction(Flink 1.xSinkFunctionAPI 기반)은 더 이상 사용되지 않으며 Flink 2.x에서는 동작하지 않아요.
옛 API (Flink 1.x – Deprecated)
// This no longer works with Flink 2.x
srcRows.addSink(new PinotSinkFunction<>(...));
새 API (Flink 2.x – 필수)
// Use this for Flink 2.2.0 and later
srcRows.sinkTo(new PinotSink<>(...));
주요 변경 사항:
.addSink()를.sinkTo()로 교체PinotSinkFunction을PinotSink으로 교체- Flink 버전을 2.2.0 이상으로 업데이트
- Java 21 지원됨
추가 자료 (Additional Resources)
- 설계 제안 (Design Proposal) - 원래 설계 동기
- PR #13107 - 업서트 테이블용 외부 파티셔닝 세그먼트
- PR #13837 - 업서트 백필용 Flink 커넥터 개선
- PR #18250 - Flink 2.2.0 업그레이드
- 테이블 설정 레퍼런스 (Table Configuration Reference)
- 스키마 설정 레퍼런스 (Schema Configuration Reference)