アプリケーションプロセスで 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;
}
}