미니언 태스크 플러그인
미니언 태스크 플러그인 (Minion Task Plugin)
태스크 생성기(generator)와 태스크 실행기(executor)를 가진 커스텀 미니언 태스크 플러그인을 작성하는 방법을 다루는 문서예요.
출처: 문서
본문
미니언 태스크는 Pinot에서 백그라운드 유지보수·데이터 처리 잡을 실행해요. 각 태스크 타입은 두 가지 컴포넌트를 가져요: 태스크 구성을 만드는 태스크 생성기 (Task Generator) (컨트롤러에서 실행)와 작업을 수행하는 태스크 실행기 (Task Executor) (미니언에서 실행).
아키텍처 (Architecture)
Controller Minion
┌────────────────────┐ ┌────────────────────┐
│ PinotTaskGenerator│ ─ Helix ─> │ PinotTaskExecutor │
│ - generateTasks() │ schedules │ - executeTask() │
│ - getTaskType() │ │ - cancel() │
└────────────────────┘ └────────────────────┘
- 컨트롤러는 스케줄 또는 ad-hoc API 호출로 태스크 생성기를 실행해요
- 생성기는 작업을 설명하는
PinotTaskConfig객체를 만들어요 - 컨트롤러는 Helix를 통해 태스크를 스케줄해요
- 미니언 인스턴스는 태그/라우팅을 기준으로 태스크를 집어요
- 미니언은 태스크 실행기를 사용해 태스크를 실행해요
- 결과는 컨트롤러로 반환돼요
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 컴팩션 태스크를 보세요.