このトピックでは、Python ソフトウェア開発キット (SDK) を使用してジョブを送信する方法について説明します。このジョブは、ログファイル内の "INFO"、"WARN"、"ERROR"、"DEBUG" の出現回数をカウントします。
ジョブの準備
データファイルの OSS へのアップロード
タスクプログラムの OSS へのアップロード
SDK を使用したジョブの作成 (送信)
結果の表示
1. ジョブの準備
このジョブは、ログファイル内の "INFO"、"WARN"、"ERROR"、"DEBUG" の出現回数をカウントします。
このジョブには、split、count、merge の 3 つのタスクが含まれます。
split タスクは、ログファイルを 3 つの部分に分割します。
count タスクは、各部分における "INFO"、"WARN"、"ERROR"、"DEBUG" の出現回数をカウントします。count タスクの InstanceCount は 3 に設定する必要があります。これにより、3 つのインスタンスが開始され、count プログラムが同時に実行されます。
merge タスクは、count タスクからの結果を結合します。
(1) データファイルの OSS へのアップロード
この例のデータファイルをダウンロードします: log-count-data.txt
log-count-data.txt を次のパスにアップロードします: oss://your-bucket/log-count/log-count-data.txt
your-bucket を、作成したバケットの名前に置き換えます。この例では、バケットが中国 (深セン) リージョンにあることを前提としています。
Object Storage Service (OSS) へのファイルのアップロード方法の詳細については、「 ファイルのアップロード」をご参照ください。
(2) タスクプログラムの OSS へのアップロード
この例では、Python で記述されたジョブプログラムを使用します。この例のプログラムをダウンロードします: log-count.tar.gz
サンプルコードを変更する必要はありません。log-count.tar.gz ファイルを直接 OSS パスにアップロードします。例: oss://your-bucket/log-count/log-count.tar.gz
アップロード方法は前のセクションで説明したとおりです。
BatchCompute は .tar.gz 拡張子の圧縮パッケージのみをサポートします。ファイルは gzip を使用してパッケージ化する必要があります。そうしないと、パッケージを解析できません。
コードを修正する場合は、パッケージを解凍し、変更を加えてから、ファイルを再パッケージ化してください。
次のコマンドを実行します:
> cd log-count #フォルダに移動します > tar -czf log-count.tar.gz * #このフォルダ内のすべてのファイルを log-count.tar.gz にパッケージ化します次のコマンドを実行して、圧縮パッケージの内容を表示します:
$ tar -tvf log-count.tar.gz次のリストが表示されます:
conf.py count.py merge.py split.py
2. SDK を使用したジョブの作成 (送信)
Python SDK のダウンロードとインストールの詳細については、「Python SDK のダウンロードとインストール」をご参照ください。
API バージョン 20151111 を使用する場合、クラスター ID を指定するか、匿名クラスターパラメーターを使用してジョブを送信する必要があります。この例では、匿名クラスターを使用します。匿名クラスターには 2 つのパラメーターを設定する必要があります:
利用可能なイメージ ID。システム提供のイメージを使用するか、カスタムイメージを作成できます。詳細については、「 カスタムイメージ」をご参照ください。
インスタンスタイプ。詳細については、「サポートされているインスタンスタイプ」をご参照ください。
OSS で、プログラム出力 (StdoutRedirectPath) とエラーログ (StderrRedirectPath) を保存するパスを作成します。この例では、パスは oss://your-bucket/log-count/logs/ です。
この例を実行するには、プログラムのコメントにある変数を実際の OSS パスに変更する必要があります。
以下のコードは、Python SDK を使用してジョブをサブミットするためのテンプレートです。プログラム内のパラメーターの詳細については、「パラメーターの説明」をご参照ください。
#encoding=utf-8
import sys
from batchcompute import Client, ClientError
from batchcompute import CN_SHENZHEN as REGION # 必要に応じてリージョンを設定します。
from batchcompute.resources import (
JobDescription, TaskDescription, DAG, AutoCluster, Configs, Networks, VPC,
)
ACCESS_KEY_ID='' # ご利用の AccessKey ID を入力します。
ACCESS_KEY_SECRET='' # ご利用の AccessKey Secret を入力します。
IMAGE_ID = 'img-ubuntu' # ご利用のイメージ ID をここに入力します。
INSTANCE_TYPE = 'ecs.sn1.medium' # リージョンでサポートされているインスタンスタイプを入力します。
WORKER_PATH = '' # 'oss://your-bucket/log-count/log-count.tar.gz' log-count.tar.gz をアップロードした OSS パスを入力します。
LOG_PATH = '' # 'oss://your-bucket/log-count/logs/' エラーフィードバックとタスク出力用に作成した OSS パスを入力します。
OSS_MOUNT= '' # 'oss://your-bucket/log-count/' /home/inputs と /home/outputs に同時にマウントします。
client = Client(REGION, ACCESS_KEY_ID, ACCESS_KEY_SECRET)
def main():
try:
job_desc = JobDescription()
# 自動クラスターを作成します。
cluster = AutoCluster()
cluster.InstanceType = INSTANCE_TYPE
cluster.ResourceType = "OnDemand"
cluster.ImageId = IMAGE_ID
configs = Configs()
networks = Networks()
vpc = VPC()
vpc.CidrBlock = '192.168.0.0/16'
# vpc.VpcId = "vpc-8vbfxdyhxxxx"
networks.VPC = vpc
configs.Networks = networks
# システムディスクのタイプ (cloud_efficiency または cloud_ssd) とサイズ (GB) を設定します。
configs.add_system_disk(size=40, type_='cloud_efficiency')
configs.InstanceCount = 1
cluster.Configs = configs
# split タスクを作成します。
split_task = TaskDescription()
split_task.Parameters.Command.CommandLine = "python split.py"
split_task.Parameters.Command.PackagePath = WORKER_PATH
split_task.Parameters.StdoutRedirectPath = LOG_PATH
split_task.Parameters.StderrRedirectPath = LOG_PATH
split_task.InstanceCount = 1
split_task.AutoCluster = cluster
split_task.InputMapping[OSS_MOUNT]='/home/input'
split_task.OutputMapping['/home/output'] = OSS_MOUNT
# map タスクを作成します。
count_task = TaskDescription(split_task)
count_task.Parameters.Command.CommandLine = "python count.py"
count_task.InstanceCount = 3
count_task.InputMapping[OSS_MOUNT] = '/home/input'
count_task.OutputMapping['/home/output'] = OSS_MOUNT
# merge タスクを作成します。
merge_task = TaskDescription(split_task)
merge_task.Parameters.Command.CommandLine = "python merge.py"
merge_task.InstanceCount = 1
merge_task.InputMapping[OSS_MOUNT] = '/home/input'
merge_task.OutputMapping['/home/output'] = OSS_MOUNT
# タスク DAG を作成します。
task_dag = DAG()
task_dag.add_task(task_name="split", task=split_task)
task_dag.add_task(task_name="count", task=count_task)
task_dag.add_task(task_name="merge", task=merge_task)
task_dag.Dependencies = {
'split': ['count'],
'count': ['merge']
}
# ジョブ記述を作成します。
job_desc.DAG = task_dag
job_desc.Priority = 99 # 0-1000
job_desc.Name = "log-count"
job_desc.Description = "PythonSDKDemo"
job_desc.JobFailOnInstanceFail = True
job_id = client.create_job(job_desc).Id
print('job created: %s' % job_id)
except ClientError, e:
print (e.get_status_code(), e.get_code(), e.get_requestid(), e.get_msg())
if __name__ == '__main__':
sys.exit(main())3. ジョブステータスの表示
SDK の Get Job メソッドを使用して、ジョブステータスを取得します。詳細については、「 ジョブ情報の取得」をご参照ください。
jobInfo = client.get_job(job_id)
print (jobInfo.State)ジョブは、Waiting、Running、Finished、Failed、Stopped のいずれかの状態になります。
4. 結果の表示
OSS コンソールにログインして、バケット内の /log-count/merge_result.json ファイルを表示します。
内容は次のようになります:
{"INFO": 2460, "WARN": 2448, "DEBUG": 2509, "ERROR": 2583}OSS SDK を使用して、結果を取得することもできます。