미니언 태스크 플러그인

미니언 태스크 플러그인 (Minion Task Plugin)

태스크 생성기(generator)와 태스크 실행기(executor)를 가진 커스텀 미니언 태스크 플러그인을 작성하는 방법을 다루는 문서예요.

출처: 문서

본문

미니언 태스크는 Pinot에서 백그라운드 유지보수·데이터 처리 잡을 실행해요. 각 태스크 타입은 두 가지 컴포넌트를 가져요: 태스크 구성을 만드는 태스크 생성기 (Task Generator) (컨트롤러에서 실행)와 작업을 수행하는 태스크 실행기 (Task Executor) (미니언에서 실행).

아키텍처 (Architecture)

Controller                          Minion
┌────────────────────┐             ┌────────────────────┐
│  PinotTaskGenerator│  ─ Helix ─> │ PinotTaskExecutor  │
│  - generateTasks() │  schedules  │ - executeTask()    │
│  - getTaskType()   │             │ - cancel()         │
└────────────────────┘             └────────────────────┘
  1. 컨트롤러는 스케줄 또는 ad-hoc API 호출로 태스크 생성기를 실행해요
  2. 생성기는 작업을 설명하는 PinotTaskConfig 객체를 만들어요
  3. 컨트롤러는 Helix를 통해 태스크를 스케줄해요
  4. 미니언 인스턴스는 태그/라우팅을 기준으로 태스크를 집어요
  5. 미니언은 태스크 실행기를 사용해 태스크를 실행해요
  6. 결과는 컨트롤러로 반환돼요

1단계: 태스크 실행기 구현 (Step 1: Implement the Task Executor)

실행기는 미니언에서 실행되고 실제 작업을 수행해요.

태스크 실행기 팩토리 (Task Executor Factory)

import org.apache.pinot.minion.executor.PinotTaskExecutorFactory;
import org.apache.pinot.minion.executor.PinotTaskExecutor;
import org.apache.pinot.spi.annotations.minion.TaskExecutorFactory;

@TaskExecutorFactory(enabled = true)
public class MyTaskExecutorFactory implements PinotTaskExecutorFactory {

  @Override
  public void init(MinionTaskZkMetadataManager zkMetadataManager, MinionConf minionConf) {
    // Initialize resources needed by the executor
  }

  @Override
  public String getTaskType() {
    return "MyCustomTask";  // Must match the generator's task type
  }

  @Override
  public PinotTaskExecutor create() {
    return new MyTaskExecutor();
  }
}

태스크 실행기 (Task Executor)

import org.apache.pinot.minion.executor.PinotTaskExecutor;
import org.apache.pinot.core.minion.PinotTaskConfig;

public class MyTaskExecutor implements PinotTaskExecutor {

  private volatile boolean _cancelled = false;

  @Override
  public Object executeTask(PinotTaskConfig pinotTaskConfig) throws Exception {
    String tableName = pinotTaskConfig.getTableName();
    Map<String, String> config = pinotTaskConfig.getConfig();

    // Perform the task work here
    // Check _cancelled periodically for long-running tasks

    return "Task completed successfully";
  }

  @Override
  public void cancel() {
    _cancelled = true;
  }
}

실행기 팩토리 클래스는 자동 등록을 위해 @TaskExecutorFactory(enabled = true)로 어노테이션해야 해요. 클래스는 org.apache.pinot.*.plugin.minion.tasks.*와 일치하는 패키지에 있어야 해요.

메타데이터 푸시로 세그먼트 안전하게 교체 (Safely replace a segment with metadata push)

커스텀 태스크가 BaseSingleSegmentConversionExecutor를 확장하고 같은 세그먼트 이름으로 세그먼트를 재빌드하며 push.mode=METADATA를 사용한다면, 컨트롤러 측 복사에 옵트인해요:

@Override
protected boolean isCopyToDeepStoreForMetadataPush() {
  return true;
}

또한 각 태스크 구성에 output.segment.dir.uri를 제공해요. 미니언은 이 위치를 사용해 변환된 세그먼트를 스테이징하고, 컨트롤러는 스테이징 위치와 테이블의 딥 스토어에 모두 접근할 수 있어야 해요.

이 옵트인 경로는 기존 세그먼트 다운로드 URL을 보존해요. 미니언은 태스크별 이름으로 교체를 스테이징하고, 원본 세그먼트 CRC와 새로고침 전용 가드를 컨트롤러에 보내며, 스테이징 데이터를 세그먼트의 기존 딥 스토어 위치로 복사하도록 컨트롤러에 요청해요. 요청이 완료되거나 실패한 후 미니언은 스테이징 파일을 삭제해요. 같은 태스크의 재시도는 자신의 남은 스테이징 파일을 덮어쓸 수 있어요.

변환 결과가 원본 세그먼트와 같은 CRC를 생성하면 Pinot는 세그먼트 데이터를 스테이징·복사하지 않아요. 기존 다운로드 URL에 대해 세그먼트 메타데이터만 새로고침해요.

이 메서드는 같은 이름 교체 태스크에만 오버라이드해요. 기본값은 false이므로 PurgeTask, RefreshSegmentTask, UpsertCompactionTask를 포함한 기존 태스크는 실행기가 명시적으로 옵트인하지 않는 한 현재 메타데이터 푸시 동작을 유지해요.

2단계: 태스크 생성기 구현 (Step 2: Implement the Task Generator)

생성기는 컨트롤러에서 실행되고 언제·어떻게 태스크를 만들지 결정해요.

import org.apache.pinot.controller.helix.core.minion.generator.PinotTaskGenerator;
import org.apache.pinot.spi.annotations.minion.TaskGenerator;
import org.apache.pinot.spi.config.table.TableConfig;
import org.apache.pinot.core.minion.PinotTaskConfig;

@TaskGenerator(enabled = true)
public class MyTaskGenerator implements PinotTaskGenerator {

  private ClusterInfoAccessor _clusterInfoAccessor;

  @Override
  public void init(ClusterInfoAccessor clusterInfoAccessor) {
    _clusterInfoAccessor = clusterInfoAccessor;
  }

  @Override
  public String getTaskType() {
    return "MyCustomTask";  // Must match the executor's task type
  }

  // Called on schedule for all tables with this task type enabled
  @Override
  public List<PinotTaskConfig> generateTasks(List<TableConfig> tableConfigs) {
    List<PinotTaskConfig> tasks = new ArrayList<>();
    for (TableConfig tableConfig : tableConfigs) {
      Map<String, String> taskConfig = new HashMap<>();
      taskConfig.put("myParam", "myValue");

      tasks.add(new PinotTaskConfig(
          getTaskType(),
          tableConfig.getTableName(),
          taskConfig));
    }
    return tasks;
  }

  // Called for ad-hoc task execution via API
  @Override
  public List<PinotTaskConfig> generateTasks(
      TableConfig tableConfig, Map<String, String> taskConfigs) throws Exception {
    return List.of(new PinotTaskConfig(
        getTaskType(),
        tableConfig.getTableName(),
        taskConfigs));
  }
}

선택적 생성기 오버라이드 (Optional Generator Overrides)

메서드 기본값 설명
getTaskTimeoutMs(String minionTag) 3,600,000 (1시간) 태스크 타임아웃(밀리초)
getNumConcurrentTasksPerInstance(String minionTag) 1 미니언당 최대 동시 태스크
getMaxAttemptsPerTask(String minionTag) 1 (재시도 없음) 최대 재시도 횟수
getMinionInstanceTag(TableConfig tableConfig) UNTAGGED_MINION_INSTANCE 특정 미니언 인스턴스로 태스크를 라우팅하는 태그
validateTaskConfigs(TableConfig, Schema, Map) no-op 테이블 생성 시 태스크별 구성을 검증

3단계: 테이블 구성에서 태스크 활성화 (Step 3: Enable the Task in Table Config)

테이블 구성에 태스크 구성을 추가해요:

{
  "tableName": "myTable_OFFLINE",
  "task": {
    "taskTypeConfigsMap": {
      "MyCustomTask": {
        "schedule": "0 0 * * * ?",
        "myParam": "myValue"
      }
    }
  }
}

schedule 속성은 CRON 표현식을 사용해 태스크 생성기가 실행되는 시점을 정의해요. 없으면 태스크는 컨트롤러 API를 통해서만 트리거될 수 있어요.

4단계: API로 태스크 트리거 (Step 4: Trigger Tasks via API)

태스크 생성을 수동으로 트리거하려면:

# Generate tasks for all tables with this task type
POST /tasks/schedule?taskType=MyCustomTask

# Generate tasks for a specific table
POST /tasks/execute?taskType=MyCustomTask&tableName=myTable_OFFLINE

내장 태스크 타입 (Built-in Task Types)

Pinot는 여러 내장 미니언 태스크를 포함해요:

태스크 타입 설명
RealtimeToOfflineSegmentsTask 실시간 세그먼트를 오프라인 세그먼트로 변환
UpsertCompactionTask upsert 지원 테이블의 세그먼트 컴팩션
UpsertCompactMergeTask upsert 세그먼트 병합·컴팩션
PurgeTask purge 함수와 일치하는 레코드 제거
MergeRollupTask 작은 세그먼트 병합·롤업
SegmentGenerationAndPushTask 외부 데이터에서 세그먼트 생성·푸시
RefreshSegmentTask 딥 스토어에서 세그먼트 새로고침

내장 태스크에 대한 자세한 내용은 미니언 Merge Rollup 태스크와 Upsert 컴팩션 태스크를 보세요.

더 알아보기 (Learn more)