Flink 커넥터

Apache Pinot 테이블에 데이터를 직접 쓰는 Apache Flink 커넥터로, offline, realtime, upsert 테이블 유형을 지원해요.

Pinot Flink Connector는 모든 Flink 스트리밍 또는 배치 작업에 연결되어 Pinot 세그먼트를 프로세스 내(in-process)에서 생성하고 클러스터에 직접 업로드하는 PinotSink를 제공해요. offline 테이블, realtime 테이블, full-upsert 테이블을 지원해요.

출처: 문서

본문

요구 사항

  • Flink 2.2.0 이상 – 커넥터는 Flink 2.x를 필요로 하고 Java 21을 지원해요.
  • Java 11+ – 이전 버전: Java 8 지원은 Flink 1.x에서 끝났어요.

참고: Flink 1.x를 사용한다면 Pinot 버전 1.5.0 이하를 사용하세요. Pinot 1.6.0+는 Flink 2.x가 필요해요.

사용 사례

  • Offline 테이블 백필(backfill) -- 데이터 레이크, 데이터베이스 내보내기, 또는 Flink가 읽을 수 있는 어떤 소스에서든 offline 테이블을 채우거나 갱신.
  • Upsert 테이블 부트스트래핑 -- 올바른 파티션 할당과 비교 컬럼 정렬을 유지하면서 히스토리 데이터로 realtime upsert 테이블을 시드.
  • ETL 및 enrichment 파이프라인 -- 데이터를 로드하기 전에 조인, 필터링, enrichment하는 더 큰 Flink DAG 안에 Pinot 쓰기를 임베드.
Capability Flink Connector Spark Batch Ingestion Standalone LaunchDataIngestionJob
Processing framework Apache Flink (streaming or batch) Apache Spark None (독립 실행형 Java 프로세스)
Upsert table backfill Yes -- 올바르게 파티셔닝된 uploaded-realtime 세그먼트 생성 기본 지원 안 함 기본 지원 안 함
Custom transformation logic 전체 Flink API (joins, windows, aggregations) 전체 Spark API 수집 구성 변환으로 제한됨
Cluster dependency Flink 클러스터 또는 로컬 Flink 환경 필요 Spark 클러스터 필요 단일 JVM 프로세스로 실행
Typical data sources Kafka, data lake 파일, JDBC, 모든 Flink 소스 HDFS, S3, GCS, 모든 Spark 소스 로컬/원격 파일 (CSV, JSON, Avro, Parquet, ORC, Thrift)
Best for 이미 Flink를 운영하는 팀; upsert 백필 시나리오 이미 Spark를 운영하는 팀; 대규모 배치 로드 처리 프레임워크 없는 간단한 일회성 또는 스케줄 로드

Maven 의존성

<dependency>
  <groupId>org.apache.pinot</groupId>
  <artifactId>pinot-flink-connector</artifactId>
  <version>${pinot.version}</version>
</dependency>

${pinot.version}을 Pinot 릴리스 버전으로 교체하세요. 최신 안정 버전은 Apache Pinot releases에서 확인하세요.

이 아티팩트는 Apache Maven 저장소에 게시되며 Pinot admin client, segment writer, Flink 2.x 핵심 의존성을 전이적으로 포함해요.

빠른 예시

StreamExecutionEnvironment execEnv = StreamExecutionEnvironment.getExecutionEnvironment();
execEnv.setParallelism(2);

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 adminClient = new PinotAdminClient(controllerAddress, properties)) {
  Schema schema = adminClient.getSchemaClient().getSchemaObject("myTable");
  TableConfig tableConfig =
      adminClient.getTableClient().getTableConfigObjectForType("myTable", TableType.OFFLINE);

  srcRows.sinkTo(new PinotSink<>(
      new FlinkRowGenericRowConverter(typeInfo),
      tableConfig,
      schema,
      controllerUrl));
}
execEnv.execute();

전체 구성 참조

업셋 파티셔닝 요구 사항, 세그먼트 플러시 제어, 세그먼트 이름 지정, realtime 테이블 지원을 포함한 완전한 구성 세부 사항은 Flink 배치 수집 참조를 참조하세요.

Deprecated: 레거시 PinotSinkFunction(Flink 1.x SinkFunction API 기반)은 deprecated되며 Flink 2.x에서 더 이상 동작하지 않아요. PinotSink와 sinkTo() API를 사용하도록 코드를 업데이트하세요.

구 API (Flink 1.x – 더 이상 지원 안 함)

// 이 코드는 Flink 2.x에서 동작하지 않음
srcRows.addSink(new PinotSinkFunction<>(...));

새 API (Flink 2.x)

// Flink 2.2.0 이상에서 이 API를 사용
srcRows.sinkTo(new PinotSink<>(...));

마이그레이션에 대한 도움이 필요하면 업데이트된 Flink 배치 수집 예시를 참조하세요.

추가 리소스

더 알아보기 (Learn more)