private class LogDirFailureHandler(name: String, haltBrokerOnDirFailure: Boolean) extends ShutdownableThread(name) { Override def doWork() { val newOfflineLogDir = logDirFailureChannel.takeNextOfflineLogDir() if (haltBrokerOnDirFailure) { fatal(s"Halting broker because dir $newOfflineLogDir is offline") Exit.halt(1) } handleLogDirFailure(newOfflineLogDir) } }
isr_change_notification
ISRに変化があったことを確認する。
kafka/server/ReplicaManager.scala:269
1 2 3 4 5 6 7 8 9 10 11 12
def maybePropagateIsrChanges() { val now = System.currentTimeMillis() isrChangeSet synchronized { if (isrChangeSet.nonEmpty && (lastIsrChangeMs.get() + ReplicaManager.IsrChangePropagationBlackOut < now || lastIsrPropagationMs.get() + ReplicaManager.IsrChangePropagationInterval < now)) { zkClient.propagateIsrChanges(isrChangeSet) isrChangeSet.clear() lastIsrPropagationMs.set(now) } } }
brokers
以下のように、ブローカに関するいくつかの情報を保持する。
1 2
ls /kafka/brokers [seqid, topics, ids]
例えば、ブローカ情報を記録するのは以下の通り。
kafka/zk/KafkaZkClient.scala:95
1 2 3 4 5 6
def registerBroker(brokerInfo: BrokerInfo): Long = { val path = brokerInfo.path val stat = checkedEphemeralCreate(path, brokerInfo.toJsonBytes) info(s"Registered broker ${brokerInfo.broker.id} at path $path with addresses: ${brokerInfo.broker.endPoints}, czxid (broker epoch): ${stat.getCzxid}") stat.getCzxid }
例えば、トピック・パーティション情報は以下の通り。
1 2
get /kafka/brokers/topics/topic/partitions/0/state {"controller_epoch":1,"leader":1001,"version":1,"leader_epoch":0,"isr":[1001]}
controller
例えば、コントローラ情報は以下の通り。
1 2
get /kafka/controller {"version":1,"brokerid":1001,"timestamp":"1551794212551"}
def feature_sum(xs): return [str(sum(x)) for x in xs]
推論
1 2 3 4 5 6 7 8 9
while True: if batch_size > 1: predict( clipper_conn.get_query_addr(), [list(np.random.random(200)) for i in range(batch_size)], batch=True) else: predict(clipper_conn.get_query_addr(), np.random.random(200)) time.sleep(0.2)
import logging, xgboost as xgb, numpy as np from clipper_admin import ClipperConnection, DockerContainerManager cl = ClipperConnection(DockerContainerManager()) cl.start_clipper()
結果
1 2
19-02-27:22:49:09 INFO [docker_container_manager.py:119] Starting managed Redis instance in Docker 19-02-27:22:49:14 INFO [clipper_admin.py:126] Clipper is running
19-02-27:22:49:55 INFO [clipper_admin.py:201] Application xgboost-test was successfully registered
ファンクション定義
1 2 3 4 5 6 7 8 9 10 11 12
def get_test_point(): return [np.random.randint(255) for _ in range(784)]
# Create a training matrix. dtrain = xgb.DMatrix(get_test_point(), label=[0])
# We then create parameters, watchlist, and specify the number of rounds # This is code that we use to build our XGBoost Model, and your code may differ. param = {'max_depth': 2, 'eta': 1, 'silent': 1, 'objective': 'binary:logistic'} watchlist = [(dtrain, 'train')] num_round = 2 bst = xgb.train(param, dtrain, num_round, watchlist)
from clipper_admin.deployers import python as python_deployer # We specify which packages to install in the pkgs_to_install arg. # For example, if we wanted to install xgboost and psycopg2, we would use # pkgs_to_install = ['xgboost', 'psycopg2'] python_deployer.deploy_python_closure(cl, name='xgboost-model', version=1, input_type="integers", func=predict, pkgs_to_install=['xgboost'])
なお、ここではコンテナをビルドするときに、xgboostをインストールするように指定している。
結果
1 2 3 4 5 6 7 8 9
19-02-27:22:54:35 INFO [deployer_utils.py:44] Saving function to /tmp/clipper/tmpincj4sg2 19-02-27:22:54:35 INFO [deployer_utils.py:54] Serialized and supplied predict function 19-02-27:22:54:35 INFO [python.py:192] Python closure saved
(snip)
19-02-27:22:54:53 INFO [docker_container_manager.py:257] Found 0 replicas for xgboost-model:1. Adding 1 19-02-27:22:55:00 INFO [clipper_admin.py:635] Successfully registered model xgboost-model:1 19-02-27:22:55:00 INFO [clipper_admin.py:553] Done deploying model xgboost-model:1.
from clipper_admin import ClipperConnection, DockerContainerManager clipper_conn = ClipperConnection(DockerContainerManager()) clipper_conn.start_clipper()
結果の例
1 2
19-02-26:22:24:42 INFO [docker_container_manager.py:119] Starting managed Redis instance in Docker 19-02-26:22:26:50 INFO [clipper_admin.py:126] Clipper is running
def feature_sum(xs): return [str(sum(x)) for x in xs]
デプロイ。
1 2
from clipper_admin.deployers import python as python_deployer python_deployer.deploy_python_closure(clipper_conn, name="sum-model", version=1, input_type="doubles", func=feature_sum)
ここから、モデル抽象化層のDockerイメージが作られる。
1 2 3 4 5 6 7 8 9 10 11
19-02-26:22:30:59 INFO [deployer_utils.py:44] Saving function to /tmp/clipper/tmp67eliqhx 19-02-26:22:30:59 INFO [deployer_utils.py:54] Serialized and supplied predict function 19-02-26:22:30:59 INFO [python.py:192] Python closure saved 19-02-26:22:30:59 INFO [python.py:206] Using Python 3.6 base image
(snip)
19-02-26:22:31:26 INFO [docker_container_manager.py:257] Found 0 replicas for sum-model:1. Adding 1 19-02-26:22:31:33 INFO [clipper_admin.py:635] Successfully registered model sum-model:1 19-02-26:22:31:33 INFO [clipper_admin.py:553] Done deploying model sum-model:1. 19-02-26:22:30:59 INFO [clipper_admin.py:452] Building model Docker image with model data from /tmp/clipper/tmp67eliqhx
File "/home/******/.conda/envs/studioml/lib/python3.7/site-packages/flask/app.py", line 1813, in full_dispatch_request rv = self.dispatch_request() File "/home/******/.conda/envs/studioml/lib/python3.7/site-packages/flask/app.py", line 1799, in dispatch_request return self.view_functions[rule.endpoint](**req.view_args) File "/home/******/.conda/envs/studioml/lib/python3.7/site-packages/studio/apiserver.py", line 37, in dashboard return _render('dashboard.html') File "/home/******/.conda/envs/studioml/lib/python3.7/site-packages/studio/apiserver.py", line 518, in _render auth = get_auth(get_auth_config()) File "/home/******/.conda/envs/studioml/lib/python3.7/site-packages/studio/apiserver.py", line 511, in get_auth_config return get_config()['server']['authentication']
requests.exceptions.ConnectionError: HTTPSConnectionPool(host='zoo.studio.ml', port=443): Max retries exceeded with url: /api/get_user_experiments (Caused by NewConnectionError('<urllib3.connection.VerifiedHTTPSConnection object at 0x7ff5a64a3da0>: Failed to establisha new connection: [Errno -2] Name or service not known'))