From ef08412b4fd8796ea880fe2e7e9bbff2d22fda2a Mon Sep 17 00:00:00 2001 From: doxlik Date: Fri, 20 Mar 2026 02:16:54 +0400 Subject: [PATCH 1/4] possible Vertx optimization by using cached tuples --- .../Java/vertx/src/main/java/vertx/App.java | 959 +++++++++--------- 1 file changed, 489 insertions(+), 470 deletions(-) diff --git a/frameworks/Java/vertx/src/main/java/vertx/App.java b/frameworks/Java/vertx/src/main/java/vertx/App.java index 039e2ee905a..eb2b7192405 100755 --- a/frameworks/Java/vertx/src/main/java/vertx/App.java +++ b/frameworks/Java/vertx/src/main/java/vertx/App.java @@ -38,509 +38,528 @@ public class App extends VerticleBase implements Handler { - private static final int NUM_PROCESSORS = Runtime.getRuntime().availableProcessors(); - private static final org.slf4j.Logger log = org.slf4j.LoggerFactory.getLogger(App.class); - - /** - * Returns the value of the "queries" getRequest parameter, which is an integer - * bound between 1 and 500 with a default value of 1. - * - * @param request the current HTTP request - * @return the value of the "queries" parameter - */ - static int getQueries(HttpServerRequest request) { - String param = request.getParam("queries"); - - if (param == null) { - return 1; + private static final int NUM_PROCESSORS = Runtime.getRuntime().availableProcessors(); + private static final org.slf4j.Logger log = org.slf4j.LoggerFactory.getLogger(App.class); + + /** + * Returns the value of the "queries" getRequest parameter, which is an integer + * bound between 1 and 500 with a default value of 1. + * + * @param request the current HTTP request + * @return the value of the "queries" parameter + */ + static int getQueries(HttpServerRequest request) { + String param = request.getParam("queries"); + + if (param == null) { + return 1; + } + try { + int parsedValue = Integer.parseInt(param); + return Math.min(500, Math.max(1, parsedValue)); + } catch (NumberFormatException e) { + return 1; + } } - try { - int parsedValue = Integer.parseInt(param); - return Math.min(500, Math.max(1, parsedValue)); - } catch (NumberFormatException e) { - return 1; + + private static Logger logger = LoggerFactory.getLogger(App.class.getName()); + + private static final Integer[] BOXED_RND = IntStream.range(1, 10001).boxed().toArray(Integer[]::new); + + private static final String PATH_PLAINTEXT = "/plaintext"; + private static final String PATH_JSON = "/json"; + private static final String PATH_DB = "/db"; + private static final String PATH_QUERIES = "/queries"; + private static final String PATH_UPDATES = "/updates"; + private static final String PATH_FORTUNES = "/fortunes"; + private static final String PATH_CACHING = "/cached-queries"; + + private static final CharSequence RESPONSE_TYPE_PLAIN = HttpHeaders.createOptimized("text/plain"); + private static final CharSequence RESPONSE_TYPE_HTML = HttpHeaders.createOptimized("text/html; charset=UTF-8"); + private static final CharSequence RESPONSE_TYPE_JSON = HttpHeaders.createOptimized("application/json"); + + private static final String HELLO_WORLD = "Hello, world!"; + private static final Buffer HELLO_WORLD_BUFFER = Buffer.buffer(HELLO_WORLD, "UTF-8"); + + private static final CharSequence HEADER_SERVER = HttpHeaders.SERVER; + private static final CharSequence HEADER_DATE = HttpHeaders.DATE; + private static final CharSequence HEADER_CONTENT_TYPE = HttpHeaders.CONTENT_TYPE; + private static final CharSequence HEADER_CONTENT_LENGTH = HttpHeaders.CONTENT_LENGTH; + + private static final CharSequence SERVER = HttpHeaders.createOptimized("vert.x"); + + private static final String SELECT_WORLD = "SELECT id, randomnumber FROM world WHERE id = $1"; + private static final String SELECT_FORTUNE = "SELECT id, message FROM fortune"; + private static final String SELECT_WORLDS = "SELECT id, randomnumber FROM world"; + + private static final Tuple[] tupleCache = new Tuple[10000]; + + public static CharSequence createDateHeader() { + return HttpHeaders.createOptimized(DateTimeFormatter.RFC_1123_DATE_TIME.format(ZonedDateTime.now())); } - } - - private static Logger logger = LoggerFactory.getLogger(App.class.getName()); - - private static final Integer[] BOXED_RND = IntStream.range(1, 10001).boxed().toArray(Integer[]::new); - - private static final String PATH_PLAINTEXT = "/plaintext"; - private static final String PATH_JSON = "/json"; - private static final String PATH_DB = "/db"; - private static final String PATH_QUERIES = "/queries"; - private static final String PATH_UPDATES = "/updates"; - private static final String PATH_FORTUNES = "/fortunes"; - private static final String PATH_CACHING = "/cached-queries"; - - private static final CharSequence RESPONSE_TYPE_PLAIN = HttpHeaders.createOptimized("text/plain"); - private static final CharSequence RESPONSE_TYPE_HTML = HttpHeaders.createOptimized("text/html; charset=UTF-8"); - private static final CharSequence RESPONSE_TYPE_JSON = HttpHeaders.createOptimized("application/json"); - - private static final String HELLO_WORLD = "Hello, world!"; - private static final Buffer HELLO_WORLD_BUFFER = Buffer.buffer(HELLO_WORLD, "UTF-8"); - - private static final CharSequence HEADER_SERVER = HttpHeaders.SERVER; - private static final CharSequence HEADER_DATE = HttpHeaders.DATE; - private static final CharSequence HEADER_CONTENT_TYPE = HttpHeaders.CONTENT_TYPE; - private static final CharSequence HEADER_CONTENT_LENGTH = HttpHeaders.CONTENT_LENGTH; - - private static final CharSequence SERVER = HttpHeaders.createOptimized("vert.x"); - - private static final String SELECT_WORLD = "SELECT id, randomnumber FROM world WHERE id = $1"; - private static final String SELECT_FORTUNE = "SELECT id, message FROM fortune"; - private static final String SELECT_WORLDS = "SELECT id, randomnumber FROM world"; - - public static CharSequence createDateHeader() { - return HttpHeaders.createOptimized(DateTimeFormatter.RFC_1123_DATE_TIME.format(ZonedDateTime.now())); - } - - /** - * Returns a random integer that is a suitable value for both the {@code id} - * and {@code randomNumber} properties of a world object. - * - * @return a random world number - */ - static Integer boxedRandomWorldNumber() { - final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); - final var boxedRnd = BOXED_RND[rndValue - 1]; - assert boxedRnd.intValue() == rndValue; - return boxedRnd; - } - - private HttpServer server; - private SqlClientInternal client; - private CharSequence dateString; - private MultiMap plaintextHeaders; - private MultiMap jsonHeaders; - - private final RockerOutputFactory factory = BufferRockerOutput.factory(ContentType.RAW); - - private Throwable databaseErr; - private PreparedQuery> SELECT_WORLD_QUERY; - private PreparedQuery>> SELECT_FORTUNE_QUERY; - @SuppressWarnings("unchecked") - private PreparedQuery>[] AGGREGATED_UPDATE_WORLD_QUERY = new PreparedQuery[500]; - private WorldCache WORLD_CACHE; - - private MultiMap plaintextHeaders() { - return HttpHeaders - .headers() - .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_PLAIN) - .add(HEADER_SERVER, SERVER) - .add(HEADER_DATE, dateString) - .copy(false); - } - - private MultiMap jsonHeaders() { - return HttpHeaders - .headers() - .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_JSON) - .add(HEADER_SERVER, SERVER) - .add(HEADER_DATE, dateString) - .copy(false); - } - - @Override - public Future start() throws Exception { - int port = 8080; - server = vertx - .createHttpServer(new HttpServerOptions() - .setHttp2ClearTextEnabled(false) - .setStrictThreadMode(true)) - .requestHandler(App.this); - dateString = createDateHeader(); - plaintextHeaders = plaintextHeaders(); - jsonHeaders = jsonHeaders(); - JsonObject config = config(); - vertx.setPeriodic(1000, id -> { - dateString = createDateHeader(); - plaintextHeaders = plaintextHeaders(); - jsonHeaders = jsonHeaders(); - }); - PgConnectOptions options = new PgConnectOptions(); - options.setDatabase(config.getString("database", "hello_world")); - options.setHost(config.getString("host", "tfb-database")); - options.setPort(config.getInteger("port", 5432)); - options.setUser(config.getString("username", "benchmarkdbuser")); - options.setPassword(config.getString("password", "benchmarkdbpass")); - options.setCachePreparedStatements(true); - options.setPreparedStatementCacheMaxSize(1024); - options.setPipeliningLimit(256); // Large pipelining means less flushing and we use a single connection anyway - Future clientsInit = initClients(options); - return clientsInit - .transform(ar -> { - databaseErr = ar.cause(); - return server.listen(port); - }); - } - - private Future initClients(PgConnectOptions options) { - return PgConnection.connect(vertx, options) - .flatMap(conn -> { - client = (SqlClientInternal) conn; - List> list = new ArrayList<>(); - Future f1 = conn.prepare(SELECT_WORLD) - .andThen(onSuccess(ps -> SELECT_WORLD_QUERY = ps.query())); - list.add(f1); - Future f2 = conn.prepare(SELECT_FORTUNE) - .andThen(onSuccess(ps -> { - SELECT_FORTUNE_QUERY = ps.query(). - collecting(Collectors.mapping(row -> new Fortune(row.getInteger(0), row.getString(1)), Collectors.toList())); - })); - list.add(f2); - Future f3 = conn.preparedQuery(SELECT_WORLDS) - .collecting(Collectors.mapping(row -> new CachedWorld(row.getInteger(0), row.getInteger(1)), Collectors.toList())) - .execute() - .map(worlds -> new WorldCache(worlds.value())) - .andThen(onSuccess(wc -> WORLD_CACHE = wc)); - list.add(f3); - for (int i = 0; i < AGGREGATED_UPDATE_WORLD_QUERY.length; i++) { - int idx = i; - Future fut = conn - .prepare(buildAggregatedUpdateQuery(1 + idx)) - .andThen(onSuccess(ps -> AGGREGATED_UPDATE_WORLD_QUERY[idx] = ps.query())); - list.add(fut); - } - return Future.join(list); - }); - } - - private static String buildAggregatedUpdateQuery(int len) { - StringBuilder sql = new StringBuilder(); - sql.append("UPDATE world SET randomnumber = CASE ID"); - for (int i = 0; i < len; i++) { - int offset = (i * 2) + 1; - sql.append(" WHEN $").append(offset).append(" THEN $").append(offset + 1); + + /** + * Returns a random integer that is a suitable value for both the {@code id} + * and {@code randomNumber} properties of a world object. + * + * @return a random world number + */ + static Integer boxedRandomWorldNumber() { + final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); + final var boxedRnd = BOXED_RND[rndValue - 1]; + assert boxedRnd.intValue() == rndValue; + return boxedRnd; } - sql.append(" ELSE randomnumber"); - sql.append(" END WHERE ID IN ($1"); - for (int i = 1; i < len; i++) { - int offset = (i * 2) + 1; - sql.append(",$").append(offset); + + static Integer primitiveRandomWorldNumber() { + final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); + return rndValue; } - sql.append(")"); - return sql.toString(); - } - - public static Handler> onSuccess(Handler handler) { - return ar -> { - if (ar.succeeded()) { - handler.handle(ar.result()); - } - }; - } - - @Override - public void handle(HttpServerRequest request) { - try { - switch (request.path()) { - case PATH_PLAINTEXT: - handlePlainText(request); - break; - case PATH_JSON: - handleJson(request); - break; - case PATH_DB: - handleDb(request); - break; - case PATH_QUERIES: - new Queries(request).handle(); - break; - case PATH_UPDATES: - new Update(request).handle(); - break; - case PATH_FORTUNES: - handleFortunes(request); - break; - case PATH_CACHING: - handleCaching(request); - break; - default: - request.response() - .setStatusCode(404) - .end(); - break; - } - } catch (Exception e) { - sendError(request, e); + + static Tuple getRandomTuple() { + final int rndValue = primitiveRandomWorldNumber(); + final Tuple tuple = tupleCache[rndValue - 1]; + return tuple; } - } - - @Override - public Future stop() throws Exception { - return server != null ? server.close() : super.stop(); - } - - private void sendError(HttpServerRequest req, Throwable cause) { - logger.error(cause.getMessage(), cause); - req.response().setStatusCode(500).end(); - } - - private void handlePlainText(HttpServerRequest request) { - HttpServerResponse response = request.response(); - response.headers().setAll(plaintextHeaders); - response.end(HELLO_WORLD_BUFFER); - } - - private void handleJson(HttpServerRequest request) { - HttpServerResponse response = request.response(); - response.headers().setAll(jsonHeaders); - response.end(new Message("Hello, World!").toJson()); - } - - private void handleDb(HttpServerRequest req) { - HttpServerResponse resp = req.response(); - SELECT_WORLD_QUERY.execute(Tuple.of(boxedRandomWorldNumber())).onComplete(res -> { - if (res.succeeded()) { - RowIterator resultSet = res.result().iterator(); - if (!resultSet.hasNext()) { - resp.setStatusCode(404).end(); - return; + + + private HttpServer server; + private SqlClientInternal client; + private CharSequence dateString; + private MultiMap plaintextHeaders; + private MultiMap jsonHeaders; + + private final RockerOutputFactory factory = BufferRockerOutput.factory(ContentType.RAW); + + private Throwable databaseErr; + private PreparedQuery> SELECT_WORLD_QUERY; + private PreparedQuery>> SELECT_FORTUNE_QUERY; + @SuppressWarnings("unchecked") + private PreparedQuery>[] AGGREGATED_UPDATE_WORLD_QUERY = new PreparedQuery[500]; + private WorldCache WORLD_CACHE; + + private MultiMap plaintextHeaders() { + return HttpHeaders + .headers() + .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_PLAIN) + .add(HEADER_SERVER, SERVER) + .add(HEADER_DATE, dateString) + .copy(false); + } + + private MultiMap jsonHeaders() { + return HttpHeaders + .headers() + .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_JSON) + .add(HEADER_SERVER, SERVER) + .add(HEADER_DATE, dateString) + .copy(false); + } + + @Override + public Future start() throws Exception { + int port = 8080; + server = vertx + .createHttpServer(new HttpServerOptions() + .setHttp2ClearTextEnabled(false) + .setStrictThreadMode(true)) + .requestHandler(App.this); + dateString = createDateHeader(); + plaintextHeaders = plaintextHeaders(); + jsonHeaders = jsonHeaders(); + JsonObject config = config(); + vertx.setPeriodic(1000, id -> { + dateString = createDateHeader(); + plaintextHeaders = plaintextHeaders(); + jsonHeaders = jsonHeaders(); + }); + + for (int i = 0; i < 10000; i++) { + tupleCache[i] = Tuple.of(i + 1); } - Row row = resultSet.next(); - World world = new World(row.getInteger(0), row.getInteger(1)); - MultiMap headers = resp.headers(); - headers.add(HttpHeaders.SERVER, SERVER); - headers.add(HttpHeaders.DATE, dateString); - headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); - resp.end(world.toJson()); - } else { - sendError(req, res.cause()); - } - }); - } - - class Queries implements Handler>> { - - boolean failed; - final World[] worlds; - final HttpServerRequest req; - final HttpServerResponse resp; - final int queries; - int worldsIndex; - - public Queries(HttpServerRequest req) { - int queries = getQueries(req); - - this.req = req; - this.resp = req.response(); - this.queries = queries; - this.worlds = new World[queries]; - this.worldsIndex = 0; + + PgConnectOptions options = new PgConnectOptions(); + options.setDatabase(config.getString("database", "hello_world")); + options.setHost(config.getString("host", "tfb-database")); + options.setPort(config.getInteger("port", 5432)); + options.setUser(config.getString("username", "benchmarkdbuser")); + options.setPassword(config.getString("password", "benchmarkdbpass")); + options.setCachePreparedStatements(true); + options.setPreparedStatementCacheMaxSize(1024); + options.setPipeliningLimit(256); // Large pipelining means less flushing and we use a single connection anyway + Future clientsInit = initClients(options); + return clientsInit + .transform(ar -> { + databaseErr = ar.cause(); + return server.listen(port); + }); + } + + private Future initClients(PgConnectOptions options) { + return PgConnection.connect(vertx, options) + .flatMap(conn -> { + client = (SqlClientInternal) conn; + List> list = new ArrayList<>(); + Future f1 = conn.prepare(SELECT_WORLD) + .andThen(onSuccess(ps -> SELECT_WORLD_QUERY = ps.query())); + list.add(f1); + Future f2 = conn.prepare(SELECT_FORTUNE) + .andThen(onSuccess(ps -> { + SELECT_FORTUNE_QUERY = ps.query(). + collecting(Collectors.mapping(row -> new Fortune(row.getInteger(0), row.getString(1)), Collectors.toList())); + })); + list.add(f2); + Future f3 = conn.preparedQuery(SELECT_WORLDS) + .collecting(Collectors.mapping(row -> new CachedWorld(row.getInteger(0), row.getInteger(1)), Collectors.toList())) + .execute() + .map(worlds -> new WorldCache(worlds.value())) + .andThen(onSuccess(wc -> WORLD_CACHE = wc)); + list.add(f3); + for (int i = 0; i < AGGREGATED_UPDATE_WORLD_QUERY.length; i++) { + int idx = i; + Future fut = conn + .prepare(buildAggregatedUpdateQuery(1 + idx)) + .andThen(onSuccess(ps -> AGGREGATED_UPDATE_WORLD_QUERY[idx] = ps.query())); + list.add(fut); + } + return Future.join(list); + }); } - private void handle() { - client.group(/*queries, */c -> { - for (int i = 0; i < queries; i++) { - c.preparedQuery(SELECT_WORLD) - .execute(Tuple.of(boxedRandomWorldNumber())) - .onComplete(this); + private static String buildAggregatedUpdateQuery(int len) { + StringBuilder sql = new StringBuilder(); + sql.append("UPDATE world SET randomnumber = CASE ID"); + for (int i = 0; i < len; i++) { + int offset = (i * 2) + 1; + sql.append(" WHEN $").append(offset).append(" THEN $").append(offset + 1); } - }); + sql.append(" ELSE randomnumber"); + sql.append(" END WHERE ID IN ($1"); + for (int i = 1; i < len; i++) { + int offset = (i * 2) + 1; + sql.append(",$").append(offset); + } + sql.append(")"); + return sql.toString(); + } + + public static Handler> onSuccess(Handler handler) { + return ar -> { + if (ar.succeeded()) { + handler.handle(ar.result()); + } + }; } @Override - public void handle(AsyncResult> ar) { - if (!failed) { - if (ar.failed()) { - failed = true; - sendError(req, ar.cause()); - return; + public void handle(HttpServerRequest request) { + try { + switch (request.path()) { + case PATH_PLAINTEXT: + handlePlainText(request); + break; + case PATH_JSON: + handleJson(request); + break; + case PATH_DB: + handleDb(request); + break; + case PATH_QUERIES: + new Queries(request).handle(); + break; + case PATH_UPDATES: + new Update(request).handle(); + break; + case PATH_FORTUNES: + handleFortunes(request); + break; + case PATH_CACHING: + handleCaching(request); + break; + default: + request.response() + .setStatusCode(404) + .end(); + break; + } + } catch (Exception e) { + sendError(request, e); } + } - // we need a final reference - final Tuple row = ar.result().iterator().next(); - worlds[worldsIndex++] = new World(row.getInteger(0), row.getInteger(1)); - - // stop condition - if (worldsIndex == queries) { - MultiMap headers = resp.headers(); - headers.add(HttpHeaders.SERVER, SERVER); - headers.add(HttpHeaders.DATE, dateString); - headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); - resp.end(World.toJson(worlds)); - } - } + @Override + public Future stop() throws Exception { + return server != null ? server.close() : super.stop(); } - } - private class Update { + private void sendError(HttpServerRequest req, Throwable cause) { + logger.error(cause.getMessage(), cause); + req.response().setStatusCode(500).end(); + } - private final HttpServerRequest request; - private final World[] worldsToUpdate; - private boolean failed; - private int selectWorldCompletedCount; + private void handlePlainText(HttpServerRequest request) { + HttpServerResponse response = request.response(); + response.headers().setAll(plaintextHeaders); + response.end(HELLO_WORLD_BUFFER); + } - public Update(HttpServerRequest request) { - this.request = request; - this.worldsToUpdate = new World[getQueries(request)]; + private void handleJson(HttpServerRequest request) { + HttpServerResponse response = request.response(); + response.headers().setAll(jsonHeaders); + response.end(new Message("Hello, World!").toJson()); + } + + private void handleDb(HttpServerRequest req) { + HttpServerResponse resp = req.response(); + SELECT_WORLD_QUERY.execute(getRandomTuple()).onComplete(res -> { + if (res.succeeded()) { + RowIterator resultSet = res.result().iterator(); + if (!resultSet.hasNext()) { + resp.setStatusCode(404).end(); + return; + } + Row row = resultSet.next(); + World world = new World(row.getInteger(0), row.getInteger(1)); + MultiMap headers = resp.headers(); + headers.add(HttpHeaders.SERVER, SERVER); + headers.add(HttpHeaders.DATE, dateString); + headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); + resp.end(world.toJson()); + } else { + sendError(req, res.cause()); + } + }); } - public void handle() { - client.group(/*worldsToUpdate.length, */c -> { - final PreparedQuery> preparedQuery = c.preparedQuery(App.SELECT_WORLD); - for (int i = 0; i < worldsToUpdate.length; i++) { - final Integer id = boxedRandomWorldNumber(); - final int index = i; - preparedQuery.execute(Tuple.of(id)).onComplete(res -> { + class Queries implements Handler>> { + + boolean failed; + final World[] worlds; + final HttpServerRequest req; + final HttpServerResponse resp; + final int queries; + int worldsIndex; + + public Queries(HttpServerRequest req) { + int queries = getQueries(req); + + this.req = req; + this.resp = req.response(); + this.queries = queries; + this.worlds = new World[queries]; + this.worldsIndex = 0; + } + + private void handle() { + client.group(/*queries, */c -> { + for (int i = 0; i < queries; i++) { + c.preparedQuery(SELECT_WORLD) + .execute(getRandomTuple()) + .onComplete(this); + } + }); + } + + @Override + public void handle(AsyncResult> ar) { if (!failed) { - if (res.failed()) { - failed = true; - sendError(request, res.cause()); - return; - } - worldsToUpdate[index] = new World(res.result().iterator().next().getInteger(0), boxedRandomWorldNumber()); - if (++selectWorldCompletedCount == worldsToUpdate.length) { - randomWorldsQueryCompleted(); - } + if (ar.failed()) { + failed = true; + sendError(req, ar.cause()); + return; + } + + // we need a final reference + final Tuple row = ar.result().iterator().next(); + worlds[worldsIndex++] = new World(row.getInteger(0), row.getInteger(1)); + + // stop condition + if (worldsIndex == queries) { + MultiMap headers = resp.headers(); + headers.add(HttpHeaders.SERVER, SERVER); + headers.add(HttpHeaders.DATE, dateString); + headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); + resp.end(World.toJson(worlds)); + } } - }); } - }); } - private void randomWorldsQueryCompleted() { - Arrays.sort(worldsToUpdate); - final List params = new ArrayList<>(worldsToUpdate.length * 2); - for (int i = 0, count = worldsToUpdate.length;i < count;i++) { - var world = worldsToUpdate[i]; - params.add(world.getId()); - params.add(world.getRandomNumber()); - } - AGGREGATED_UPDATE_WORLD_QUERY[worldsToUpdate.length - 1] - .execute(Tuple.wrap(params)) - .onComplete(updateResult -> { - if (updateResult.failed()) { - sendError(request, updateResult.cause()); - return; + private class Update { + + private final HttpServerRequest request; + private final World[] worldsToUpdate; + private boolean failed; + private int selectWorldCompletedCount; + + public Update(HttpServerRequest request) { + this.request = request; + this.worldsToUpdate = new World[getQueries(request)]; + } + + public void handle() { + client.group(/*worldsToUpdate.length, */c -> { + final PreparedQuery> preparedQuery = c.preparedQuery(App.SELECT_WORLD); + for (int i = 0; i < worldsToUpdate.length; i++) { + final int index = i; + preparedQuery.execute(getRandomTuple()).onComplete(res -> { + if (!failed) { + if (res.failed()) { + failed = true; + sendError(request, res.cause()); + return; + } + worldsToUpdate[index] = new World(res.result().iterator().next().getInteger(0), primitiveRandomWorldNumber()); + if (++selectWorldCompletedCount == worldsToUpdate.length) { + randomWorldsQueryCompleted(); + } + } + }); + } + }); + } + + private void randomWorldsQueryCompleted() { + Arrays.sort(worldsToUpdate); + final List params = new ArrayList<>(worldsToUpdate.length * 2); + for (int i = 0, count = worldsToUpdate.length; i < count; i++) { + var world = worldsToUpdate[i]; + params.add(world.getId()); + params.add(world.getRandomNumber()); + } + AGGREGATED_UPDATE_WORLD_QUERY[worldsToUpdate.length - 1] + .execute(Tuple.wrap(params)) + .onComplete(updateResult -> { + if (updateResult.failed()) { + sendError(request, updateResult.cause()); + return; + } + sendResponse(); + }); + } + + private void sendResponse() { + var res = request.response(); + MultiMap headers = res.headers(); + headers.add(HttpHeaders.SERVER, App.SERVER); + headers.add(HttpHeaders.DATE, dateString); + headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); + Buffer buff = WorldJsonSerializer.toJsonBuffer(worldsToUpdate); + res.end(buff); } - sendResponse(); - }); } - private void sendResponse() { - var res = request.response(); - MultiMap headers = res.headers(); - headers.add(HttpHeaders.SERVER, App.SERVER); - headers.add(HttpHeaders.DATE, dateString); - headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); - Buffer buff = WorldJsonSerializer.toJsonBuffer(worldsToUpdate); - res.end(buff); + private void handleFortunes(HttpServerRequest req) { + SELECT_FORTUNE_QUERY + .execute() + .onComplete(ar -> { + HttpServerResponse response = req.response(); + if (ar.succeeded()) { + SqlResult> result = ar.result(); + if (result.size() == 0) { + response.setStatusCode(404).end("No results"); + return; + } + List fortunes = result.value(); + fortunes.add(new Fortune(0, "Additional fortune added at request time.")); + Collections.sort(fortunes); + MultiMap headers = response.headers(); + headers.add(HttpHeaders.SERVER, SERVER); + headers.add(HttpHeaders.DATE, dateString); + headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_HTML); + FortunesTemplate template = FortunesTemplate.template(fortunes); + response.end(template.render(factory).buffer()); + } else { + sendError(req, ar.cause()); + } + }); } - } - - private void handleFortunes(HttpServerRequest req) { - SELECT_FORTUNE_QUERY - .execute() - .onComplete(ar -> { - HttpServerResponse response = req.response(); - if (ar.succeeded()) { - SqlResult> result = ar.result(); - if (result.size() == 0) { - response.setStatusCode(404).end("No results"); - return; + + private void handleCaching(HttpServerRequest req) { + int count = 1; + try { + String countStr = req.getParam("count"); + if (countStr != null) { + count = Integer.parseInt(countStr); + } + } catch (NumberFormatException ignore) { } - List fortunes = result.value(); - fortunes.add(new Fortune(0, "Additional fortune added at request time.")); - Collections.sort(fortunes); + count = Math.max(1, count); + count = Math.min(500, count); + List worlds = WORLD_CACHE.getCachedWorld(count); + HttpServerResponse response = req.response(); MultiMap headers = response.headers(); - headers.add(HttpHeaders.SERVER, SERVER); - headers.add(HttpHeaders.DATE, dateString); - headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_HTML); - FortunesTemplate template = FortunesTemplate.template(fortunes); - response.end(template.render(factory).buffer()); - } else { - sendError(req, ar.cause()); - } - }); - } - - private void handleCaching(HttpServerRequest req) { - int count = 1; - try { - String countStr = req.getParam("count"); - if (countStr != null) { - count = Integer.parseInt(countStr); - } - } catch (NumberFormatException ignore) { - } - count = Math.max(1, count); - count = Math.min(500, count); - List worlds = WORLD_CACHE.getCachedWorld(count); - HttpServerResponse response = req.response(); - MultiMap headers = response.headers(); - headers - .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_JSON) - .add(HEADER_SERVER, SERVER) - .add(HEADER_DATE, dateString); - response.end(CachedWorld.toJson(worlds)); - } - - public static void main(String[] args) throws Exception { - int eventLoopPoolSize = NUM_PROCESSORS; - String sizeProp = System.getProperty("vertx.eventLoopPoolSize"); - if (sizeProp != null) { - try { - eventLoopPoolSize = Integer.parseInt(sizeProp); - } catch (NumberFormatException e) { - e.printStackTrace(); - } + headers + .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_JSON) + .add(HEADER_SERVER, SERVER) + .add(HEADER_DATE, dateString); + response.end(CachedWorld.toJson(worlds)); } - JsonObject config = new JsonObject(new String(Files.readAllBytes(new File(args[0]).toPath()))); - Vertx vertx = Vertx.vertx(new VertxOptions() - .setEventLoopPoolSize(eventLoopPoolSize) - .setPreferNativeTransport(true) - .setDisableTCCL(true) - ); - vertx.exceptionHandler(err -> { - err.printStackTrace(); - }); - printConfig((VertxInternal) vertx); - vertx.deployVerticle( - App.class.getName(), - new DeploymentOptions().setInstances(eventLoopPoolSize).setConfig(config)) - .onComplete(event -> { - if (event.succeeded()) { - logger.info("Server listening on port " + 8080); - } else { - logger.error("Unable to start your application", event.cause()); - } - }); - } - - private static void printConfig(VertxInternal vertx) { - boolean nativeTransport = vertx.isNativeTransportEnabled(); - String transport = vertx.transport().getClass().getSimpleName(); - String version = "unknown"; - try { - InputStream in = Vertx.class.getClassLoader().getResourceAsStream("META-INF/vertx/vertx-version.txt"); - if (in == null) { - in = Vertx.class.getClassLoader().getResourceAsStream("vertx-version.txt"); - } - ByteArrayOutputStream out = new ByteArrayOutputStream(); - byte[] buffer = new byte[256]; - while (true) { - int amount = in.read(buffer); - if (amount == -1) { - break; + + public static void main(String[] args) throws Exception { + int eventLoopPoolSize = NUM_PROCESSORS; + String sizeProp = System.getProperty("vertx.eventLoopPoolSize"); + if (sizeProp != null) { + try { + eventLoopPoolSize = Integer.parseInt(sizeProp); + } catch (NumberFormatException e) { + e.printStackTrace(); + } } - out.write(buffer, 0, amount); - } - version = out.toString(); - } catch (IOException e) { - logger.error("Could not read Vertx version", e);; + JsonObject config = new JsonObject(new String(Files.readAllBytes(new File(args[0]).toPath()))); + Vertx vertx = Vertx.vertx(new VertxOptions() + .setEventLoopPoolSize(eventLoopPoolSize) + .setPreferNativeTransport(true) + .setDisableTCCL(true) + ); + vertx.exceptionHandler(err -> { + err.printStackTrace(); + }); + printConfig((VertxInternal) vertx); + vertx.deployVerticle( + App.class.getName(), + new DeploymentOptions().setInstances(eventLoopPoolSize).setConfig(config)) + .onComplete(event -> { + if (event.succeeded()) { + logger.info("Server listening on port " + 8080); + } else { + logger.error("Unable to start your application", event.cause()); + } + }); } - logger.info("Vertx: " + version); - logger.info("Processors: " + NUM_PROCESSORS); - logger.info("Event Loop Size: " + ((MultithreadEventExecutorGroup)vertx.nettyEventLoopGroup()).executorCount()); - logger.info("Native transport : " + nativeTransport); - logger.info("Transport : " + transport); - logger.info("Netty buffer bound check : " + System.getProperty("io.netty.buffer.checkBounds")); - logger.info("Netty buffer accessibility check : " + System.getProperty("io.netty.buffer.checkAccessible")); - for (SysProps sysProp : SysProps.values()) { - logger.info(sysProp.name + " : " + sysProp.get()); + + private static void printConfig(VertxInternal vertx) { + boolean nativeTransport = vertx.isNativeTransportEnabled(); + String transport = vertx.transport().getClass().getSimpleName(); + String version = "unknown"; + try { + InputStream in = Vertx.class.getClassLoader().getResourceAsStream("META-INF/vertx/vertx-version.txt"); + if (in == null) { + in = Vertx.class.getClassLoader().getResourceAsStream("vertx-version.txt"); + } + ByteArrayOutputStream out = new ByteArrayOutputStream(); + byte[] buffer = new byte[256]; + while (true) { + int amount = in.read(buffer); + if (amount == -1) { + break; + } + out.write(buffer, 0, amount); + } + version = out.toString(); + } catch (IOException e) { + logger.error("Could not read Vertx version", e); + ; + } + logger.info("Vertx: " + version); + logger.info("Processors: " + NUM_PROCESSORS); + logger.info("Event Loop Size: " + ((MultithreadEventExecutorGroup) vertx.nettyEventLoopGroup()).executorCount()); + logger.info("Native transport : " + nativeTransport); + logger.info("Transport : " + transport); + logger.info("Netty buffer bound check : " + System.getProperty("io.netty.buffer.checkBounds")); + logger.info("Netty buffer accessibility check : " + System.getProperty("io.netty.buffer.checkAccessible")); + for (SysProps sysProp : SysProps.values()) { + logger.info(sysProp.name + " : " + sysProp.get()); + } } - } } From 937373b87e99f324ede35471f2a05a2ba6f89c30 Mon Sep 17 00:00:00 2001 From: doxlik Date: Fri, 20 Mar 2026 02:27:41 +0400 Subject: [PATCH 2/4] possible Vertx optimization by using cached tuples --- frameworks/Java/vertx/src/main/java/vertx/App.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/frameworks/Java/vertx/src/main/java/vertx/App.java b/frameworks/Java/vertx/src/main/java/vertx/App.java index eb2b7192405..c0146235bd1 100755 --- a/frameworks/Java/vertx/src/main/java/vertx/App.java +++ b/frameworks/Java/vertx/src/main/java/vertx/App.java @@ -111,7 +111,7 @@ static Integer boxedRandomWorldNumber() { return boxedRnd; } - static Integer primitiveRandomWorldNumber() { + static int primitiveRandomWorldNumber() { final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); return rndValue; } From 5b0e89d8067b349affeeaddc6c510080462c3f06 Mon Sep 17 00:00:00 2001 From: doxlik Date: Fri, 20 Mar 2026 18:26:28 +0400 Subject: [PATCH 3/4] remove automatic idents --- .../Java/vertx/src/main/java/vertx/App.java | 972 +++++++++--------- 1 file changed, 484 insertions(+), 488 deletions(-) diff --git a/frameworks/Java/vertx/src/main/java/vertx/App.java b/frameworks/Java/vertx/src/main/java/vertx/App.java index c0146235bd1..510e2ff30b9 100755 --- a/frameworks/Java/vertx/src/main/java/vertx/App.java +++ b/frameworks/Java/vertx/src/main/java/vertx/App.java @@ -38,528 +38,524 @@ public class App extends VerticleBase implements Handler { - private static final int NUM_PROCESSORS = Runtime.getRuntime().availableProcessors(); - private static final org.slf4j.Logger log = org.slf4j.LoggerFactory.getLogger(App.class); - - /** - * Returns the value of the "queries" getRequest parameter, which is an integer - * bound between 1 and 500 with a default value of 1. - * - * @param request the current HTTP request - * @return the value of the "queries" parameter - */ - static int getQueries(HttpServerRequest request) { - String param = request.getParam("queries"); - - if (param == null) { - return 1; - } - try { - int parsedValue = Integer.parseInt(param); - return Math.min(500, Math.max(1, parsedValue)); - } catch (NumberFormatException e) { - return 1; - } + private static final int NUM_PROCESSORS = Runtime.getRuntime().availableProcessors(); + private static final org.slf4j.Logger log = org.slf4j.LoggerFactory.getLogger(App.class); + + /** + * Returns the value of the "queries" getRequest parameter, which is an integer + * bound between 1 and 500 with a default value of 1. + * + * @param request the current HTTP request + * @return the value of the "queries" parameter + */ + static int getQueries(HttpServerRequest request) { + String param = request.getParam("queries"); + + if (param == null) { + return 1; } - - private static Logger logger = LoggerFactory.getLogger(App.class.getName()); - - private static final Integer[] BOXED_RND = IntStream.range(1, 10001).boxed().toArray(Integer[]::new); - - private static final String PATH_PLAINTEXT = "/plaintext"; - private static final String PATH_JSON = "/json"; - private static final String PATH_DB = "/db"; - private static final String PATH_QUERIES = "/queries"; - private static final String PATH_UPDATES = "/updates"; - private static final String PATH_FORTUNES = "/fortunes"; - private static final String PATH_CACHING = "/cached-queries"; - - private static final CharSequence RESPONSE_TYPE_PLAIN = HttpHeaders.createOptimized("text/plain"); - private static final CharSequence RESPONSE_TYPE_HTML = HttpHeaders.createOptimized("text/html; charset=UTF-8"); - private static final CharSequence RESPONSE_TYPE_JSON = HttpHeaders.createOptimized("application/json"); - - private static final String HELLO_WORLD = "Hello, world!"; - private static final Buffer HELLO_WORLD_BUFFER = Buffer.buffer(HELLO_WORLD, "UTF-8"); - - private static final CharSequence HEADER_SERVER = HttpHeaders.SERVER; - private static final CharSequence HEADER_DATE = HttpHeaders.DATE; - private static final CharSequence HEADER_CONTENT_TYPE = HttpHeaders.CONTENT_TYPE; - private static final CharSequence HEADER_CONTENT_LENGTH = HttpHeaders.CONTENT_LENGTH; - - private static final CharSequence SERVER = HttpHeaders.createOptimized("vert.x"); - - private static final String SELECT_WORLD = "SELECT id, randomnumber FROM world WHERE id = $1"; - private static final String SELECT_FORTUNE = "SELECT id, message FROM fortune"; - private static final String SELECT_WORLDS = "SELECT id, randomnumber FROM world"; - - private static final Tuple[] tupleCache = new Tuple[10000]; - - public static CharSequence createDateHeader() { - return HttpHeaders.createOptimized(DateTimeFormatter.RFC_1123_DATE_TIME.format(ZonedDateTime.now())); - } - - /** - * Returns a random integer that is a suitable value for both the {@code id} - * and {@code randomNumber} properties of a world object. - * - * @return a random world number - */ - static Integer boxedRandomWorldNumber() { - final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); - final var boxedRnd = BOXED_RND[rndValue - 1]; - assert boxedRnd.intValue() == rndValue; - return boxedRnd; + try { + int parsedValue = Integer.parseInt(param); + return Math.min(500, Math.max(1, parsedValue)); + } catch (NumberFormatException e) { + return 1; } - - static int primitiveRandomWorldNumber() { - final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); - return rndValue; + } + + private static Logger logger = LoggerFactory.getLogger(App.class.getName()); + + private static final Integer[] BOXED_RND = IntStream.range(1, 10001).boxed().toArray(Integer[]::new); + + private static final String PATH_PLAINTEXT = "/plaintext"; + private static final String PATH_JSON = "/json"; + private static final String PATH_DB = "/db"; + private static final String PATH_QUERIES = "/queries"; + private static final String PATH_UPDATES = "/updates"; + private static final String PATH_FORTUNES = "/fortunes"; + private static final String PATH_CACHING = "/cached-queries"; + + private static final CharSequence RESPONSE_TYPE_PLAIN = HttpHeaders.createOptimized("text/plain"); + private static final CharSequence RESPONSE_TYPE_HTML = HttpHeaders.createOptimized("text/html; charset=UTF-8"); + private static final CharSequence RESPONSE_TYPE_JSON = HttpHeaders.createOptimized("application/json"); + + private static final String HELLO_WORLD = "Hello, world!"; + private static final Buffer HELLO_WORLD_BUFFER = Buffer.buffer(HELLO_WORLD, "UTF-8"); + + private static final CharSequence HEADER_SERVER = HttpHeaders.SERVER; + private static final CharSequence HEADER_DATE = HttpHeaders.DATE; + private static final CharSequence HEADER_CONTENT_TYPE = HttpHeaders.CONTENT_TYPE; + private static final CharSequence HEADER_CONTENT_LENGTH = HttpHeaders.CONTENT_LENGTH; + + private static final CharSequence SERVER = HttpHeaders.createOptimized("vert.x"); + + private static final String SELECT_WORLD = "SELECT id, randomnumber FROM world WHERE id = $1"; + private static final String SELECT_FORTUNE = "SELECT id, message FROM fortune"; + private static final String SELECT_WORLDS = "SELECT id, randomnumber FROM world"; + + private static final Tuple[] tupleCache = new Tuple[10000]; + + public static CharSequence createDateHeader() { + return HttpHeaders.createOptimized(DateTimeFormatter.RFC_1123_DATE_TIME.format(ZonedDateTime.now())); + } + + /** + * Returns a random integer that is a suitable value for both the {@code id} + * and {@code randomNumber} properties of a world object. + * + * @return a random world number + */ + static Integer boxedRandomWorldNumber() { + final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); + final var boxedRnd = BOXED_RND[rndValue - 1]; + assert boxedRnd.intValue() == rndValue; + return boxedRnd; + } + + static int primitiveRandomWorldNumber() { + final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); + return rndValue; + } + + static Tuple getRandomTuple() { + final int rndValue = primitiveRandomWorldNumber(); + final Tuple tuple = tupleCache[rndValue - 1]; + return tuple; + } + + private HttpServer server; + private SqlClientInternal client; + private CharSequence dateString; + private MultiMap plaintextHeaders; + private MultiMap jsonHeaders; + + private final RockerOutputFactory factory = BufferRockerOutput.factory(ContentType.RAW); + + private Throwable databaseErr; + private PreparedQuery> SELECT_WORLD_QUERY; + private PreparedQuery>> SELECT_FORTUNE_QUERY; + @SuppressWarnings("unchecked") + private PreparedQuery>[] AGGREGATED_UPDATE_WORLD_QUERY = new PreparedQuery[500]; + private WorldCache WORLD_CACHE; + + private MultiMap plaintextHeaders() { + return HttpHeaders + .headers() + .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_PLAIN) + .add(HEADER_SERVER, SERVER) + .add(HEADER_DATE, dateString) + .copy(false); + } + + private MultiMap jsonHeaders() { + return HttpHeaders + .headers() + .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_JSON) + .add(HEADER_SERVER, SERVER) + .add(HEADER_DATE, dateString) + .copy(false); + } + + @Override + public Future start() throws Exception { + int port = 8080; + server = vertx + .createHttpServer(new HttpServerOptions() + .setHttp2ClearTextEnabled(false) + .setStrictThreadMode(true)) + .requestHandler(App.this); + dateString = createDateHeader(); + plaintextHeaders = plaintextHeaders(); + jsonHeaders = jsonHeaders(); + JsonObject config = config(); + vertx.setPeriodic(1000, id -> { + dateString = createDateHeader(); + plaintextHeaders = plaintextHeaders(); + jsonHeaders = jsonHeaders(); + }); + for (int i = 0; i < 10000; i++) { + tupleCache[i] = Tuple.of(i + 1); } - - static Tuple getRandomTuple() { - final int rndValue = primitiveRandomWorldNumber(); - final Tuple tuple = tupleCache[rndValue - 1]; - return tuple; + PgConnectOptions options = new PgConnectOptions(); + options.setDatabase(config.getString("database", "hello_world")); + options.setHost(config.getString("host", "tfb-database")); + options.setPort(config.getInteger("port", 5432)); + options.setUser(config.getString("username", "benchmarkdbuser")); + options.setPassword(config.getString("password", "benchmarkdbpass")); + options.setCachePreparedStatements(true); + options.setPreparedStatementCacheMaxSize(1024); + options.setPipeliningLimit(256); // Large pipelining means less flushing and we use a single connection anyway + Future clientsInit = initClients(options); + return clientsInit + .transform(ar -> { + databaseErr = ar.cause(); + return server.listen(port); + }); + } + + private Future initClients(PgConnectOptions options) { + return PgConnection.connect(vertx, options) + .flatMap(conn -> { + client = (SqlClientInternal) conn; + List> list = new ArrayList<>(); + Future f1 = conn.prepare(SELECT_WORLD) + .andThen(onSuccess(ps -> SELECT_WORLD_QUERY = ps.query())); + list.add(f1); + Future f2 = conn.prepare(SELECT_FORTUNE) + .andThen(onSuccess(ps -> { + SELECT_FORTUNE_QUERY = ps.query(). + collecting(Collectors.mapping(row -> new Fortune(row.getInteger(0), row.getString(1)), Collectors.toList())); + })); + list.add(f2); + Future f3 = conn.preparedQuery(SELECT_WORLDS) + .collecting(Collectors.mapping(row -> new CachedWorld(row.getInteger(0), row.getInteger(1)), Collectors.toList())) + .execute() + .map(worlds -> new WorldCache(worlds.value())) + .andThen(onSuccess(wc -> WORLD_CACHE = wc)); + list.add(f3); + for (int i = 0; i < AGGREGATED_UPDATE_WORLD_QUERY.length; i++) { + int idx = i; + Future fut = conn + .prepare(buildAggregatedUpdateQuery(1 + idx)) + .andThen(onSuccess(ps -> AGGREGATED_UPDATE_WORLD_QUERY[idx] = ps.query())); + list.add(fut); + } + return Future.join(list); + }); + } + + private static String buildAggregatedUpdateQuery(int len) { + StringBuilder sql = new StringBuilder(); + sql.append("UPDATE world SET randomnumber = CASE ID"); + for (int i = 0; i < len; i++) { + int offset = (i * 2) + 1; + sql.append(" WHEN $").append(offset).append(" THEN $").append(offset + 1); } - - - private HttpServer server; - private SqlClientInternal client; - private CharSequence dateString; - private MultiMap plaintextHeaders; - private MultiMap jsonHeaders; - - private final RockerOutputFactory factory = BufferRockerOutput.factory(ContentType.RAW); - - private Throwable databaseErr; - private PreparedQuery> SELECT_WORLD_QUERY; - private PreparedQuery>> SELECT_FORTUNE_QUERY; - @SuppressWarnings("unchecked") - private PreparedQuery>[] AGGREGATED_UPDATE_WORLD_QUERY = new PreparedQuery[500]; - private WorldCache WORLD_CACHE; - - private MultiMap plaintextHeaders() { - return HttpHeaders - .headers() - .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_PLAIN) - .add(HEADER_SERVER, SERVER) - .add(HEADER_DATE, dateString) - .copy(false); + sql.append(" ELSE randomnumber"); + sql.append(" END WHERE ID IN ($1"); + for (int i = 1; i < len; i++) { + int offset = (i * 2) + 1; + sql.append(",$").append(offset); } - - private MultiMap jsonHeaders() { - return HttpHeaders - .headers() - .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_JSON) - .add(HEADER_SERVER, SERVER) - .add(HEADER_DATE, dateString) - .copy(false); + sql.append(")"); + return sql.toString(); + } + + public static Handler> onSuccess(Handler handler) { + return ar -> { + if (ar.succeeded()) { + handler.handle(ar.result()); + } + }; + } + + @Override + public void handle(HttpServerRequest request) { + try { + switch (request.path()) { + case PATH_PLAINTEXT: + handlePlainText(request); + break; + case PATH_JSON: + handleJson(request); + break; + case PATH_DB: + handleDb(request); + break; + case PATH_QUERIES: + new Queries(request).handle(); + break; + case PATH_UPDATES: + new Update(request).handle(); + break; + case PATH_FORTUNES: + handleFortunes(request); + break; + case PATH_CACHING: + handleCaching(request); + break; + default: + request.response() + .setStatusCode(404) + .end(); + break; + } + } catch (Exception e) { + sendError(request, e); } - - @Override - public Future start() throws Exception { - int port = 8080; - server = vertx - .createHttpServer(new HttpServerOptions() - .setHttp2ClearTextEnabled(false) - .setStrictThreadMode(true)) - .requestHandler(App.this); - dateString = createDateHeader(); - plaintextHeaders = plaintextHeaders(); - jsonHeaders = jsonHeaders(); - JsonObject config = config(); - vertx.setPeriodic(1000, id -> { - dateString = createDateHeader(); - plaintextHeaders = plaintextHeaders(); - jsonHeaders = jsonHeaders(); - }); - - for (int i = 0; i < 10000; i++) { - tupleCache[i] = Tuple.of(i + 1); + } + + @Override + public Future stop() throws Exception { + return server != null ? server.close() : super.stop(); + } + + private void sendError(HttpServerRequest req, Throwable cause) { + logger.error(cause.getMessage(), cause); + req.response().setStatusCode(500).end(); + } + + private void handlePlainText(HttpServerRequest request) { + HttpServerResponse response = request.response(); + response.headers().setAll(plaintextHeaders); + response.end(HELLO_WORLD_BUFFER); + } + + private void handleJson(HttpServerRequest request) { + HttpServerResponse response = request.response(); + response.headers().setAll(jsonHeaders); + response.end(new Message("Hello, World!").toJson()); + } + + private void handleDb(HttpServerRequest req) { + HttpServerResponse resp = req.response(); + SELECT_WORLD_QUERY.execute(getRandomTuple()).onComplete(res -> { + if (res.succeeded()) { + RowIterator resultSet = res.result().iterator(); + if (!resultSet.hasNext()) { + resp.setStatusCode(404).end(); + return; } - - PgConnectOptions options = new PgConnectOptions(); - options.setDatabase(config.getString("database", "hello_world")); - options.setHost(config.getString("host", "tfb-database")); - options.setPort(config.getInteger("port", 5432)); - options.setUser(config.getString("username", "benchmarkdbuser")); - options.setPassword(config.getString("password", "benchmarkdbpass")); - options.setCachePreparedStatements(true); - options.setPreparedStatementCacheMaxSize(1024); - options.setPipeliningLimit(256); // Large pipelining means less flushing and we use a single connection anyway - Future clientsInit = initClients(options); - return clientsInit - .transform(ar -> { - databaseErr = ar.cause(); - return server.listen(port); - }); - } - - private Future initClients(PgConnectOptions options) { - return PgConnection.connect(vertx, options) - .flatMap(conn -> { - client = (SqlClientInternal) conn; - List> list = new ArrayList<>(); - Future f1 = conn.prepare(SELECT_WORLD) - .andThen(onSuccess(ps -> SELECT_WORLD_QUERY = ps.query())); - list.add(f1); - Future f2 = conn.prepare(SELECT_FORTUNE) - .andThen(onSuccess(ps -> { - SELECT_FORTUNE_QUERY = ps.query(). - collecting(Collectors.mapping(row -> new Fortune(row.getInteger(0), row.getString(1)), Collectors.toList())); - })); - list.add(f2); - Future f3 = conn.preparedQuery(SELECT_WORLDS) - .collecting(Collectors.mapping(row -> new CachedWorld(row.getInteger(0), row.getInteger(1)), Collectors.toList())) - .execute() - .map(worlds -> new WorldCache(worlds.value())) - .andThen(onSuccess(wc -> WORLD_CACHE = wc)); - list.add(f3); - for (int i = 0; i < AGGREGATED_UPDATE_WORLD_QUERY.length; i++) { - int idx = i; - Future fut = conn - .prepare(buildAggregatedUpdateQuery(1 + idx)) - .andThen(onSuccess(ps -> AGGREGATED_UPDATE_WORLD_QUERY[idx] = ps.query())); - list.add(fut); - } - return Future.join(list); - }); + Row row = resultSet.next(); + World world = new World(row.getInteger(0), row.getInteger(1)); + MultiMap headers = resp.headers(); + headers.add(HttpHeaders.SERVER, SERVER); + headers.add(HttpHeaders.DATE, dateString); + headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); + resp.end(world.toJson()); + } else { + sendError(req, res.cause()); + } + }); + } + + class Queries implements Handler>> { + + boolean failed; + final World[] worlds; + final HttpServerRequest req; + final HttpServerResponse resp; + final int queries; + int worldsIndex; + + public Queries(HttpServerRequest req) { + int queries = getQueries(req); + + this.req = req; + this.resp = req.response(); + this.queries = queries; + this.worlds = new World[queries]; + this.worldsIndex = 0; } - private static String buildAggregatedUpdateQuery(int len) { - StringBuilder sql = new StringBuilder(); - sql.append("UPDATE world SET randomnumber = CASE ID"); - for (int i = 0; i < len; i++) { - int offset = (i * 2) + 1; - sql.append(" WHEN $").append(offset).append(" THEN $").append(offset + 1); + private void handle() { + client.group(/*queries, */c -> { + for (int i = 0; i < queries; i++) { + c.preparedQuery(SELECT_WORLD) + .execute(getRandomTuple()) + .onComplete(this); } - sql.append(" ELSE randomnumber"); - sql.append(" END WHERE ID IN ($1"); - for (int i = 1; i < len; i++) { - int offset = (i * 2) + 1; - sql.append(",$").append(offset); - } - sql.append(")"); - return sql.toString(); - } - - public static Handler> onSuccess(Handler handler) { - return ar -> { - if (ar.succeeded()) { - handler.handle(ar.result()); - } - }; + }); } @Override - public void handle(HttpServerRequest request) { - try { - switch (request.path()) { - case PATH_PLAINTEXT: - handlePlainText(request); - break; - case PATH_JSON: - handleJson(request); - break; - case PATH_DB: - handleDb(request); - break; - case PATH_QUERIES: - new Queries(request).handle(); - break; - case PATH_UPDATES: - new Update(request).handle(); - break; - case PATH_FORTUNES: - handleFortunes(request); - break; - case PATH_CACHING: - handleCaching(request); - break; - default: - request.response() - .setStatusCode(404) - .end(); - break; - } - } catch (Exception e) { - sendError(request, e); + public void handle(AsyncResult> ar) { + if (!failed) { + if (ar.failed()) { + failed = true; + sendError(req, ar.cause()); + return; } - } - - @Override - public Future stop() throws Exception { - return server != null ? server.close() : super.stop(); - } - private void sendError(HttpServerRequest req, Throwable cause) { - logger.error(cause.getMessage(), cause); - req.response().setStatusCode(500).end(); + // we need a final reference + final Tuple row = ar.result().iterator().next(); + worlds[worldsIndex++] = new World(row.getInteger(0), row.getInteger(1)); + + // stop condition + if (worldsIndex == queries) { + MultiMap headers = resp.headers(); + headers.add(HttpHeaders.SERVER, SERVER); + headers.add(HttpHeaders.DATE, dateString); + headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); + resp.end(World.toJson(worlds)); + } + } } + } - private void handlePlainText(HttpServerRequest request) { - HttpServerResponse response = request.response(); - response.headers().setAll(plaintextHeaders); - response.end(HELLO_WORLD_BUFFER); - } + private class Update { - private void handleJson(HttpServerRequest request) { - HttpServerResponse response = request.response(); - response.headers().setAll(jsonHeaders); - response.end(new Message("Hello, World!").toJson()); - } + private final HttpServerRequest request; + private final World[] worldsToUpdate; + private boolean failed; + private int selectWorldCompletedCount; - private void handleDb(HttpServerRequest req) { - HttpServerResponse resp = req.response(); - SELECT_WORLD_QUERY.execute(getRandomTuple()).onComplete(res -> { - if (res.succeeded()) { - RowIterator resultSet = res.result().iterator(); - if (!resultSet.hasNext()) { - resp.setStatusCode(404).end(); - return; - } - Row row = resultSet.next(); - World world = new World(row.getInteger(0), row.getInteger(1)); - MultiMap headers = resp.headers(); - headers.add(HttpHeaders.SERVER, SERVER); - headers.add(HttpHeaders.DATE, dateString); - headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); - resp.end(world.toJson()); - } else { - sendError(req, res.cause()); - } - }); + public Update(HttpServerRequest request) { + this.request = request; + this.worldsToUpdate = new World[getQueries(request)]; } - class Queries implements Handler>> { - - boolean failed; - final World[] worlds; - final HttpServerRequest req; - final HttpServerResponse resp; - final int queries; - int worldsIndex; - - public Queries(HttpServerRequest req) { - int queries = getQueries(req); - - this.req = req; - this.resp = req.response(); - this.queries = queries; - this.worlds = new World[queries]; - this.worldsIndex = 0; - } - - private void handle() { - client.group(/*queries, */c -> { - for (int i = 0; i < queries; i++) { - c.preparedQuery(SELECT_WORLD) - .execute(getRandomTuple()) - .onComplete(this); - } - }); - } - - @Override - public void handle(AsyncResult> ar) { + public void handle() { + client.group(/*worldsToUpdate.length, */c -> { + final PreparedQuery> preparedQuery = c.preparedQuery(App.SELECT_WORLD); + for (int i = 0; i < worldsToUpdate.length; i++) { + final int index = i; + preparedQuery.execute(getRandomTuple()).onComplete(res -> { if (!failed) { - if (ar.failed()) { - failed = true; - sendError(req, ar.cause()); - return; - } - - // we need a final reference - final Tuple row = ar.result().iterator().next(); - worlds[worldsIndex++] = new World(row.getInteger(0), row.getInteger(1)); - - // stop condition - if (worldsIndex == queries) { - MultiMap headers = resp.headers(); - headers.add(HttpHeaders.SERVER, SERVER); - headers.add(HttpHeaders.DATE, dateString); - headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); - resp.end(World.toJson(worlds)); - } + if (res.failed()) { + failed = true; + sendError(request, res.cause()); + return; + } + worldsToUpdate[index] = new World(res.result().iterator().next().getInteger(0), primitiveRandomWorldNumber()); + if (++selectWorldCompletedCount == worldsToUpdate.length) { + randomWorldsQueryCompleted(); + } } + }); } + }); } - private class Update { - - private final HttpServerRequest request; - private final World[] worldsToUpdate; - private boolean failed; - private int selectWorldCompletedCount; - - public Update(HttpServerRequest request) { - this.request = request; - this.worldsToUpdate = new World[getQueries(request)]; - } - - public void handle() { - client.group(/*worldsToUpdate.length, */c -> { - final PreparedQuery> preparedQuery = c.preparedQuery(App.SELECT_WORLD); - for (int i = 0; i < worldsToUpdate.length; i++) { - final int index = i; - preparedQuery.execute(getRandomTuple()).onComplete(res -> { - if (!failed) { - if (res.failed()) { - failed = true; - sendError(request, res.cause()); - return; - } - worldsToUpdate[index] = new World(res.result().iterator().next().getInteger(0), primitiveRandomWorldNumber()); - if (++selectWorldCompletedCount == worldsToUpdate.length) { - randomWorldsQueryCompleted(); - } - } - }); - } - }); - } - - private void randomWorldsQueryCompleted() { - Arrays.sort(worldsToUpdate); - final List params = new ArrayList<>(worldsToUpdate.length * 2); - for (int i = 0, count = worldsToUpdate.length; i < count; i++) { - var world = worldsToUpdate[i]; - params.add(world.getId()); - params.add(world.getRandomNumber()); - } - AGGREGATED_UPDATE_WORLD_QUERY[worldsToUpdate.length - 1] - .execute(Tuple.wrap(params)) - .onComplete(updateResult -> { - if (updateResult.failed()) { - sendError(request, updateResult.cause()); - return; - } - sendResponse(); - }); - } - - private void sendResponse() { - var res = request.response(); - MultiMap headers = res.headers(); - headers.add(HttpHeaders.SERVER, App.SERVER); - headers.add(HttpHeaders.DATE, dateString); - headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); - Buffer buff = WorldJsonSerializer.toJsonBuffer(worldsToUpdate); - res.end(buff); + private void randomWorldsQueryCompleted() { + Arrays.sort(worldsToUpdate); + final List params = new ArrayList<>(worldsToUpdate.length * 2); + for (int i = 0, count = worldsToUpdate.length;i < count;i++) { + var world = worldsToUpdate[i]; + params.add(world.getId()); + params.add(world.getRandomNumber()); + } + AGGREGATED_UPDATE_WORLD_QUERY[worldsToUpdate.length - 1] + .execute(Tuple.wrap(params)) + .onComplete(updateResult -> { + if (updateResult.failed()) { + sendError(request, updateResult.cause()); + return; } + sendResponse(); + }); } - private void handleFortunes(HttpServerRequest req) { - SELECT_FORTUNE_QUERY - .execute() - .onComplete(ar -> { - HttpServerResponse response = req.response(); - if (ar.succeeded()) { - SqlResult> result = ar.result(); - if (result.size() == 0) { - response.setStatusCode(404).end("No results"); - return; - } - List fortunes = result.value(); - fortunes.add(new Fortune(0, "Additional fortune added at request time.")); - Collections.sort(fortunes); - MultiMap headers = response.headers(); - headers.add(HttpHeaders.SERVER, SERVER); - headers.add(HttpHeaders.DATE, dateString); - headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_HTML); - FortunesTemplate template = FortunesTemplate.template(fortunes); - response.end(template.render(factory).buffer()); - } else { - sendError(req, ar.cause()); - } - }); + private void sendResponse() { + var res = request.response(); + MultiMap headers = res.headers(); + headers.add(HttpHeaders.SERVER, App.SERVER); + headers.add(HttpHeaders.DATE, dateString); + headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_JSON); + Buffer buff = WorldJsonSerializer.toJsonBuffer(worldsToUpdate); + res.end(buff); } - - private void handleCaching(HttpServerRequest req) { - int count = 1; - try { - String countStr = req.getParam("count"); - if (countStr != null) { - count = Integer.parseInt(countStr); - } - } catch (NumberFormatException ignore) { + } + + private void handleFortunes(HttpServerRequest req) { + SELECT_FORTUNE_QUERY + .execute() + .onComplete(ar -> { + HttpServerResponse response = req.response(); + if (ar.succeeded()) { + SqlResult> result = ar.result(); + if (result.size() == 0) { + response.setStatusCode(404).end("No results"); + return; } - count = Math.max(1, count); - count = Math.min(500, count); - List worlds = WORLD_CACHE.getCachedWorld(count); - HttpServerResponse response = req.response(); + List fortunes = result.value(); + fortunes.add(new Fortune(0, "Additional fortune added at request time.")); + Collections.sort(fortunes); MultiMap headers = response.headers(); - headers - .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_JSON) - .add(HEADER_SERVER, SERVER) - .add(HEADER_DATE, dateString); - response.end(CachedWorld.toJson(worlds)); + headers.add(HttpHeaders.SERVER, SERVER); + headers.add(HttpHeaders.DATE, dateString); + headers.add(HttpHeaders.CONTENT_TYPE, RESPONSE_TYPE_HTML); + FortunesTemplate template = FortunesTemplate.template(fortunes); + response.end(template.render(factory).buffer()); + } else { + sendError(req, ar.cause()); + } + }); + } + + private void handleCaching(HttpServerRequest req) { + int count = 1; + try { + String countStr = req.getParam("count"); + if (countStr != null) { + count = Integer.parseInt(countStr); + } + } catch (NumberFormatException ignore) { } - - public static void main(String[] args) throws Exception { - int eventLoopPoolSize = NUM_PROCESSORS; - String sizeProp = System.getProperty("vertx.eventLoopPoolSize"); - if (sizeProp != null) { - try { - eventLoopPoolSize = Integer.parseInt(sizeProp); - } catch (NumberFormatException e) { - e.printStackTrace(); - } - } - JsonObject config = new JsonObject(new String(Files.readAllBytes(new File(args[0]).toPath()))); - Vertx vertx = Vertx.vertx(new VertxOptions() - .setEventLoopPoolSize(eventLoopPoolSize) - .setPreferNativeTransport(true) - .setDisableTCCL(true) - ); - vertx.exceptionHandler(err -> { - err.printStackTrace(); - }); - printConfig((VertxInternal) vertx); - vertx.deployVerticle( - App.class.getName(), - new DeploymentOptions().setInstances(eventLoopPoolSize).setConfig(config)) - .onComplete(event -> { - if (event.succeeded()) { - logger.info("Server listening on port " + 8080); - } else { - logger.error("Unable to start your application", event.cause()); - } - }); + count = Math.max(1, count); + count = Math.min(500, count); + List worlds = WORLD_CACHE.getCachedWorld(count); + HttpServerResponse response = req.response(); + MultiMap headers = response.headers(); + headers + .add(HEADER_CONTENT_TYPE, RESPONSE_TYPE_JSON) + .add(HEADER_SERVER, SERVER) + .add(HEADER_DATE, dateString); + response.end(CachedWorld.toJson(worlds)); + } + + public static void main(String[] args) throws Exception { + int eventLoopPoolSize = NUM_PROCESSORS; + String sizeProp = System.getProperty("vertx.eventLoopPoolSize"); + if (sizeProp != null) { + try { + eventLoopPoolSize = Integer.parseInt(sizeProp); + } catch (NumberFormatException e) { + e.printStackTrace(); + } } - - private static void printConfig(VertxInternal vertx) { - boolean nativeTransport = vertx.isNativeTransportEnabled(); - String transport = vertx.transport().getClass().getSimpleName(); - String version = "unknown"; - try { - InputStream in = Vertx.class.getClassLoader().getResourceAsStream("META-INF/vertx/vertx-version.txt"); - if (in == null) { - in = Vertx.class.getClassLoader().getResourceAsStream("vertx-version.txt"); - } - ByteArrayOutputStream out = new ByteArrayOutputStream(); - byte[] buffer = new byte[256]; - while (true) { - int amount = in.read(buffer); - if (amount == -1) { - break; - } - out.write(buffer, 0, amount); - } - version = out.toString(); - } catch (IOException e) { - logger.error("Could not read Vertx version", e); - ; - } - logger.info("Vertx: " + version); - logger.info("Processors: " + NUM_PROCESSORS); - logger.info("Event Loop Size: " + ((MultithreadEventExecutorGroup) vertx.nettyEventLoopGroup()).executorCount()); - logger.info("Native transport : " + nativeTransport); - logger.info("Transport : " + transport); - logger.info("Netty buffer bound check : " + System.getProperty("io.netty.buffer.checkBounds")); - logger.info("Netty buffer accessibility check : " + System.getProperty("io.netty.buffer.checkAccessible")); - for (SysProps sysProp : SysProps.values()) { - logger.info(sysProp.name + " : " + sysProp.get()); + JsonObject config = new JsonObject(new String(Files.readAllBytes(new File(args[0]).toPath()))); + Vertx vertx = Vertx.vertx(new VertxOptions() + .setEventLoopPoolSize(eventLoopPoolSize) + .setPreferNativeTransport(true) + .setDisableTCCL(true) + ); + vertx.exceptionHandler(err -> { + err.printStackTrace(); + }); + printConfig((VertxInternal) vertx); + vertx.deployVerticle( + App.class.getName(), + new DeploymentOptions().setInstances(eventLoopPoolSize).setConfig(config)) + .onComplete(event -> { + if (event.succeeded()) { + logger.info("Server listening on port " + 8080); + } else { + logger.error("Unable to start your application", event.cause()); + } + }); + } + + private static void printConfig(VertxInternal vertx) { + boolean nativeTransport = vertx.isNativeTransportEnabled(); + String transport = vertx.transport().getClass().getSimpleName(); + String version = "unknown"; + try { + InputStream in = Vertx.class.getClassLoader().getResourceAsStream("META-INF/vertx/vertx-version.txt"); + if (in == null) { + in = Vertx.class.getClassLoader().getResourceAsStream("vertx-version.txt"); + } + ByteArrayOutputStream out = new ByteArrayOutputStream(); + byte[] buffer = new byte[256]; + while (true) { + int amount = in.read(buffer); + if (amount == -1) { + break; } + out.write(buffer, 0, amount); + } + version = out.toString(); + } catch (IOException e) { + logger.error("Could not read Vertx version", e);; + } + logger.info("Vertx: " + version); + logger.info("Processors: " + NUM_PROCESSORS); + logger.info("Event Loop Size: " + ((MultithreadEventExecutorGroup)vertx.nettyEventLoopGroup()).executorCount()); + logger.info("Native transport : " + nativeTransport); + logger.info("Transport : " + transport); + logger.info("Netty buffer bound check : " + System.getProperty("io.netty.buffer.checkBounds")); + logger.info("Netty buffer accessibility check : " + System.getProperty("io.netty.buffer.checkAccessible")); + for (SysProps sysProp : SysProps.values()) { + logger.info(sysProp.name + " : " + sysProp.get()); } + } } From cceb3cd83ff571d29587bcbad6e6bf8480d154b1 Mon Sep 17 00:00:00 2001 From: doxlik Date: Mon, 23 Mar 2026 01:53:15 +0400 Subject: [PATCH 4/4] adding tuple cache to Quarkus --- .../benchmark/repository/WorldRepository.java | 8 ++--- .../benchmark/resource/DbResource.java | 26 +++++++++++++++-- .../benchmark/repository/WorldRepository.java | 29 +++++++++++++++---- 3 files changed, 51 insertions(+), 12 deletions(-) diff --git a/frameworks/Java/quarkus/reactive-routes-pgclient/src/main/java/io/quarkus/benchmark/repository/WorldRepository.java b/frameworks/Java/quarkus/reactive-routes-pgclient/src/main/java/io/quarkus/benchmark/repository/WorldRepository.java index 630be2b3861..fabc69285b2 100644 --- a/frameworks/Java/quarkus/reactive-routes-pgclient/src/main/java/io/quarkus/benchmark/repository/WorldRepository.java +++ b/frameworks/Java/quarkus/reactive-routes-pgclient/src/main/java/io/quarkus/benchmark/repository/WorldRepository.java @@ -28,18 +28,18 @@ public JsonWorld(final Integer id, final Integer randomNumber) { PgClients clients; - public Uni findAsJsonWorld(final Integer id) { + public Uni findAsJsonWorld(final Tuple id) { return clients.getClient().preparedQuery("SELECT id, randomNumber FROM World WHERE id = $1") - .execute(Tuple.of(id)) + .execute(id) .map(rowset -> { final Row row = rowset.iterator().next(); return new JsonWorld(row.getInteger(0), row.getInteger(1)); }); } - public Uni find(final Integer id) { + public Uni find(final Tuple id) { return clients.getClient().preparedQuery("SELECT id, randomNumber FROM World WHERE id = $1") - .execute(Tuple.of(id)) + .execute(id) .map(rowset -> { final Row row = rowset.iterator().next(); return new World(row.getInteger(0), row.getInteger(1)); diff --git a/frameworks/Java/quarkus/reactive-routes-pgclient/src/main/java/io/quarkus/benchmark/resource/DbResource.java b/frameworks/Java/quarkus/reactive-routes-pgclient/src/main/java/io/quarkus/benchmark/resource/DbResource.java index 4237a81f8dc..f6f83b9ff19 100644 --- a/frameworks/Java/quarkus/reactive-routes-pgclient/src/main/java/io/quarkus/benchmark/resource/DbResource.java +++ b/frameworks/Java/quarkus/reactive-routes-pgclient/src/main/java/io/quarkus/benchmark/resource/DbResource.java @@ -12,6 +12,7 @@ import io.smallrye.mutiny.Uni; import io.vertx.core.json.JsonArray; import io.vertx.ext.web.RoutingContext; +import io.vertx.mutiny.sqlclient.Tuple; import jakarta.inject.Inject; import jakarta.inject.Singleton; @@ -24,7 +25,7 @@ public class DbResource extends BaseResource { @Route(path = "db") public void db(final RoutingContext rc) { - worldRepository.findAsJsonWorld(boxedRandomWorldNumber()) + worldRepository.findAsJsonWorld(getRandomTuple()) .subscribe().with(world -> sendJson(rc, world), t -> handleFail(rc, t)); } @@ -36,7 +37,7 @@ public void queries(final RoutingContext rc) { final var ret = new JsonWorld[worlds.length]; // replace below with a for loop Arrays.setAll(worlds, i -> { - return worldRepository.findAsJsonWorld(boxedRandomWorldNumber()).map(w -> ret[i] = w); + return worldRepository.findAsJsonWorld(getRandomTuple()).map(w -> ret[i] = w); }); Uni.combine().all().unis(worlds) @@ -75,10 +76,18 @@ public void updates(final RoutingContext rc) { } private Uni randomWorld() { - return worldRepository.find(boxedRandomWorldNumber()); + return worldRepository.find(getRandomTuple()); } private static final Integer[] BOXED_RND = IntStream.range(1, 10001).boxed().toArray(Integer[]::new); + private static final Tuple[] tupleCache = new Tuple[10000]; + + static { + for (int i = 0; i < 10000; i++) { + tupleCache[i] = Tuple.of(i + 1); + } + } + private static Integer boxedRandomWorldNumber() { final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); @@ -87,6 +96,17 @@ private static Integer boxedRandomWorldNumber() { return boxedRnd; } + private static int primitiveRandomWorldNumber() { + final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); + return rndValue; + } + + private static Tuple getRandomTuple() { + final int rndValue = primitiveRandomWorldNumber(); + final Tuple tuple = tupleCache[rndValue - 1]; + return tuple; + } + private static int parseQueryCount(final String textValue) { if (textValue == null) { return 1; diff --git a/frameworks/Java/quarkus/vertx/src/main/java/io/quarkus/benchmark/repository/WorldRepository.java b/frameworks/Java/quarkus/vertx/src/main/java/io/quarkus/benchmark/repository/WorldRepository.java index 3095137ae5e..e6532742720 100644 --- a/frameworks/Java/quarkus/vertx/src/main/java/io/quarkus/benchmark/repository/WorldRepository.java +++ b/frameworks/Java/quarkus/vertx/src/main/java/io/quarkus/benchmark/repository/WorldRepository.java @@ -28,6 +28,14 @@ public class WorldRepository { private PgConnectionPool connectionPool; private static final Integer[] BOXED_RND = IntStream.range(1, 10001).boxed().toArray(Integer[]::new); + private static final Tuple[] tupleCache = new Tuple[10000]; + + static { + for (int i = 0; i < 10000; i++) { + tupleCache[i] = Tuple.of(i + 1); + } + } + private static Integer boxedRandomWorldNumber() { final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); @@ -36,8 +44,20 @@ private static Integer boxedRandomWorldNumber() { return boxedRnd; } + private static int primitiveRandomWorldNumber() { + final int rndValue = ThreadLocalRandom.current().nextInt(1, 10001); + return rndValue; + } + + private static Tuple getRandomTuple() { + final int rndValue = primitiveRandomWorldNumber(); + final Tuple tuple = tupleCache[rndValue - 1]; + return tuple; + } + + public void loadRandomJsonWorld(final Handler> worldHandler) { - connectionPool.pgConnection().selectWorldQuery().execute(Tuple.of(boxedRandomWorldNumber()), randomWorldRow -> { + connectionPool.pgConnection().selectWorldQuery().execute(getRandomTuple(), randomWorldRow -> { if (randomWorldRow.succeeded()) { final RowIterator resultSet = randomWorldRow.result().iterator(); if (!resultSet.hasNext()) { @@ -101,16 +121,15 @@ private void run() { connection.rawConnection().group(c -> { final PreparedQuery> preparedQuery = c.preparedQuery(PgConnectionPool.SELECT_WORLD); for (int i = 0; i < worldsToUpdate.length; i++) { - final Integer id = boxedRandomWorldNumber(); final int index = i; - preparedQuery.execute(Tuple.of(id), worldId -> { + preparedQuery.execute(getRandomTuple(), worldId -> { if (!failed) { if (worldId.failed()) { failed = true; resultHandler.handle(Future.failedFuture(worldId.cause())); return; } - worldsToUpdate[index] = new World(worldId.result().iterator().next().getInteger(0), boxedRandomWorldNumber()); + worldsToUpdate[index] = new World(worldId.result().iterator().next().getInteger(0), primitiveRandomWorldNumber()); if (++selectWorldCompletedCount == worldsToUpdate.length) { randomWorldsQueryCompleted(); } @@ -159,7 +178,7 @@ public static void execute(final PgConnectionPool connectionPool, final int quer private void run() { connection.rawConnection().group(c -> { for (int i = 0; i < count; i++) { - c.preparedQuery(PgConnectionPool.SELECT_WORLD).execute(Tuple.of(boxedRandomWorldNumber()), this); + c.preparedQuery(PgConnectionPool.SELECT_WORLD).execute(getRandomTuple(), this); } }); }