Hydrosphere.ioについて
参考
メモ
Hydroshpereブログの最も古い記事 が投稿されたのは、2016/6/14である。 上記ブログで強調されていたのは、Big DataプロジェクトのDevOps対応。 DockerやAnsibleといった道具を使いながら。
Hydrosphere Serving
参考情報
開発状況
2019/03/21現在でも、わりと活発に開発されているようだ。 https://github.com/Hydrospheredata/hydro-serving/graphs/commit-activity
開発元は、hydrosphere.io。パロアルトにある企業らしい。
hydrosphere.ioのプロダクト
ここで取り上げているServingを含む、以下のプロダクトがある。
- Serving: モデルのサーブ、アドホック分析への対応
- Sonar: モデルやパイプラインの品質管理
- Mist: Sparkの計算リソースをREST API経由で提供する仕組み(マルチテナンシーの実現)
概要
GitHubのREADMEに特徴が書かれていが、その中でも個人的にポイントと思ったのは以下の通り。
- Envoyプロキシを用いてサービスメッシュ型のサービングを実現
- 複数の機械学習モデルに対応し、パイプライン化も可能
- UIがある
Hydro Servingの公式ウェブサイト に掲載されていた動画を見る限り、 コマンドライン経由でモデルを登録することもでき、モデルを登録したあとは、 処理フロー(ただし、シリアルなフローに見える)を定義可能。 例えば、機械学習モデルの推論器に渡す前処理・後処理も組み込めるようだ。
セットアップ方法
Docker Composeを使う方法とk8sを使う方法があるようだ。
利用方法
モデルを学習させ、出力する。(例では、h5形式で出力していた) モデルや必要なライブラリを示したテキストなどを、ペイロードとして登楼する。
必要なライブラリの指定は以下のようにする。
requirements.txt 1
2
3Keras==2.2.0
tensorflow==1.8.0
numpy==1.13.3
また、それらの規約事項は、contract(定義ファイル)として保存する。 フォルダ内の構成は以下のようになる。
公式ドキュメントから引用 1
2
3
4
5
6
7linear_regression
├── model.h5
├── model.py
├── requirements.txt
├── serving.yaml
└── src
└── func_main.py
上記のようなファイル群を作り、
hs uploadコマンドを使ってアップロードする。
アップロードしたあとは、ウェブフロントエンドから処理フローを定義する。
クエリはREST APIで以下のように投げる。
公式ウェブサイトから引用 1
2$ curl -X POST --header 'Content-Type: application/json' --header 'Accept: application/json' -d '{
"x": [[1, 1],[1, 1]]}' 'http://localhost/gateway/applications/linear_regression/infer'
また上記のREST APIの他にも、gRPCを用いて推論結果を取得することも可能。 また、データ構造としてはTensorProtoを使える。
なお、TensorFlowモデルのサーブについては、 公式ウェブサイトの Hydro ServingでTensorFlowモデルをサーブ がわかりやすい。 また、単純にPython関数を渡すこともできる。(必ずしもTensorFlowにロックインではない)
コンセプト
- モデル
- Hydro Servingに渡されたモデルはバージョン管理される。
- フレーワークは多数対応
- ただしフレームワークによって、出力されるメタデータの情報量に差がある。
- アプリケーション
- 単一で動かす方法とパイプラインを構成する方法がある
- ランタイム
- 予め実行環境を整えたDockerイメージが提供されている
- Python、TensorFlow、Spark
動作確認
Hydro Servingの公式ドキュメント に従って動作確認する。
1 | $ mkdir HydroServing |
起動したコンテナを確認する。
1 | $ sudo docker-compose ps |
その他CLIを導入する。
1 | $ conda create -n hydro-serving python=3.6 python |
エラーが生じた。
1 | $ hs cluster add --name local --server http://localhost |
エラーが出ているのに登録されたように見える。 念の為、クラスタ情報を確認する。
1 | $ hs cluster |
いったんこのまま進める。 まずはアプリを作成。
1 | $ mkdir -p ~/Sources/linear_regression |
学習し、モデルを出力。
1 | $ python model.py |
サーブ対象となる関数のアプリを作成する。 1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28$ mkdir src
$ cd src
$ cat << EOF > func_main.py
import numpy as np
import hydro_serving_grpc as hs
from keras.models import load_model
# 0. Load model once
model = load_model('/model/files/model.h5')
def infer(x):
# 1. Retrieve tensor's content and put it to numpy array
data = np.array(x.double_val)
data = data.reshape([dim.size for dim in x.tensor_shape.dim])
# 2. Make a prediction
result = model.predict(data)
# 3. Pack the answer
y_shape = hs.TensorShapeProto(dim=[hs.TensorShapeProto.Dim(size=-1)])
y_tensor = hs.TensorProto(
dtype=hs.DT_DOUBLE,
double_val=result.flatten(),
tensor_shape=y_shape)
# 4. Return the result
return hs.PredictResponse(outputs={"y": y_tensor})
EOF
1 | $ cd .. |
つづいて必要ライブラリを定義する。 1
2
3
4
5$ cat << EOF > requirements.txt
Keras==2.2.0
tensorflow==1.8.0
numpy==1.13.3
EOF
最終的に以下のような構成になった。
1 | $ tree |
つづいてモデルなどをアップロード。 1
$ hs upload
ここで以下のようなエラーが生じた。 1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16Using 'local' cluster
['/home/centos/Sources/linear_regression/src', '/home/centos/Sources/linear_regression/requirements.txt', '/home/centos/Sources/linear_regression/model.h5']
Packing the model [####################################] 100%
Assembling the model [####################################] 100%
Uploading to http://localhost
Uploading model assembly [####################################] 100%
Traceback (most recent call last):
File "/home/centos/.conda/envs/hydro-serving/lib/python3.6/site-packages/hydroserving/httpclient/remote_connection.py", line 82, in postprocess_response
response.raise_for_status()
File "/home/centos/.conda/envs/hydro-serving/lib/python3.6/site-packages/requests/models.py", line 940, in raise_for_status
raise HTTPError(http_error_msg, response=self)
requests.exceptions.HTTPError: 404 Client Error: Not Found for url: http://localhost/api/v1/model/upload
During handling of the above exception, another exception occurred:
(snip)
軽く調べると ISSUE255 が 関係していそうである。
ここでは動作確認のため、pip install hs==2.0.0rc2としてインストールして試した。
1 | $ pip uninstall hs |
また、 serving.yaml に以下のエントリを追加した。
1 | runtime: "hydrosphere/serving-runtime-python:3.6-latest" |
再び hs upload したところ完了のように見えた。
ウェブUIを確認したところ以下の通り。
それではテストを実行してみる。 1
$ curl -X POST --header 'Content-Type: application/json' --header 'Accept: application/json' -d '{ "x": [ [ 1, 1 ] ] }' 'http://10.0.0.209/gateway/application/linear_regression'
以下のようなエラーが生じた。 1
{"error":"InternalUncaught","information":"UNKNOWN: Exception calling application: No module named 'numpy'"}
ランタイム上にnumpyがインストールされなかったのだろうか・・・
WIP
Hydro sonar
参考
- Hydro Sonarの公式ウェブサイト
- [公式のGUI画面]
メモ
公式のSonarのGUI画面 を見ると、モデルの性能を観測するための仕組みに見える。 モデル選択、メンテナンス、リアイアメントのために用いられる。 例えば、入力データ変化を起因したモデル精度の劣化など。
Hydro Mist
参考
- Hydro Mistの公式ウェブサイト
- 公式のMistのアーキテクチャイメージ
- Hydro Mistの公式ドキュメント
- Hydro Mistの公式クイックスタート
- Hydro Mistがどうやってジョブを起動するか
- Hydro Mistのジョブのステート
- Hydro Mistのジョブローンチの流れ
- Hydro Mistのコンテキスト管理
- Hydro MIstのmist-cliのGitHub
- Hydro Mistのコンフィグの例
- Hydro MistのScala API
メモ
公式のMistのアーキテクチャイメージ を見ると、 Sparkのマルチテナンシーを実現するものに見える。
特徴として挙げられていたものの中で興味深いのは以下。
- Spark Function as a Service
- ユーザのAPIをSparkのコンフィグレーションから分離する
- HTTP、Kafka、MQTTによるやりとり
Hydro Mistの公式クイックスタートを見ると、Scala、Java、Pythonのサンプルアプリが掲載されている。 Mistはライブラリとしてインポートし、フレームワークに則ってアプリを作成すると、 アプリ内で定義された関数を実行するジョブをローンチできるようになる。
また実際にローンチするときには、REST API等で起動することになるが、 そのときに引数を渡すこともできる。
動作確認
Hydro Mistの公式クイックスタート に従い動作を確認する。 Dockerで実行するか、バイナリで実行するかすれば良い。
Dockerでの起動
Dockerの場合は以下の通り。
1 | $ sudo docker run -p 2004:2004 \ |
バイナリを配備しての起動
バイナリをダウンロードして実行する場合は以下の通り。 まずSparkをダウンロードする。(なお、JDKは予めインストールされていることを前提とする)
1 | $ mkdir ~/Spark |
これで、~/Spark/default以下にSparkが配備された。 つづいて、Hydro Mistのバイナリをダウンロードし、配備する。
1 | $ mkdir ~/HydroMist |
以上でHydro Mistが配備された。 それではHydro Mistのマスタを起動する。
1 | $ SPARK_HOME=${HOME}/Spark/default ./bin/mist-master start --debug true |
以上で、バイナリを配備したHydro Mistが起動する。
Mistの動作確認
以降、サンプルアプリを実行している。
まずmistのCLIを導入する。
1 | $ conda create -n mist python=3.6 python |
続いて、サンプルプロジェクトをcloneする。
1 | $ mkdir -p ~/Sources |
試しにScala版を動かす。 (なお、SBTがインストールされていることを前提とする)
1 | $ cd hello_mist/scala |
上記コマンドの結果、以下のような出力が得られる。
1 | Process 5 file entries |
試しにstandaloneのコンテキストを確認してみる。
(なお、上記メッセージではPOSTを使うよう書かれているが、GETでないとエラーになった)
1 | $ curl -H 'Content-Type: application/json' -X GET http://localhost:2004/v2/api/contexts/standalone | jq |
なお、サンプルアプリは以下の通り。 Piの値を簡易的に計算するものである。
1 | import mist.api._ |
結果として、以下のようにジョブが登録される。
登録されたジョブを実行する。
また上記で表示されていたcurl ...をコマンドで実行することでも、ジョブを走らせることができる。
ジョブの一覧は、ウェブUIから以下のように確認できる。
ジョブ一覧からジョブを選ぶと、そのメタデータ、渡されたパラメータ、結果などを確認できる。 さらにジョブ実行時のログもウェブUIから確認できる。
Sparkのドライバログと思われるものも含まれている。
アーキテクチャと動作
Hydro Mistがどうやってジョブを起動するか によると以下の通り。
Mistで実行する関数は、mist-workerにラップされており、当該ワーカがSparkContextを保持しているようだ。
またワーカを通じてジョブを実行する。 (参考:Hydro
Mistのジョブのステート )
ユーザがジョブを実行させようとしたとき、Mistのマスタはそのリクエストをキューに入れ、ワーカが空くのを待つ。 ワーカのモードには2種類がある。
- exclusive
- ジョブごとにワーカを立ち上げる
- shared
- ひとつのジョブが完了してもワーカを立ち上げたままにし再利用する。
これにより、ジョブの並列度、ジョブの耐障害性の面で有利になる。
また、 Hydro Mistのジョブローンチの流れ からジョブ登録から実行までの大まかな流れがわかる。
コンテキストの管理
Sparkのコンテキスト管理やワーカのモード設定は、 mist-cliを通じてできるようだ。 Hydro Mistのコンテキスト管理 参照。
mist-cli
Hydro
MIstのmist-cliのGitHub 参照。 mist-cli
は、コンフィグファイル(やコンフィグファイルが配備されたディレクトリ)を指定しながら実行する。
ディレクトリにコンフィグを置く場合は、ファイル名の先頭に2桁の数字を入れることで、
プライオリティを指定することができる。
コンフィグには、Artifact、Context、Functionの設定を記載する。 参考:Hydro Mistのコンフィグの例
Hydro
Mistの公式クイックスタート を実行したあとの状態で確認してみる。
1
2
3
4
5
6
7
8
9
10
11
12$ mist-cli list contexts
ID WORKER MODE
default exclusiveemr_ctx exclusiveemr_autoscale_ctx exclusivestandalone exclusive
$ mist-cli list functions
FUNCTION DEFAULT CONTEXT PATH CLASS NAME
hello-mist-scala default hello-mist-scala_0.0.1.jar HelloMist$
$ mist-cli status
Mist version: 1.1.1
Spark version: 2.4.0
Java version: 1.8.0_201-b09
いくつかのコンテキストと、関数が登録されていることがわかる。 なお、GitHub上のバイナリを用いたので使用するSparkのバージョンが2.4.0になっていた。
Scala API
Hydro
MistのScala API
とサンプルアプリ(mist\examples\examples\src\main\scala\PiExample.scala)を確認してみる。
エントリポイント
MistFnがエントリポイントのようだ。
PiExample.scala:6 1
2
3
4
5
6object PiExample extends MistFn {
override def handle: Handle = {
val samples = arg[Int]("samples").validated(_ > 0, "Samples should be positive")
(snip)
引数定義
また関数の引数設定はmist.api.ArgsInstances#arg[A]などで行う。
下記にサンプルに記載の例を示す。
PiExample.scala:9 1
val samples = arg[Int]("samples").validated(_ > 0, "Samples should be positive")
[Int]により引数として渡される値の型を指定する。argメソッドの戻り値はNamedUserArgクラスだが、
当該クラスはUserArg[A]トレートを拡張している。
UserArg#validatedメソッドを用いることで、値の検証を行う。
複数の引数を渡すには、withArgsを使うか、combineや&を使う。
例: 1
2
3val three = withArgs(arg[Int]("n"), arg[String]("str"), arg[Boolean]("flag"))
val three = arg[Int]("n") & arg[String]("str") & arg[Boolean]("flag")
またドキュメントには、case classと各種Extractorを用いることで、JSON形式の入力データを 扱う方法が説明されていた。
Sparkのコンテキスト
引数を定義したのちは、Mistのコンテキスト管理のAPIを利用し、
SparkSessionなどを取得する。
onSparkContextやonSparkSessionなどを利用可能。
mist/api/MistFnSyntax.scala:91 1
2
3
4
5def onSparkSession[F, Cmb, Out](f: F)(
implicit
cmb: ArgCombiner.Aux[A, SparkSession, Cmb],
fnT: FnForTuple.Aux[Cmb, F, Out]
): RawHandle[Out] = args.combine(sparkSessionArg).apply(f)
サンプルでは以下のような使い方を示している。
PiExample.scala:10 1
2
3
4 withArgs(samples).onSparkContext((n: Int, sc: SparkContext) => {
val count = sc.parallelize(1 to n).filter(_ => {
(snip)
もしSparkSessionを使うならば以下の通りか。
1 | import org.apache.spark.sql.SparkSession |
なお、onStreamingContextというのもあり、Spark
Streamingも実行可能のようだ。
結果のハンドリング
上記onSparkContextメソッドなどの戻り値の型はRawHandle[Out]である。
mist/api/MistFnSyntax.scala:58 1
2
3
4
5def onSparkContext[F, Cmb, Out](f: F)(
implicit
cmb: ArgCombiner.Aux[A, SparkContext, Cmb],
fnT: FnForTuple.Aux[Cmb, F, Out]
): RawHandle[Out] = args.combine(sparkContextArg).apply(f)
最終的に結果をJSON形式で返すために、asHandleメソッドを利用する。
PiExample.scala:10 1
2
3
4
5withArgs(samples).onSparkContext((n: Int, sc: SparkContext) => {
(snip)
}).asHandle
asHandleメソッドは以下の通り。
mist/api/MistFnSyntax.scala:48 1
2
3implicit class AsHandleOps[A](raw: RawHandle[A]) {
def asHandle(implicit enc: JsEncoder[A]): Handle = raw.toHandle(enc)
}
余談:implicitクラスの使用
mist.api以下で複数のクラスがimplicit定義されて利用されており、一見して解析しづらい印象を覚えた。
例えばonSparkContextやonSparkSessionがimplicitクラスContextsOpsクラスに定義されている。
考察:フレームワークの良し悪し
Sparkの単純なプロキシではなく、フレームワーク(ライブラリ)化されていることで、 ユーザはSparkのコンテキストの管理から開放される、という利点がある。
一方で、Hydro Mistのフレームワークに従って実装する必要があり、多少なり
- Sparkに加えて、Hydro Mistの学習コストがある
- Mistのフレームワークでは実装しづらいケースが存在するかもしれない
- トラブルシュートの際に、Mistの実装まで含めて確認が必要になるかもしれない
という懸念点・欠点が挙げられる。
HTTP API
Hydro MistのHTTP API に一覧が載っている。
例えば関数一覧を取得する。 1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24$ curl -X GET 'http://10.0.0.209:2004/v2/api/functions' | jq
% Total % Received % Xferd Average Speed Time Time Time Current
Dload Upload Total Spent Left Speed
100 208 100 208 0 0 13256 0 --:--:-- --:--:-- --:--:-- 13866
[
{
"name": "hello-mist-scala",
"execute": {
"samples": {
"type": "MOption",
"args": [
{
"type": "MInt"
}
]
}
},
"path": "hello-mist-scala_0.0.1.jar",
"tags": [],
"className": "HelloMist$",
"defaultContext": "default",
"lang": "scala"
}
]
続いて、当該関数についてジョブ一覧を取得する。 1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27$ curl -X GET 'http://10.0.0.209:2004/v2/api/functions/hello-mist-scala/jobs' | jq
% Total % Received % Xferd Average Speed Time Time Time Current
Dload Upload Total Spent Left Speed
100 1682 100 1682 0 0 29484 0 --:--:-- --:--:-- --:--:-- 30035
[
{
"source": "Http",
"startTime": 1553358512167,
"createTime": 1553358507406,
"context": "default",
"params": {
"filePath": "hello-mist-scala_0.0.1.jar",
"className": "HelloMist$",
"arguments": {
"samples": 9
},
"action": "execute"
},
"endTime": 1553358513162,
"jobResult": 2.2222222222222223,
"status": "finished",
"function": "hello-mist-scala",
"jobId": "b044d08f-9554-4be9-8e22-1d687e58c52e",
"workerId": "default_1cb3b66d-99c3-400b-ac3a-f11d72ab8124_4"
},
(snip)
上記のように、ジョブのメタデータと結果が取得される。 ジョブの情報であれば、直接jobs APIを用いても取得可能。
1 | $ curl -X GET 'http://10.0.0.209:2004/v2/api/jobs' | jq |
またジョブのログを出力できる。
1 | $ curl -X GET 'http://10.0.0.209:2004/v2/api/jobs/a87197c8-1692-48bc-b151-978ea89b058a/logs' | head -n 20 |
Reactie API
Hydro MistのReactive API を見ると、MQTTやKafkaと連携して動くAPIがあるようだが、 まだドキュメントが成熟していない。 デフォルトでは無効になっている。
EMRとの連係
Hydro
MistのEMR連係
を眺めるとhello_mistプロジェクト等でEMRとの連係の仕方を示してくれているようだが、まだ情報が足りない。