このドキュメントでは、Java SDK を使用して、ログファイル内の "INFO"、"WARN"、"ERROR"、"DEBUG" の出現回数をカウントするジョブを送信する方法について説明します。
手順
ジョブの準備
データファイルの OSS へのアップロード
サンプルコードの使用
コードのコンパイルとパッケージ化
パッケージの OSS へのアップロード
SDK を使用したジョブの作成と送信
結果の表示
1. ジョブの準備
このジョブは、ログファイル内の "INFO"、"WARN"、"ERROR"、"DEBUG" の出現回数をカウントします。
このジョブは、分割 (split)、カウント (count)、マージ (merge) の 3 つのタスクで構成されます。
分割タスクは、ログファイルを 3 つの部分に分割します。
カウントタスクは、各部分の出現回数をカウントします。このタスクには 3 のインスタンス数が必要で、3 つのインスタンスが同時にカウントプログラムを実行することを意味します。
マージタスクは、カウントタスクの結果を統合します。
(1) データファイルの OSS へのアップロード
この例のデータファイルをダウンロードします: log-count-data.txt
log-count-data.txt を次のパスにアップロードします: oss://your-bucket/log-count/log-count-data.txt
your-bucketをご利用のバケット名に置き換えてください。このクイックスタートでは、リージョンが中国 (深セン) であることを前提としています。
(2) サンプルコードの使用
この例では、Java を使用してジョブタスクを作成し、Maven を使用してコードをコンパイルします。無料の IntelliJ IDEA Community Edition の使用を推奨します: https://www.jetbrains.com/idea/download/。
サンプルプログラムをダウンロードします: java-log-count.zip
これは Maven プロジェクトです。
コードを変更する必要はありません。
(3) コンパイルとパッケージ化
次のコマンドを実行して、コードをコンパイルおよびパッケージ化します。
mvn packageこのコマンドは、target ディレクトリに次の 3 つの JAR ファイルを生成します。
batchcompute-job-log-count-1.0-SNAPSHOT-Split.jar
batchcompute-job-log-count-1.0-SNAPSHOT-Count.jar
batchcompute-job-log-count-1.0-SNAPSHOT-Merge.jar次に、次のコマンドを使用して、3 つの JAR ファイルを tar.gz アーカイブに圧縮します。
> cd target # target ディレクトリに移動します。
> tar -czf worker.tar.gz *SNAPSHOT-*.jar # ファイルをパッケージ化します。次のコマンドを実行して、パッケージの内容を確認します。
> tar -tvf worker.tar.gz
batchcompute-job-log-count-1.0-SNAPSHOT-Split.jar
batchcompute-job-log-count-1.0-SNAPSHOT-Count.jar
batchcompute-job-log-count-1.0-SNAPSHOT-Merge.jarBatchCompute は、.tar.gz 拡張子を持つ圧縮パッケージのみをサポートします。例に示すように、gzip を使用してファイルをパッケージ化する必要があります。そうしないと、パッケージを解析できません。
(4) パッケージの OSS へのアップロード
worker.tar.gz を OSS のバケットにアップロードします: oss://your-bucket/log-count/worker.tar.gz
この例を実行するには、独自のバケットを作成し、そこに worker.tar.gz ファイルをアップロードしてください。
2. SDK を使用したジョブの送信
(1) Maven プロジェクトの作成
プロジェクトの pom.xml ファイルに次の依存関係を追加します。
<dependencies>
<dependency>
<groupId>com.aliyun</groupId>
<artifactId>aliyun-java-sdk-batchcompute</artifactId>
<version>5.2.0</version>
</dependency>
<dependency>
<groupId>com.aliyun</groupId>
<artifactId>aliyun-java-sdk-core</artifactId>
<version>3.2.3</version>
</dependency>
</dependencies>最新バージョンの Java SDK を使用してください。
(2) Demo.java の作成
ジョブを送信する際には、クラスター ID または自動クラスターのパラメーターのいずれかを指定する必要があります。この例では自動クラスターを使用しており、次の 2 つのパラメーターが必要です。
利用可能なイメージ ID。システム提供のイメージまたはカスタムイメージを使用できます。
インスタンスタイプ。詳細については、「現在サポートされているタイプ」をご参照ください。
OSS で、プログラム出力 (StdoutRedirectPath) とエラーログ (StderrRedirectPath) を保存するためのパスを作成します。この例では、パスは oss://your-bucket/log-count/logs/ です。
この例を実行するには、コード内のプレースホルダー変数を、ご自身の認証情報と作成した OSS パスに置き換えてください。
次のテンプレートは、Java SDK を使用してジョブを送信する方法を示しています。パラメーターの詳細については、「SDK インターフェイスの説明」をご参照ください。
Demo.java:
/*
* IMAGE_ID: ECS イメージ ID。
* INSTANCE_TYPE: インスタンスタイプ。
* REGION_ID: ジョブが送信されるリージョン。worker パッケージのバケットと同じリージョンである必要があります。
* ACCESS_KEY_ID: ご自身の AccessKey ID。
* ACCESS_KEY_SECRET: ご自身の AccessKey Secret。
* WORKER_PATH: worker パッケージの OSS パス。
* LOG_PATH: エラーログとタスク出力を保存するための OSS パス。このディレクトリは事前に作成しておく必要があります。
*/
import com.aliyuncs.batchcompute.main.v20151111.*;
import com.aliyuncs.batchcompute.model.v20151111.*;
import com.aliyuncs.batchcompute.pojo.v20151111.*;
import com.aliyuncs.exceptions.ClientException;
import java.util.ArrayList;
import java.util.List;
public class Demo {
static String IMAGE_ID = "img-ubuntu";; // ご自身の ECS イメージ ID を入力します。
static String INSTANCE_TYPE = "ecs.sn1.medium"; // ご利用のリージョンで利用可能なインスタンスタイプを入力します。
static String REGION_ID = "cn-shenzhen"; // リージョン ID を入力します。
static String ACCESS_KEY_ID = ""; // ご自身の AccessKey ID を入力します。
static String ACCESS_KEY_SECRET = ""; // ご自身の AccessKey Secret を入力します。
static String WORKER_PATH = ""; // worker.tar.gz ファイルの OSS パスを入力します (例: "oss://your-bucket/log-count/worker.tar.gz")。
static String LOG_PATH = ""; // ログ用の OSS パスを入力します (例: "oss://your-bucket/log-count/logs/")。
static String MOUNT_PATH = ""; // 例: "oss://your-bucket/log-count/"。
public static void main(String[] args){
/** BatchCompute クライアントを作成します。 */
BatchCompute client = new BatchComputeClient(REGION_ID, ACCESS_KEY_ID, ACCESS_KEY_SECRET);
try{
/** JobDescription オブジェクトを構築します。 */
JobDescription jobDescription = genJobDescription();
// ジョブを作成します。
CreateJobResponse response = client.createJob(jobDescription);
// ジョブが作成されると、ジョブ ID が返されます。
String jobId = response.getJobId();
System.out.println("Job created successfully. Job ID: "+jobId);
// ジョブステータスを照会します。
GetJobResponse getJobResponse = client.getJob(jobId);
Job job = getJobResponse.getJob();
System.out.println("Job state:"+job.getState());
} catch (ClientException e) {
e.printStackTrace();
System.out.println("Job creation failed. Error code: "+ e.getErrCode()+", Error message: "+e.getErrMsg());
}
}
private static JobDescription genJobDescription(){
JobDescription jobDescription = new JobDescription();
jobDescription.setName("java-log-count");
jobDescription.setPriority(0);
jobDescription.setDescription("log-count demo");
jobDescription.setJobFailOnInstanceFail(true);
jobDescription.setType("DAG");
DAG taskDag = new DAG();
/** 分割タスクを追加します。 */
TaskDescription splitTask = genTaskDescription();
splitTask.setTaskName("split");
splitTask.setInstanceCount(1);
splitTask.getParameters().getCommand().setCommandLine("java -jar batchcompute-job-log-count-1.0-SNAPSHOT-Split.jar");
taskDag.addTask(splitTask);
/** カウントタスクを追加します。 */
TaskDescription countTask = genTaskDescription();
countTask.setTaskName("count");
countTask.setInstanceCount(3);
countTask.getParameters().getCommand().setCommandLine("java -jar batchcompute-job-log-count-1.0-SNAPSHOT-Count.jar");
taskDag.addTask(countTask);
/** マージタスクを追加します。 */
TaskDescription mergeTask = genTaskDescription();
mergeTask.setTaskName("merge");
mergeTask.setInstanceCount(1);
mergeTask.getParameters().getCommand().setCommandLine("java -jar batchcompute-job-log-count-1.0-SNAPSHOT-Merge.jar");
taskDag.addTask(mergeTask);
/** タスクの依存関係を追加します: split-->count-->merge */
List<String> taskNameTargets = new ArrayList();
taskNameTargets.add("merge");
taskDag.addDependencies("count", taskNameTargets);
List<String> taskNameTargets2 = new ArrayList();
taskNameTargets2.add("count");
taskDag.addDependencies("split", taskNameTargets2);
//dag
jobDescription.setDag(taskDag);
return jobDescription;
}
private static TaskDescription genTaskDescription(){
AutoCluster autoCluster = new AutoCluster();
autoCluster.setInstanceType(INSTANCE_TYPE);
autoCluster.setImageId(IMAGE_ID);
//autoCluster.setResourceType("OnDemand");
TaskDescription task = new TaskDescription();
//task.setTaskName("Find");
// VPC を使用する場合は、cidrBlock を設定する必要があります。IP アドレス範囲が他の範囲と競合しないようにしてください。
Configs configs = new Configs();
Networks networks = new Networks();
VPC vpc = new VPC();
vpc.setCidrBlock("192.168.0.0/16");
networks.setVpc(vpc);
configs.setNetworks(networks);
autoCluster.setConfigs(configs);
// アップロードされたジョブパッケージの完全な OSS パス。
Parameters p = new Parameters();
Command cmd = new Command();
//cmd.setCommandLine("");
// アップロードされたジョブパッケージの完全な OSS パス。
cmd.setPackagePath(WORKER_PATH);
p.setCommand(cmd);
// エラーログを保存するパス。
p.setStderrRedirectPath(LOG_PATH);
// 最終出力を保存するパス。
p.setStdoutRedirectPath(LOG_PATH);
task.setParameters(p);
task.addInputMapping(MOUNT_PATH, "/home/input");
task.addOutputMapping("/home/output",MOUNT_PATH);
task.setAutoCluster(autoCluster);
//task.setClusterId(clusterId);
task.setTimeout(30000); /* 30,000 秒 */
task.setInstanceCount(1); /** 1 つのインスタンスで実行します。 */
return task;
}
}出力例:
Job created successfully. Job ID: job-01010100010192397211
Job state:Waiting3. ジョブステータスの確認
SDK の ジョブ情報の取得 メソッドを呼び出すことで、ジョブステータスを確認できます。
// ジョブステータスを照会します。
GetJobResponse getJobResponse = client.getJob(jobId);
Job job = getJobResponse.getJob();
System.out.println("Job state:"+job.getState());ジョブの状態は、Waiting、Running、Finished、Failed、Stopped のいずれかです。
4. 結果の表示
BatchCompute コンソールにログインして、ジョブステータスを表示できます。
ジョブが完了したら、OSS コンソールにログインし、バケット内の /log-count/merge_result.json ファイルを表示します。
ファイルの内容は次のとおりです。
{"INFO": 2460, "WARN": 2448, "DEBUG": 2509, "ERROR": 2583}また、OSS SDK を使用して結果を取得することもできます。