すべてのプロダクト
Search
ドキュメントセンター

SchedulerX:Java ジョブ

最終更新日:Aug 28, 2026

アプリケーションプロセスで Java ジョブを実行できます。このトピックでは、Java ジョブの管理方法について説明します。

実行モード

Java ジョブは、次の実行モードをサポートしています。

  • スタンドアロン: ジョブは、同じ groupId 内のランダムなマシンで実行されます。

  • ブロードキャスト:ジョブは、同じ groupId 内のすべてのマシンで同時に実行されます。

  • ビジュアル MapReduce:Professional Edition が必要な MapReduce タスクモデルです。最大 1,000 個のサブタスクをサポートします。ビジネスキーワードを使用して、詳細な実行レコード、操作ログ、およびサブタスクスタックをクエリできます。

  • MapReduce:多数のサブタスクの並列処理をサポートする通常の MapReduce タスクモデルです。サブタスクの数が 100 万未満の場合は、このモードを選択してください。サブタスク実行の概要情報のみをクエリできます。

  • シャーディング実行:このモードでは、静的シャーディングと動的バッチ処理を使用してビッグデータを処理します。

スタンドアロンモードとブロードキャストモードでは JavaProcessor を実装します。ビジュアル MapReduce、MapReduce、およびシャーディング実行モードでは MapJobProcessor を実装します。

プロセッサクラスパスは、実装クラスの完全修飾名です。たとえば、com.apache.armon.test.schedulerx.processor.MySimpleJob のようになります。

  • JAR パッケージをアップロードしない場合、SchedulerX はアプリケーションのクラスパスでプロセッサの実装クラスを検索します。そのため、変更のたびにアプリケーションを再コンパイルして公開する必要があります。

  • JAR パッケージをアップロードすると、SchedulerX は JAR パッケージとプロセッサをホットリロードします。アプリケーションを再公開する必要はありません。

プログラミングモデル

Java ジョブは、JavaProcessor と MapJobProcessor の 2 つのプログラミングモデルをサポートしています。

  • JavaProcessor

    • 任意: public void preProcess(JobContext context) throws Exception

    • 必須: public ProcessResult process(JobContext context) throws Exception

    • 任意: public void postProcess(JobContext context)

    • 任意: public void kill(JobContext context)

  • MapJobProcessor

    • 必須: public ProcessResult process(JobContext context) throws Exception

    • 任意: public void postProcess(JobContext context)

    • 任意: public void kill(JobContext context)

    • 必須: public ProcessResult map(List<? extends Object> taskList, String taskName)

ProcessResult

各プロセスは ProcessResult オブジェクトを返す必要があります。このオブジェクトは、ジョブの実行ステータス、結果、およびエラーメッセージを示します。

  • タスクは正常に実行され、new ProcessResult(true) を返します。

  • タスクは、new ProcessResult(false, ErrorMsg) を返すか、例外をスローすると失敗します。

  • タスクが正常に実行されると、return new ProcessResult(true, result) が返されます。ここで、result は 1000 バイトを超えない文字列です。

HelloSchedulerX2.0 ジョブの例

@Component
public class MyProcessor1 extends JavaProcessor {

    @Override
    public ProcessResult process(JobContext context) throws Exception {
        // TODO: ここにビジネスロジックを実装します。
        System.out.println("Hello, schedulerx2.0!");
        return new ProcessResult(true);
    }
}            

キル機能をサポートするジョブの例

@Component
public class MyProcessor2 extends JavaProcessor {
    private volatile boolean stop = false;

    @Override
    public ProcessResult process(JobContext context) throws Exception {
        int N = 10000;
        while (!stop && N >= 0) {
            // TODO: ここにビジネスロジックを実装します。
            N--;
        }
        return new ProcessResult(true);
    }

    @Override
    public void kill(JobContext context) {
        stop = true;
    }

    @Override
    public void preProcess(JobContext context) {
        stop = false;  // ジョブが Spring を使用して起動され、Bean がシングルトンの場合、preProcess を使用してフラグをリセットします。
    }
}          

Map モデルを使用したバッチ処理のジョブ例

/**
 * 非パーティション化テーブルの分散バッチ処理。
 * 1. ルートタスクはテーブルをクエリして minId と maxId を取得します。
 * 2. PageTask を構築し、map メソッドを使用して分散します。
 * 3. 次のレベルが PageTask を受信した場合、データを処理します。
 *
 */
@Component
public class ScanSingleTableJobProcessor extends MapJobProcessor {
    private static final int pageSize = 100;

    static class PageTask {
        private int startId;
        private int endId;
        public PageTask(int startId, int endId) {
             this.startId = startId;
             this.endId = endId;
        }
        public int getStartId() {
              return startId;
        }
        public int getEndId() {
              return endId;
        }
    }

    @Override
    public ProcessResult process(JobContext context) {
        String taskName = context.getTaskName();
        Object task = context.getTask();
        if (isRootTask(context)) {
            System.out.println("start root task");
            Pair<Integer, Integer> idPair = queryMinAndMaxId();
            int minId = idPair.getFirst();
            int maxId = idPair.getSecond();
            List<PageTask> taskList = Lists.newArrayList();
            for (int i = minId; i < maxId; i += pageSize) {
                taskList.add(new PageTask(i, (i + pageSize > maxId ? maxId : i + pageSize)));
            }
            return map(taskList, "Level1Dispatch"); // map メソッドを呼び出してサブタスクを分散します。
        } else if (taskName.equals("Level1Dispatch")) {
            PageTask record = (PageTask)task;
            long startId = record.getStartId();
            long endId = record.getEndId();
            // TODO: ここにビジネスロジックを実装します。
            return new ProcessResult(true);
        }

        return new ProcessResult(true);
    }

    @Override
    public void postProcess(JobContext context) {
        // TODO: ここにビジネスロジックを実装します。
        System.out.println("All tasks are finished.");
    }

    private Pair<Integer, Integer> queryMinAndMaxId() {
        // TODO: xxx から min(id) と max(id) を選択するクエリを実装します。
        return null;
    }

}