Flink

Apache Flink를 이용해 Apache Pinot에 배치 수집을 수행하는 방법을 다루는 페이지예요.

출처: Flink

본문

Apache Pinot은 Apache Flink를 처리 프레임워크로 사용해 세그먼트를 생성하고 업로드하는 것을 지원해요. Pinot 배포판에는 Flink 애플리케이션(스트리밍 또는 배치)에 통합해 데이터를 바로 세그먼트로 만들어 Pinot 테이블에 쓸 수 있는 PinotSink이 포함되어 있어요.

PinotSink은 오프라인 테이블, 실시간 테이블, 그리고 업서트 테이블(전체 업서트만)을 지원해요. 데이터는 메모리에 버퍼링되고 설정된 임계값에 도달하면 세그먼트로 플러시되어 Pinot 클러스터에 업로드돼요.

요구 사항 (Requirements)

  • Flink 2.2.0 이상 – 새 Flink 2.x Sink API를 사용해요. 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) 값에 따라 업서트 시맨틱에 올바르게 참여해요.

요구 사항:

  1. 파티셔닝: 데이터는 업스트림 스트림(예: Kafka)과 같은 전략으로 파티셔닝되어야 해요
  2. 병렬도: Flink 잡 병렬도는 업스트림 스트림/테이블 파티션 수와 일치해야 해요
  3. 비교 컬럼: 비교 컬럼의 값 순서는 업스트림 스트림과 일관되어야 해요. 그래야 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 2.2.0 이상과 Java 11+을 요구해요. 옛 PinotSinkFunction (Flink 1.x SinkFunction API 기반)은 더 이상 사용되지 않으며 Flink 2.x에서는 동작하지 않아요.

// This no longer works with Flink 2.x
srcRows.addSink(new PinotSinkFunction<>(...));
// 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)

더 알아보기 (Learn more)