メインモジュールの作成
RestVerticleで変換されたイベントを受信し、DBへの登録や更新などを行います。今回は単純な認証処理風の実装をしてみます。
実装の前にh2dbのjarファイルをコンパイル用のクラスパスに追加する必要があります。build.gradleのdependenciesに一行追加します。忘れないうちに1.8対応の設定も追加しておきます。
compile "com.h2database:h2:1.4.187" } sourceCompatibility = 1.8 targetCompatibility = 1.8
今度はBusModBaseを継承し、ExampleBusクラスを作成します。
package com.dmm.vertx.example;
import org.vertx.java.busmods.BusModBase;
public class ExampleBus extends BusModBase {
}
startメソッドを実装します。RestVerticleから発生するイベントのハンドラーを登録し、mod-jdbc-persisterを使用したトランザクション処理を実行します。このときにH2DBのWebコンソールの起動を行い、テーブルの作成も行います。
private long timeout; // JDBCタイムアウト値
private String jdbcAddress; // mod-jdbc-persistorアドレス
private static Server h2Server; // H2DB WebConsole
@Override
public void start(Future<Void> startedResult) {
super.start(startedResult);
this.timeout = this.getOptionalLongConfig("timeout", 15 * 1000);
this.jdbcAddress = this.getOptionalStringConfig("address", "jdbc");
this.eb.registerHandler("http.get", this::show).registerHandler("http.post", this::register)
.registerHandler("http.put", this::modify).registerHandler("http.delete", this::unregister)
.registerHandler("example.ready", (Message<JsonObject> event) -> {
if (event.body().getString("status").equals("ok")) { // RESTインターフェースのデプロイ
this.container.deployVerticle("com.dmm.vertx.example.RestVerticle");
} else {
this.container.exit();
}
});
try {
startH2Console();
this.createTable(); // テーブル作成
} catch (SQLException e) {
startedResult.setFailure(e);
}
}
private void createTable() {
JsonObject query = makeQuery("execute", "CREATE TABLE IF NOT EXISTS EXAMPLE_TBL "
+ "("
+ " ID VARCHAR(10) NOT NULL,"
+ " PASSWORD VARCHAR(100) NOT NULL,"
+ " PRIMARY KEY (ID)"
+ ")");
this.eb.send(jdbcAddress, query, (Message<JsonObject> event) -> {
if (event.body().getString("status").equals("ok")) {
this.eb.send("example.ready", new JsonObject("{\"status\":\"ok\"}"));
} else {
this.eb.send("example.ready", new JsonObject("{\"status\":\"error\"}"));
}
});
}
private JsonObject makeQuery(String action) {
return this.makeQuery(action, null, (Object[])null);
}
private JsonObject makeQuery(String action, String stmt) {
return this.makeQuery(action, stmt, (Object[])null);
}
private JsonObject makeQuery(String action, String stmt, Object... values) {
JsonObject query = new JsonObject().putString("action", action);
if (stmt != null) query.putString("stmt", stmt);
if (values != null) query.putArray("values", new JsonArray(values));
return query;
}
H2DBのWebコンソールに関わる部分はstatic、かつsynchronizedなメソッドとして定義します。Vert.xではこういった実装は不要なはずと思われるかもしれませんが、モジュールのデプロイ時にインスタンスの数を複数にした場合の対策として、こういった実装も必要な場合があります。
今回の場合、このH2DBのWebコンソールはVert.xとは依存性のない機能で、その名のとおり独立したWebサーバーとして起動します。このWebサーバーインスタンスをExampleBusのインスタンス変数として定義し、それぞれのstartメソッドで起動した場合、当然、後に起動されたインスタンスでポートが重複してエラーとなってしまいます。インスタンス数を増やす場合、対象Verticleのモジュールクラスローダーは同一になるため、staticな変数も利用することができます。
RestVerticleで使用したHttpServerの場合はVert.xのコンテナー側で管理されるため、複数インスタンスで起動しても問題ありません。
private synchronized static void startH2Console() throws SQLException {
if (h2Server == null) {
h2Server = Server.createWebServer("-webPort", "8082", "-webAllowOthers", "-webDaemon");
h2Server.start(); // H2DB Webコンソールの起動
}
}
private synchronized static void stopH2Console() {
if (h2Server != null) {
h2Server.stop();
h2Server = null;
}
}
引き続き各処理の実装として、まずはGETに対応した参照処理を実装します。リクエストイベントからパラメータを取得し、クエリーオブジェクトを作成した後、mod-jdbc-persistor向けにイベントを送信しています。そのイベントへの応答を受け、リクエストイベント発生元のRestVerticleへ応答を返しています。
private void show(final Message<JsonObject> event) {
JsonObject params = event.body();
JsonObject query = makeQuery("select", "SELECT * FROM EXAMPLE_TBL WHERE ID = ?", params.getString("id"));
this.eb.sendWithTimeout(jdbcAddress, query, this.timeout, (AsyncResult<Message<JsonObject>> result) -> {
Message<JsonObject> record = result.result();
if (result.succeeded() && record.body().getString("status").equals("ok")
&& record.body().getArray("result").size() == 1) {
sendOK(event, record.body());
} else {
sendStatus("ng", event, record != null ? record.body() : new JsonObject());
}
});
}
次からは更新系の処理を作成します。まずはPOSTに対応したregisterメソッドを実装します。
private void register(final Message<JsonObject> event) {
JsonObject params = event.body();
JsonObject query = makeQuery("insert", "INSERT INTO EXAMPLE_TBL (ID, PASSWORD) VALUES (?, ?)",
params.getString("id"), params.getString("password"));
this.eb.sendWithTimeout(jdbcAddress, query, this.timeout, (AsyncResult<Message<JsonObject>> result) -> {
Message<JsonObject> record = result.result();
if (result.succeeded() && record.body().getString("status").equals("ok")) {
sendOK(event, record.body());
} else {
sendStatus("ng", event, record != null ? record.body() : new JsonObject());
}
});
}
次にPUTに対応したmodifyメソッドを作成します。この場合はトランザクション処理が必要となるため、ネストが深くなってしまいます。
private void modify(final Message<JsonObject> event) {
JsonObject params = event.body();
// トランザクション開始
eb.sendWithTimeout(jdbcAddress, makeQuery("transaction"), this.timeout,
(AsyncResult<Message<JsonObject>> begin) -> {
Message<JsonObject> transaction = begin.result();
if (begin.succeeded() && transaction.body().getString("status").equals("ok")) {
// レコードロック
JsonObject lockQuery = makeQuery("select",
"SELECT * FROM EXAMPLE_TBL WHERE ID = ? FOR UPDATE", params.getString("id"));
transaction.replyWithTimeout(lockQuery, this.timeout, (AsyncResult<Message<JsonObject>> lock) -> {
Message<JsonObject> lockRecord = lock.result();
JsonObject lockResult = lockRecord.body();
if (lock.succeeded() && lockResult.getString("status").equals("ok")
&& lockResult.getArray("result").size() == 1) {
// レコード更新
JsonObject updateQuery = makeQuery("update",
"UPDATE EXAMPLE_TBL SET PASSWORD = ? WHERE ID = ? AND PASSWORD = ?",
params.getString("newpass")), params.getString("id"),
params.getString("password")));
lockRecord.replyWithTimeout(updateQuery,this.timeout,(AsyncResult<Message<JsonObject>> update) -> {
Message<JsonObject> updatedRecord = update.result();
JsonObject updateResult = updatedRecord.body();
if (update.succeeded() && updateResult.getString("status").equals("ok")
&& updateResult.getInteger("updated") == 1) {
// コミット
updatedRecord.replyWithTimeout(makeQuery("commit"), this.timeout,
(AsyncResult<Message<JsonObject>> commit) -> {
if (commit.succeeded()) sendOK(event);
else sendStatus("ng", event);
});
} else { // ロールバック
updatedRecord.reply(new JsonObject("{\"action\":\"rollback\"}"));
sendStatus("ng", event);
}
});
} else { // ロールバック
lockRecord.reply(new JsonObject("{\"action\":\"rollback\"}"));
sendStatus("ng", event);
}
});
} else { // トランザクション開始失敗
sendError(event, "transaction failure.");
}
});
}
最後にDELETEに対応したunregisterメソッドを実装します。
private void unregister(final Message<JsonObject> event) {
JsonObject params = event.body();
JsonObject query = makeQuery("update", "DELETE EXAMPLE_TBL WHERE ID = ?", params.getString("id"));
this.eb.sendWithTimeout(jdbcAddress, query, this.timeout, (AsyncResult<Message<JsonObject>> result) -> {
Message<JsonObject> record = result.result();
if (result.succeeded() && record.body().getString("status").equals("ok")) {
sendOK(event, record.body());
} else {
sendStatus("ng", event, record != null ? record.body() : new JsonObject());
}
});
}
