From 4ac6bb3cfaaf9440102687427c638bc08d030087 Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Fri, 25 Sep 2026 15:52:23 -0700 Subject: [PATCH 1/2] Implemented running with remote server. Added tests with compression --- performance/README.md | 32 +++++ .../clickhouse/benchmark/BenchmarkRunner.java | 6 +- .../clickhouse/benchmark/TestEnvironment.java | 104 +++++++++----- .../benchmark/clients/BenchmarkBase.java | 133 ++++++++---------- .../clients/ConcurrentInsertClient.java | 4 +- .../clients/ConcurrentQueryClient.java | 4 +- .../benchmark/clients/DataTypes.java | 2 +- .../benchmark/clients/Deserializers.java | 2 +- .../benchmark/clients/InsertClient.java | 68 +++++++-- .../benchmark/clients/JDBCInsert.java | 22 ++- .../benchmark/clients/JDBCQuery.java | 46 +++++- .../benchmark/clients/MixedWorkload.java | 4 +- .../benchmark/clients/Serializers.java | 2 +- .../clickhouse/benchmark/data/DataSet.java | 8 ++ 14 files changed, 295 insertions(+), 142 deletions(-) diff --git a/performance/README.md b/performance/README.md index ce0f72e17..434f31d8f 100644 --- a/performance/README.md +++ b/performance/README.md @@ -16,6 +16,38 @@ mvn compile exec:exec -Dexec.executable=java -Dexec.args="-classpath %classpath -input sample_dataset.sql -name default -rows 10" ``` +#### Target Server + +By default the benchmarks start a local ClickHouse Docker container and run +against it. To run against an existing remote server instead (ClickHouse Cloud +or any self-hosted instance), set `CLICKHOUSE_URL`: + +```shell +export CLICKHOUSE_URL="https://default:my-password@abc123.clickhouse.cloud:8443?cluster=true" +mvn compile exec:exec +``` + +```shell +export CLICKHOUSE_URL="http://default@localhost:8123" +mvn compile exec:exec +``` + +The URL is parsed as: `://[:]@[:][?cluster=true]` +- scheme (`http`/`https`) selects whether SSL is used +- username/password come from the URL's user-info; if the password is omitted + the client connects without one, and if the username is omitted it defaults + to `default` +- port is optional; if omitted it defaults to `8443` for `https` and `8123` + for `http` +- `cluster` is optional (default `false`) and tells the benchmarks whether the + remote server is part of a replicated cluster (e.g. ClickHouse Cloud). When + `true`, a `SYSTEM SYNC REPLICA` is issued after writes. Leave it unset for a + plain standalone remote server — its tables aren't replicated and it will + reject that statement. + +When `CLICKHOUSE_URL` is unset, the local Docker container is started +automatically and no other configuration is needed. + #### Running Benchmarks With default settings : diff --git a/performance/src/main/java/com/clickhouse/benchmark/BenchmarkRunner.java b/performance/src/main/java/com/clickhouse/benchmark/BenchmarkRunner.java index 81fee5556..ea58ca5b6 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/BenchmarkRunner.java +++ b/performance/src/main/java/com/clickhouse/benchmark/BenchmarkRunner.java @@ -28,7 +28,7 @@ import java.util.TreeSet; import java.util.concurrent.TimeUnit; -import static com.clickhouse.benchmark.TestEnvironment.isCloud; +import static com.clickhouse.benchmark.TestEnvironment.isRemote; public class BenchmarkRunner { @@ -38,11 +38,11 @@ public static void main(String[] args) throws Exception { LOGGER.info("Starting Benchmarks"); Map options = parseArgs(args); System.out.println("Start Benchmarks with options: " + options); - final String env = isCloud() ? "cloud" : "local"; + final String env = isRemote() ? "remote" : "local"; final long time = System.currentTimeMillis(); final int measurementIterations = Integer.parseInt(options.getOrDefault("-m", "10")); - final int measurementTime = Integer.parseInt(options.getOrDefault("-t", "" + (isCloud() ? 30 : 10))); + final int measurementTime = Integer.parseInt(options.getOrDefault("-t", "" + (isRemote() ? 30 : 10))); final String resultFile = String.format("jmh-results-%s-%s.json", env, time); final String outputFile = String.format("jmh-results-%s-%s.out", env, time); final String datasetName = options.getOrDefault("-d", "file://default.csv"); diff --git a/performance/src/main/java/com/clickhouse/benchmark/TestEnvironment.java b/performance/src/main/java/com/clickhouse/benchmark/TestEnvironment.java index 0f0d6a4eb..0121a446e 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/TestEnvironment.java +++ b/performance/src/main/java/com/clickhouse/benchmark/TestEnvironment.java @@ -10,6 +10,7 @@ import org.testcontainers.containers.wait.strategy.Wait; import java.net.InetSocketAddress; +import java.net.URI; import java.time.Duration; import java.util.Collections; @@ -25,50 +26,75 @@ public class TestEnvironment { //Environment Variables - public static boolean isCloud() { - return System.getenv("CLICKHOUSE_HOST") != null; + // Set CLICKHOUSE_URL to point the benchmarks at an existing remote ClickHouse + // server (ClickHouse Cloud or any self-hosted instance), e.g.: + // https://default:my-password@abc123.clickhouse.cloud:8443?cluster=true + // http://default@localhost:8123 + // The scheme selects HTTP vs HTTPS/SSL, and credentials come from the URL's + // user-info (username[:password]). When unset, a local Docker container is + // started automatically instead. + // + // The optional "cluster" query parameter (default false) tells the benchmarks + // whether the remote server is part of a replicated cluster (e.g. ClickHouse + // Cloud) and therefore needs a SYSTEM SYNC REPLICA after writes. Leave it + // unset/false for a plain standalone remote server, whose tables aren't + // replicated and would reject that statement. + private static URI getRemoteUrl() { + String url = System.getenv("CLICKHOUSE_URL"); + return url == null ? null : URI.create(url); } - public static String getHost() { - String host = System.getenv("CLICKHOUSE_HOST"); - if (host == null) { - host = container.getHost(); - } - - return host; + public static boolean isRemote() { + return getRemoteUrl() != null; } - public static int getPort() { - String port = System.getenv("CLICKHOUSE_PORT"); - if (port == null) { - if (isCloud()) {//Default handling for ClickHouse Cloud - port = "8443"; - } else { - port = String.valueOf(container.getMappedPort(8123)); + public static boolean isSsl() { + URI url = getRemoteUrl(); + return url != null && "https".equalsIgnoreCase(url.getScheme()); + } + public static boolean isCluster() { + URI url = getRemoteUrl(); + String query = url == null ? null : url.getQuery(); + if (query == null) { + return false; + } + for (String param : query.split("&")) { + String[] kv = param.split("=", 2); + if (kv.length == 2 && kv[0].equalsIgnoreCase("cluster")) { + return Boolean.parseBoolean(kv[1]); } } - - return Integer.parseInt(port); + return false; } - public static String getPassword() { - String password = System.getenv("CLICKHOUSE_PASSWORD"); - if (password == null) { - if (isCloud()) { - password = System.getenv("CLICKHOUSE_PASSWORD"); - } else { - password = container.getPassword(); - } + public static String getHost() { + URI url = getRemoteUrl(); + return url != null ? url.getHost() : container.getHost(); + } + public static int getPort() { + URI url = getRemoteUrl(); + if (url != null) { + return url.getPort() != -1 ? url.getPort() : (isSsl() ? 8443 : 8123); } - return password; + return container.getMappedPort(8123); } public static String getUsername() { - String username = System.getenv("CLICKHOUSE_USERNAME"); - if (username == null) { - if (isCloud()) { - username = "default"; - } else { - username = container.getUsername(); - } + URI url = getRemoteUrl(); + if (url == null) { + return container.getUsername(); + } + String userInfo = url.getUserInfo(); + if (userInfo == null) { + return "default"; + } + int sep = userInfo.indexOf(':'); + return sep == -1 ? userInfo : userInfo.substring(0, sep); + } + public static String getPassword() { + URI url = getRemoteUrl(); + if (url == null) { + return container.getPassword(); } - return username; + String userInfo = url.getUserInfo(); + int sep = userInfo == null ? -1 : userInfo.indexOf(':'); + return sep == -1 ? null : userInfo.substring(sep + 1); } public static ClickHouseNode getServer() { return serverNode; @@ -79,8 +105,8 @@ public static ClickHouseNode getServer() { public static void setupEnvironment() { LOGGER.info("Initializing ClickHouse test environment..."); - if (isCloud()) { - LOGGER.info("Using ClickHouse Cloud"); + if (isRemote()) { + LOGGER.info("Using remote ClickHouse server at {}:{}", getHost(), getPort()); container = null; } else { LOGGER.info("Using ClickHouse Docker container"); @@ -95,7 +121,7 @@ public static void setupEnvironment() { serverNode = ClickHouseNode.builder(ClickHouseNode.builder().build()) .address(ClickHouseProtocol.HTTP, new InetSocketAddress(getHost(), getPort())) .credentials(ClickHouseCredentials.fromUserAndPassword(getUsername(), getPassword())) - .options(Collections.singletonMap(ClickHouseClientOption.SSL.getKey(), isCloud() ? "true" : "false")) + .options(Collections.singletonMap(ClickHouseClientOption.SSL.getKey(), isSsl() ? "true" : "false")) .database(DB_NAME) .build(); createDatabase(); @@ -103,7 +129,7 @@ public static void setupEnvironment() { public static void cleanupEnvironment() { LOGGER.info("Cleaning up ClickHouse test environment..."); - if (isCloud()) { + if (isRemote()) { dropDatabase(); } diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/BenchmarkBase.java b/performance/src/main/java/com/clickhouse/benchmark/clients/BenchmarkBase.java index a74f5dbe9..4bfe2d1bb 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/BenchmarkBase.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/BenchmarkBase.java @@ -4,17 +4,13 @@ import com.clickhouse.benchmark.data.FileDataSet; import com.clickhouse.benchmark.data.SimpleDataSet; import com.clickhouse.benchmark.data.SyntheticDataSet; -import com.clickhouse.client.ClickHouseClient; -import com.clickhouse.client.ClickHouseCredentials; -import com.clickhouse.client.ClickHouseNode; -import com.clickhouse.client.ClickHouseNodeSelector; -import com.clickhouse.client.ClickHouseProtocol; -import com.clickhouse.client.ClickHouseResponse; +import com.clickhouse.client.*; import com.clickhouse.client.api.Client; import com.clickhouse.client.api.ClientConfigProperties; import com.clickhouse.client.api.enums.Protocol; import com.clickhouse.client.api.insert.InsertResponse; import com.clickhouse.client.api.query.GenericRecord; +import com.clickhouse.client.config.ClickHouseClientOption; import com.clickhouse.client.config.ClickHouseDefaults; import com.clickhouse.data.ClickHouseDataProcessor; import com.clickhouse.data.ClickHouseFormat; @@ -23,12 +19,7 @@ import com.clickhouse.data.format.ClickHouseRowBinaryProcessor; import com.clickhouse.jdbc.ClickHouseDriver; import com.clickhouse.jdbc.DriverProperties; -import org.openjdk.jmh.annotations.Level; -import org.openjdk.jmh.annotations.Param; -import org.openjdk.jmh.annotations.Scope; -import org.openjdk.jmh.annotations.Setup; -import org.openjdk.jmh.annotations.State; -import org.openjdk.jmh.annotations.TearDown; +import org.openjdk.jmh.annotations.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -37,63 +28,50 @@ import java.math.BigInteger; import java.nio.ByteBuffer; import java.sql.Connection; -import java.sql.SQLException; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; -import java.util.Properties; - -import static com.clickhouse.benchmark.TestEnvironment.DB_NAME; -import static com.clickhouse.benchmark.TestEnvironment.cleanupEnvironment; -import static com.clickhouse.benchmark.TestEnvironment.getPassword; -import static com.clickhouse.benchmark.TestEnvironment.getServer; -import static com.clickhouse.benchmark.TestEnvironment.getUsername; -import static com.clickhouse.benchmark.TestEnvironment.isCloud; -import static com.clickhouse.benchmark.TestEnvironment.setupEnvironment; +import java.util.*; + +import static com.clickhouse.benchmark.TestEnvironment.*; @State(Scope.Benchmark) public class BenchmarkBase { private static final Logger LOGGER = LoggerFactory.getLogger(BenchmarkBase.class); protected ClickHouseClient clientV1; protected Client clientV2; - protected static Connection jdbcV1; - protected static Connection jdbcV2; + protected ClickHouseClient clientV1Compressed; + protected Client clientV2Compressed; + protected Connection jdbcV1RowBinary; + protected Connection jdbcV2RowBinary; + protected Connection jdbcV1Compressed; + protected Connection jdbcV2Compressed; + protected Connection jdbcV2CompressedText; + + private List closeables = new ArrayList<>(); @Setup(Level.Iteration) public void setUpIteration() { - LOGGER.info("BenchmarkBase::setUpIteration"); - clientV1 = getClientV1(); - clientV2 = getClientV2(); - jdbcV1 = getJdbcV1(); - jdbcV2 = getJdbcV2(); + LOGGER.info("BenchmarkBase::setUpIteration: pid: " + ProcessHandle.current().pid()); + clientV1 = getClientV1(false); + clientV2 = getClientV2IncludeDb(false); + clientV1Compressed = getClientV1(true); + clientV2Compressed = getClientV2IncludeDb(true); + jdbcV1RowBinary = getJdbcV1(false); + jdbcV2RowBinary = getJdbcV2(false, true); + jdbcV1Compressed = getJdbcV1(true); + jdbcV2Compressed = getJdbcV2(true, true); + jdbcV2CompressedText = getJdbcV2(true, false); + closeables = Arrays.asList(clientV1, clientV2, jdbcV1RowBinary, jdbcV2RowBinary, jdbcV1Compressed, + jdbcV2Compressed, jdbcV2CompressedText); } @TearDown(Level.Iteration) public void tearDownIteration() { LOGGER.info("BenchmarkBase::tearDownIteration"); - if (clientV1 != null) { - clientV1.close(); - clientV1 = null; - } - if (clientV2 != null) { - clientV2.close(); - clientV2 = null; - } - if (jdbcV1 != null) { - try { - jdbcV1.close(); - } catch (SQLException e) { - LOGGER.error(e.getMessage()); - } - jdbcV1 = null; - } - if (jdbcV2 != null) { + for (AutoCloseable closeable : closeables) { try { - jdbcV2.close(); - } catch (SQLException e) { - LOGGER.error(e.getMessage()); + closeable.close(); + } catch (Exception e) { + LOGGER.error("failed to close", e); } - jdbcV2 = null; } } @@ -154,6 +132,8 @@ public void setup(DataState dataState) { LOGGER.info("Loading data from file " + dataState.datasetSourceName + " with limit " + dataState.limit); dataState.dataSet = new FileDataSet(dataState.datasetSourceName.substring("file://".length()), dataState.limit); } + System.out.println("Dataset " + dataState.dataSet.getName() + ": " + dataState.dataSet.getSize() + " rows, " + + dataState.dataSet.getSizeInBytes(dataState.dataSet.getFormat()) + " bytes"); initializeTables(dataState); } @@ -191,7 +171,7 @@ public static List runQuery(String query) { return runQuery(query, true); } public static List runQuery(String query, boolean useDatabase) { - try (Client client = getClientV2(useDatabase)) { + try (Client client = getClientV2(useDatabase, true)) { return client.queryAll(query); } } @@ -202,7 +182,7 @@ public static void runAndSyncQuery(String query, String tableName) { public static void syncQuery(String tableName) { - if (isCloud()) { + if (isCluster()) { LOGGER.debug("Syncing: {}", tableName); runQuery(getSyncQuery(tableName)); } @@ -220,7 +200,7 @@ public static void dropTable(String tableName) { } public static void insertData(String tableName, InputStream dataStream, ClickHouseFormat format) { - try (Client client = getClientV2(); + try (Client client = getClientV2IncludeDb(false); InsertResponse ignored = client.insert(tableName, dataStream, format).get()) { syncQuery(tableName); List count = runQuery(getSelectCountQuery(tableName)); @@ -244,51 +224,55 @@ public static boolean verifyCount(String tableName, long expectedCount) { return true; } - protected static ClickHouseClient getClientV1() { + protected static ClickHouseClient getClientV1(boolean serverCompression) { // We get a new client so that closing won't affect other subsequent calls return ClickHouseClient.builder() + .option(ClickHouseClientOption.COMPRESS, serverCompression) .defaultCredentials(ClickHouseCredentials.fromUserAndPassword(getUsername(), getPassword())) .nodeSelector(ClickHouseNodeSelector.of(ClickHouseProtocol.HTTP)) .build(); } - protected static Client getClientV2() { - return getClientV2(true); + protected static Client getClientV2IncludeDb(boolean serverCompression) { + return getClientV2(true, serverCompression); } - protected static Client getClientV2(boolean includeDb) { + protected static Client getClientV2(boolean includeDb, boolean serverCompression) { ClickHouseNode node = getServer(); //We get a new client so that closing won't affect other subsequent calls return new Client.Builder() - .addEndpoint(Protocol.HTTP, node.getHost(), node.getPort(), isCloud()) + .addEndpoint(Protocol.HTTP, node.getHost(), node.getPort(), isSsl()) .setUsername(getUsername()) .setPassword(getPassword()) .setMaxRetries(0) + .compressServerResponse(serverCompression) .setDefaultDatabase(includeDb ? DB_NAME : "default") .build(); } - private static String jdbcURLV1(boolean isCloud) { + private static String jdbcURLV1(boolean ssl) { ClickHouseNode node = getServer(); - if (isCloud) { + if (ssl) { return String.format("jdbc:clickhouse://%s:%s?clickhouse.jdbc.v1=true&ssl=true", node.getHost(), node.getPort()); } else return String.format("jdbc:clickhouse://%s:%s?clickhouse.jdbc.v1=true", node.getHost(), node.getPort()); } - private static String jdbcURLV2(boolean isCloud) { + private static String jdbcURLV2(boolean ssl) { ClickHouseNode node = getServer(); - if (isCloud) { + if (ssl) { return String.format("jdbc:clickhouse:https://%s:%s?ssl=true", node.getHost(), node.getPort()); } else return String.format("jdbc:clickhouse://%s:%s", node.getHost(), node.getPort()); } - protected static Connection getJdbcV1() { + protected static Connection getJdbcV1(boolean isCompressed) { Properties properties = new Properties(); properties.put(ClickHouseDefaults.USER.getKey(), getUsername()); properties.put(ClickHouseDefaults.PASSWORD.getKey(), getPassword()); properties.put(ClickHouseDefaults.DATABASE.getKey(), DB_NAME); - + if (isCompressed) { + properties.put(ClickHouseClientOption.DECOMPRESS.getKey(), "1"); + } Connection jdbcV1 = null; - String jdbcURL = jdbcURLV1(isCloud()); + String jdbcURL = jdbcURLV1(isSsl()); LOGGER.warn("JDBC URL V1: " + jdbcURL); try { jdbcV1 = new ClickHouseDriver().connect(jdbcURL, properties); @@ -298,15 +282,20 @@ protected static Connection getJdbcV1() { return jdbcV1; } - protected static Connection getJdbcV2() { + protected static Connection getJdbcV2(boolean isCompressed, boolean isRowBinary) { Properties properties = new Properties(); properties.put(ClientConfigProperties.USER.getKey(), getUsername()); properties.put(ClientConfigProperties.PASSWORD.getKey(), getPassword()); - properties.put(DriverProperties.BETA_ROW_BINARY_WRITER.getKey(), "true"); + properties.put(DriverProperties.BETA_ROW_BINARY_WRITER.getKey(), String.valueOf(isRowBinary)); + if (isCompressed) { + properties.put(ClientConfigProperties.COMPRESS_CLIENT_REQUEST.getKey(), "true"); + properties.put(ClientConfigProperties.COMPRESS_SERVER_RESPONSE.getKey(), "true"); + } + properties.put(ClientConfigProperties.CLIENT_NETWORK_BUFFER_SIZE.getKey(), String.valueOf(10 * 1024 * 1024)); properties.put(ClientConfigProperties.DATABASE.getKey(), DB_NAME); Connection jdbcV2 = null; - String jdbcURL = jdbcURLV2(isCloud()); + String jdbcURL = jdbcURLV2(isSsl()); LOGGER.warn("JDBC URL V2: " + jdbcURL); try { @@ -322,7 +311,7 @@ protected static Connection getJdbcV2() { public static void loadClickHouseRecords(DataState dataState) { syncQuery(dataState.tableNameFilled); - try (ClickHouseClient clientV1 = getClientV1(); + try (ClickHouseClient clientV1 = getClientV1(false); ClickHouseResponse response = clientV1.read(getServer()) .query(getSelectQuery(dataState.tableNameFilled)) .format(ClickHouseFormat.RowBinaryWithNamesAndTypes) diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/ConcurrentInsertClient.java b/performance/src/main/java/com/clickhouse/benchmark/clients/ConcurrentInsertClient.java index 3352d6569..28ae94e35 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/ConcurrentInsertClient.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/ConcurrentInsertClient.java @@ -33,8 +33,8 @@ public class ConcurrentInsertClient extends BenchmarkBase { private static Client clientV2Shared; @Setup(Level.Trial) public void setUpIteration() { - clientV1Shared = getClientV1(); - clientV2Shared = getClientV2(); + clientV1Shared = getClientV1(false); + clientV2Shared = getClientV2IncludeDb(false); } @TearDown(Level.Trial) public void tearDownIteration() { diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/ConcurrentQueryClient.java b/performance/src/main/java/com/clickhouse/benchmark/clients/ConcurrentQueryClient.java index a7e5a5cf6..c6218b1c0 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/ConcurrentQueryClient.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/ConcurrentQueryClient.java @@ -30,8 +30,8 @@ public class ConcurrentQueryClient extends BenchmarkBase { private Client clientV2Shared; @Setup(Level.Trial) public void setUpIteration() { - clientV1Shared = getClientV1(); - clientV2Shared = getClientV2(); + clientV1Shared = getClientV1(false); + clientV2Shared = getClientV2IncludeDb(false); } @TearDown(Level.Trial) diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/DataTypes.java b/performance/src/main/java/com/clickhouse/benchmark/clients/DataTypes.java index a6803a11c..a49f271da 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/DataTypes.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/DataTypes.java @@ -35,7 +35,7 @@ public class DataTypes extends BenchmarkBase { public void setUpIteration(DataState dataState) { super.setUpIteration(); - try (Client c = getClientV2(); QueryResponse r = c.query("SELECT * FROM " + dataState.tableNameFilled, new QuerySettings() + try (Client c = getClientV2IncludeDb(false); QueryResponse r = c.query("SELECT * FROM " + dataState.tableNameFilled, new QuerySettings() .setFormat(ClickHouseFormat.RowBinaryWithNamesAndTypes)).get()) { dataState.datasetAsRowBinaryWithNamesAndTypes = ByteBuffer.wrap(r.getInputStream().readAllBytes()); LOGGER.info("Loaded {} from dataset", dataState.datasetAsRowBinaryWithNamesAndTypes.capacity()); diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/Deserializers.java b/performance/src/main/java/com/clickhouse/benchmark/clients/Deserializers.java index c9b86a8f4..29673e547 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/Deserializers.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/Deserializers.java @@ -33,7 +33,7 @@ public class Deserializers extends BenchmarkBase { public void setUpIteration(DataState dataState) { super.setUpIteration(); - try (Client c = getClientV2(); QueryResponse r = c.query("SELECT * FROM " + dataState.tableNameFilled, new QuerySettings() + try (Client c = getClientV2IncludeDb(false); QueryResponse r = c.query("SELECT * FROM " + dataState.tableNameFilled, new QuerySettings() .setFormat(ClickHouseFormat.RowBinaryWithNamesAndTypes)).get()){ dataState.datasetAsRowBinaryWithNamesAndTypes = ByteBuffer.wrap(r.getInputStream().readAllBytes()); LOGGER.info("Loaded {} from dataset", dataState.datasetAsRowBinaryWithNamesAndTypes.capacity()); diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/InsertClient.java b/performance/src/main/java/com/clickhouse/benchmark/clients/InsertClient.java index 100120691..928829e23 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/InsertClient.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/InsertClient.java @@ -9,11 +9,7 @@ import com.clickhouse.data.ClickHouseFormat; import com.clickhouse.data.ClickHouseRecord; import com.clickhouse.data.ClickHouseSerializer; -import org.openjdk.jmh.annotations.Benchmark; -import org.openjdk.jmh.annotations.Level; -import org.openjdk.jmh.annotations.Scope; -import org.openjdk.jmh.annotations.State; -import org.openjdk.jmh.annotations.TearDown; +import org.openjdk.jmh.annotations.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -29,6 +25,7 @@ public class InsertClient extends BenchmarkBase { public void verifyRowsInsertedAndCleanup(DataState dataState) { boolean success; int count = 0; + do { success = verifyCount(dataState.tableNameEmpty, dataState.dataSet.getSize()); if (!success) { @@ -47,7 +44,7 @@ public void verifyRowsInsertedAndCleanup(DataState dataState) { truncateTable(dataState.tableNameEmpty); } - @Benchmark +// @Benchmark public void insertV1(DataState dataState) { try { ClickHouseFormat format = dataState.dataSet.getFormat(); @@ -68,7 +65,7 @@ public void insertV1(DataState dataState) { } } - @Benchmark +// @Benchmark public void insertV2(DataState dataState) { try { ClickHouseFormat format = dataState.dataSet.getFormat(); @@ -85,7 +82,7 @@ public void insertV2(DataState dataState) { } } - @Benchmark +// @Benchmark public void insertV1Compressed(DataState dataState) { try { ClickHouseFormat format = dataState.dataSet.getFormat(); @@ -107,7 +104,7 @@ public void insertV1Compressed(DataState dataState) { } } - @Benchmark +// @Benchmark public void insertV2Compressed(DataState dataState) { try { ClickHouseFormat format = dataState.dataSet.getFormat(); @@ -125,7 +122,7 @@ public void insertV2Compressed(DataState dataState) { } } - @Benchmark +// @Benchmark public void insertV1RowBinary(DataState dataState) { try { ClickHouseFormat format = ClickHouseFormat.RowBinary; @@ -175,4 +172,55 @@ public void insertV2RowBinary(DataState dataState) { } } + @Benchmark + public void insertV2RowBinaryWithHttpCompression(DataState dataState) { + try { + final ClickHouseFormat format = ClickHouseFormat.RowBinary; + try (InsertResponse response = clientV2.insert(dataState.tableNameEmpty, out -> { + RowBinaryFormatWriter w = new RowBinaryFormatWriter(out, dataState.dataSet.getSchema(), format); + for (List row : dataState.dataSet.getRowsOrdered()) { + int index = 1; + for (Object value : row) { + w.setValue(index, value); + index++; + } + w.commitRow(); + } + out.flush(); + + }, format, new InsertSettings() + .compressClientRequest(true) + .useHttpCompression(true)).get()) { + response.getWrittenRows(); + } + } catch (Exception e) { + LOGGER.error("Error: ", e); + } + } + + @Benchmark + public void insertV2RowBinaryWithCompression(DataState dataState) { + try { + final ClickHouseFormat format = ClickHouseFormat.RowBinary; + try (InsertResponse response = clientV2.insert(dataState.tableNameEmpty, out -> { + RowBinaryFormatWriter w = new RowBinaryFormatWriter(out, dataState.dataSet.getSchema(), format); + for (List row : dataState.dataSet.getRowsOrdered()) { + int index = 1; + for (Object value : row) { + w.setValue(index, value); + index++; + } + w.commitRow(); + } + out.flush(); + + }, format, new InsertSettings() + .compressClientRequest(true)).get()) { + response.getWrittenRows(); + } + } catch (Exception e) { + LOGGER.error("Error: ", e); + } + } + } diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/JDBCInsert.java b/performance/src/main/java/com/clickhouse/benchmark/clients/JDBCInsert.java index c74fd2c6b..54364ac45 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/JDBCInsert.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/JDBCInsert.java @@ -51,13 +51,27 @@ void insetData(Connection connection, DataState dataState) throws SQLException { } @Benchmark - public void insertJDBCV1(DataState dataState) throws SQLException { - insetData(jdbcV1, dataState); + public void insertJDBCRowBinaryV1(DataState dataState) throws SQLException { + insetData(jdbcV1RowBinary, dataState); } @Benchmark - public void insertJDBCV2(DataState dataState) throws SQLException { - insetData(jdbcV2, dataState); + public void insertJDBCRowBinaryV2(DataState dataState) throws SQLException { + insetData(jdbcV2RowBinary, dataState); } + @Benchmark + public void insertJDBCRowBinaryCompressedV1(DataState dataState) throws SQLException { + insetData(jdbcV1Compressed, dataState); + } + + @Benchmark + public void insertJDBCRowBinaryCompressedV2(DataState dataState) throws SQLException { + insetData(jdbcV2Compressed, dataState); + } + + @Benchmark + public void insertJDBCTextCompressedV2(DataState dataState) throws SQLException { + insetData(jdbcV2CompressedText, dataState); + } } \ No newline at end of file diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/JDBCQuery.java b/performance/src/main/java/com/clickhouse/benchmark/clients/JDBCQuery.java index c1c7e3503..731fb7a5f 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/JDBCQuery.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/JDBCQuery.java @@ -29,12 +29,22 @@ void selectData(Connection connection, DataState dataState, Blackhole blackhole) @Benchmark public void selectJDBCV1(DataState dataState, Blackhole blackhole) throws SQLException { - selectData(jdbcV1, dataState, blackhole); + selectData(jdbcV1RowBinary, dataState, blackhole); } @Benchmark public void selectJDBCV2(DataState dataState, Blackhole blackhole) throws SQLException { - selectData(jdbcV2, dataState, blackhole); + selectData(jdbcV2RowBinary, dataState, blackhole); + } + + @Benchmark + public void selectJDBCV1Compressed(DataState dataState, Blackhole blackhole) throws SQLException { + selectData(jdbcV1Compressed, dataState, blackhole); + } + + @Benchmark + public void selectJDBCV2Compressed(DataState dataState, Blackhole blackhole) throws SQLException { + selectData(jdbcV2Compressed, dataState, blackhole); } void selectDataUseNames(Connection connection, DataState dataState, Blackhole blackhole) throws SQLException { @@ -43,7 +53,32 @@ void selectDataUseNames(Connection connection, DataState dataState, Blackhole bl ResultSet rs = stmt.executeQuery(sql)) { while (rs.next()) { for (ClickHouseColumn col : dataState.dataSet.getSchema().getColumns()) { - blackhole.consume(rs.getObject(col.getColumnName())); + switch (col.getDataType()) { + case Int8: + case UInt8: + case Int16: + case UInt16: + case Int32: + case UInt32: + blackhole.consume(rs.getLong(col.getColumnName())); + break; + case Date: + case Date32: + blackhole.consume(rs.getDate(col.getColumnName())); + break; + case DateTime: + case DateTime32: + case DateTime64: + // slower than just getObject + blackhole.consume(rs.getTimestamp(col.getColumnName())); + break; + case String: + case FixedString: + blackhole.consume(rs.getString(col.getColumnName())); + break; + default: + blackhole.consume(rs.getObject(col.getColumnName())); + } } } } @@ -51,11 +86,12 @@ void selectDataUseNames(Connection connection, DataState dataState, Blackhole bl @Benchmark public void selectJDBCV1UseNames(DataState dataState, Blackhole blackhole) throws SQLException { - selectDataUseNames(jdbcV1, dataState, blackhole); + selectDataUseNames(jdbcV1RowBinary, dataState, blackhole); } @Benchmark public void selectJDBCV2UseName(DataState dataState, Blackhole blackhole) throws SQLException { - selectDataUseNames(jdbcV2, dataState, blackhole); + selectDataUseNames(jdbcV2RowBinary, dataState, blackhole); } + } \ No newline at end of file diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/MixedWorkload.java b/performance/src/main/java/com/clickhouse/benchmark/clients/MixedWorkload.java index 010bde562..d36045ea3 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/MixedWorkload.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/MixedWorkload.java @@ -39,8 +39,8 @@ public class MixedWorkload extends BenchmarkBase { private Client clientV2Shared; @Setup(Level.Trial) public void setUpTrial() { - clientV1Shared = getClientV1(); - clientV2Shared = getClientV2(); + clientV1Shared = getClientV1(false); + clientV2Shared = getClientV2IncludeDb(false); } @TearDown(Level.Trial) diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/Serializers.java b/performance/src/main/java/com/clickhouse/benchmark/clients/Serializers.java index e4a57db82..2b263f537 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/Serializers.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/Serializers.java @@ -23,7 +23,7 @@ public void SerializerOutputStreamV1(DataState dataState, Blackhole blackhole) { try { ClickHouseOutputStream chos = ClickHouseOutputStream.of(empty); ClickHouseDataProcessor p = dataState.dataSet.getClickHouseDataProcessor(); - ClickHouseSerializer[] serializers = p.getSerializers(getClientV1().getConfig(), p.getColumns()); + ClickHouseSerializer[] serializers = p.getSerializers(getClientV1(false).getConfig(), p.getColumns()); for (ClickHouseRecord record : dataState.dataSet.getClickHouseRecords()) { for (int i = 0; i < serializers.length; i++) { serializers[i].serialize(record.getValue(i), chos); diff --git a/performance/src/main/java/com/clickhouse/benchmark/data/DataSet.java b/performance/src/main/java/com/clickhouse/benchmark/data/DataSet.java index b098c8f1e..8f734dfd0 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/data/DataSet.java +++ b/performance/src/main/java/com/clickhouse/benchmark/data/DataSet.java @@ -40,6 +40,14 @@ default InputStream getInputStream(int rowId, ClickHouseFormat format) { List getBytesList(ClickHouseFormat format); + default long getSizeInBytes(ClickHouseFormat format) { + long total = 0; + for (byte[] bytes : getBytesList(format)) { + total += bytes.length; + } + return total; + } + List> getRows(); List getClickHouseRecords(); From 7cf119c1572712c75a6219d011b98b572f7974de Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Fri, 25 Sep 2026 17:15:54 -0700 Subject: [PATCH 2/2] Added compression benchmarks --- performance/README.md | 16 +- performance/pom.xml | 20 ++ .../clickhouse/benchmark/BenchmarkRunner.java | 65 ++++- .../benchmark/clients/BenchmarkBase.java | 8 +- .../benchmark/clients/Compression.java | 225 +++++++++++++++++- 5 files changed, 326 insertions(+), 8 deletions(-) diff --git a/performance/README.md b/performance/README.md index 434f31d8f..369f875d8 100644 --- a/performance/README.md +++ b/performance/README.md @@ -71,9 +71,21 @@ Other options: - "q" - QueryClient - query operation benchmarks - "ci" - ConcurrentInsertClient - concurrent version of insert benchmarks - "cq" - ConcurrentQueryClient - concurrent version of query benchmarks - - "lz" - Compression - compression related benchmarks + - "lz" - Compression - LZ4 output stream benchmarks (no server involved) + - "comp" - Compression - query/insert compression matrix of clients, methods, algorithms and formats (see below) - "writer" - Serializer - serialization only logic benchmarks - "reader" - DeSerilalizer - deserialization only logic benchmarks - "mixed" - MixedWorkload - "jq" - JDBCQuery - query operations using JDBC - - "ji" - JDBCInsert - insert operation using JDBC \ No newline at end of file + - "ji" - JDBCInsert - insert operation using JDBC + +Compression matrix filters (used with `-b comp` or `-b all`, each defaults to all values): +- "-cc" - clients: `v1,v2` +- "-cm" - compression methods: `http` (`Content-Encoding`/`Accept-Encoding`), `native` (ClickHouse block compression, `use_http_compression = false`) +- "-ca" - algorithms: `lz4,zstd,snappy,brotli` +- "-cf" - formats: `RowBinaryWithNamesAndTypes,JSONEachRow` (any ClickHouse format name, e.g. `RowBinary`, is accepted) + +Ex.: `-b comp -cc v2 -cm http -ca zstd,lz4 -cf JSONEachRow`. + +Combinations a client doesn't support are skipped: V1 sends LZ4 only natively and other algorithms only over HTTP, +V2 native mode is LZ4 only, V2 can't compress requests with brotli, and snappy works with neither client. \ No newline at end of file diff --git a/performance/pom.xml b/performance/pom.xml index 8b91cf3f6..68cf5c785 100644 --- a/performance/pom.xml +++ b/performance/pom.xml @@ -18,6 +18,9 @@ 0.11.0-rc1-SNAPSHOT 1.37 2.0.2 + 1.5.7-6 + 0.1.2 + 1.12.0 3.1.0 3.6.0 @@ -94,6 +97,23 @@ all + + + com.github.luben + zstd-jni + ${zstd-jni.version} + + + org.brotli + dec + ${brotli.version} + + + com.aayushatharva.brotli4j + brotli4j + ${brotli4j.version} + + diff --git a/performance/src/main/java/com/clickhouse/benchmark/BenchmarkRunner.java b/performance/src/main/java/com/clickhouse/benchmark/BenchmarkRunner.java index ea58ca5b6..e18639241 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/BenchmarkRunner.java +++ b/performance/src/main/java/com/clickhouse/benchmark/BenchmarkRunner.java @@ -10,6 +10,7 @@ import com.clickhouse.benchmark.clients.MixedWorkload; import com.clickhouse.benchmark.clients.QueryClient; import com.clickhouse.benchmark.clients.Serializers; +import org.openjdk.jmh.annotations.Benchmark; import org.openjdk.jmh.annotations.Mode; import org.openjdk.jmh.profile.GCProfiler; import org.openjdk.jmh.profile.MemPoolProfiler; @@ -21,12 +22,20 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.lang.reflect.Method; +import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; +import java.util.List; import java.util.Map; +import java.util.Set; import java.util.SortedSet; import java.util.TreeSet; import java.util.concurrent.TimeUnit; +import java.util.regex.Matcher; +import java.util.regex.Pattern; import static com.clickhouse.benchmark.TestEnvironment.isRemote; @@ -77,7 +86,8 @@ public static void main(String[] args) throws Exception { String[] testMaskParts = testMask.split(","); SortedSet benchmarks = new TreeSet<>(); - if (testMaskParts[0].equalsIgnoreCase("all")) { + boolean runAll = testMaskParts[0].equalsIgnoreCase("all"); + if (runAll) { BENCHMARK_FLAGS.values().forEach((b) -> { optBuilder.include(b); benchmarks.add(b); @@ -92,6 +102,20 @@ public static void main(String[] args) throws Exception { } } + if (runAll || Arrays.asList(testMaskParts).contains(COMPRESSION_MATRIX_FLAG)) { + List matrix = compressionMatrix(options); + if (matrix.isEmpty()) { + System.out.println("No compression benchmark matches the selected clients/methods/algorithms"); + } + for (String benchmark : matrix) { + optBuilder.include(Pattern.quote(benchmark) + "$"); + benchmarks.add(benchmark); + } + if (options.containsKey("-cf")) { + optBuilder.param("format", options.get("-cf").split(",")); + } + } + System.out.println("Running benchmarks: " + benchmarks); new Runner(optBuilder.build()).run(); } @@ -104,7 +128,7 @@ private static Map buildBenchmarkFlags() { map.put("i", InsertClient.class.getName()); map.put("cq", ConcurrentQueryClient.class.getName()); map.put("ci", ConcurrentInsertClient.class.getName()); - map.put("lz", Compression.class.getName()); + map.put("lz", Compression.class.getName() + ".CompressingOutputStream"); map.put("reader", Deserializers.class.getName()); map.put("writer", Serializers.class.getName()); map.put("mixed", MixedWorkload.class.getName()); @@ -113,6 +137,43 @@ private static Map buildBenchmarkFlags() { return map; } + private static final String COMPRESSION_MATRIX_FLAG = "comp"; + + private static final Pattern COMPRESSION_BENCHMARK_NAME = + Pattern.compile("(query|insert)(V1|V2)(Native|Http)(Lz4|Zstd|Snappy|Brotli)"); + + /** + * Selects {@link Compression} matrix benchmarks by client ({@code -cc v1,v2}), compression method + * ({@code -cm http,native}) and algorithm ({@code -ca lz4,zstd,snappy,brotli}). Each filter defaults to all + * values. Combinations a client does not support have no benchmark method and are skipped. + */ + private static List compressionMatrix(Map options) { + Set clients = filterValues(options, "-cc", "v1,v2"); + Set methods = filterValues(options, "-cm", "http,native"); + Set algorithms = filterValues(options, "-ca", "lz4,zstd,snappy,brotli"); + + List selected = new ArrayList<>(); + for (Method method : Compression.class.getMethods()) { + Matcher m = COMPRESSION_BENCHMARK_NAME.matcher(method.getName()); + if (method.isAnnotationPresent(Benchmark.class) && m.matches() + && clients.contains(m.group(2).toLowerCase()) + && methods.contains(m.group(3).toLowerCase()) + && algorithms.contains(m.group(4).toLowerCase())) { + selected.add(Compression.class.getName() + "." + method.getName()); + } + } + Collections.sort(selected); + return selected; + } + + private static Set filterValues(Map options, String key, String defaultValues) { + Set values = new HashSet<>(); + for (String value : options.getOrDefault(key, defaultValues).split(",")) { + values.add(value.trim().toLowerCase()); + } + return values; + } + private static Map parseArgs(String[] args) { Map options = new HashMap<>(); for (int i = 0; i < args.length; i+=2) { diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/BenchmarkBase.java b/performance/src/main/java/com/clickhouse/benchmark/clients/BenchmarkBase.java index 4bfe2d1bb..30512c17f 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/BenchmarkBase.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/BenchmarkBase.java @@ -236,16 +236,18 @@ protected static Client getClientV2IncludeDb(boolean serverCompression) { return getClientV2(true, serverCompression); } protected static Client getClientV2(boolean includeDb, boolean serverCompression) { - ClickHouseNode node = getServer(); //We get a new client so that closing won't affect other subsequent calls + return getClientV2Builder(includeDb, serverCompression).build(); + } + protected static Client.Builder getClientV2Builder(boolean includeDb, boolean serverCompression) { + ClickHouseNode node = getServer(); return new Client.Builder() .addEndpoint(Protocol.HTTP, node.getHost(), node.getPort(), isSsl()) .setUsername(getUsername()) .setPassword(getPassword()) .setMaxRetries(0) .compressServerResponse(serverCompression) - .setDefaultDatabase(includeDb ? DB_NAME : "default") - .build(); + .setDefaultDatabase(includeDb ? DB_NAME : "default"); } private static String jdbcURLV1(boolean ssl) { ClickHouseNode node = getServer(); diff --git a/performance/src/main/java/com/clickhouse/benchmark/clients/Compression.java b/performance/src/main/java/com/clickhouse/benchmark/clients/Compression.java index 016724315..657dde518 100644 --- a/performance/src/main/java/com/clickhouse/benchmark/clients/Compression.java +++ b/performance/src/main/java/com/clickhouse/benchmark/clients/Compression.java @@ -1,24 +1,60 @@ package com.clickhouse.benchmark.clients; import com.clickhouse.benchmark.data.DataSet; +import com.clickhouse.client.ClickHouseClient; +import com.clickhouse.client.ClickHouseResponse; +import com.clickhouse.client.api.Client; +import com.clickhouse.client.api.insert.InsertResponse; +import com.clickhouse.client.api.insert.InsertSettings; import com.clickhouse.client.api.internal.ClickHouseLZ4OutputStream; +import com.clickhouse.client.api.query.QueryResponse; +import com.clickhouse.client.api.query.QuerySettings; +import com.clickhouse.client.config.ClickHouseClientOption; import com.clickhouse.client.internal.net.jpountz.lz4.LZ4Factory; +import com.clickhouse.data.ClickHouseCompression; +import com.clickhouse.data.ClickHouseFormat; import com.clickhouse.data.ClickHouseOutputStream; import com.clickhouse.data.stream.Lz4OutputStream; import org.openjdk.jmh.annotations.Benchmark; import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Param; +import org.openjdk.jmh.annotations.Scope; import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.TearDown; +import org.openjdk.jmh.infra.Blackhole; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; +import java.io.InputStream; +import static com.clickhouse.benchmark.TestEnvironment.getServer; + +/** + * Compares compression setups of client V1 and V2 for queries (server compresses the response) and inserts + * (client compresses the request). + * + *

Matrix benchmarks are named {@code }. {@code Native} is + * ClickHouse block compression ({@code compress}/{@code decompress} parameters), {@code Http} is compression + * negotiated with {@code Content-Encoding}/{@code Accept-Encoding}. Only combinations supported by a client + * have a method: V1 always sends LZ4 natively and any other algorithm over HTTP, V2 native mode is LZ4 only + * and V2 cannot compress requests with brotli. Snappy does not work with either client.

+ * + *

Payloads are raw bytes of the tested format, so the results reflect transport and compression cost + * rather than (de)serialization.

+ */ public class Compression extends BenchmarkBase { private static final Logger LOGGER = LoggerFactory.getLogger(Compression.class); static final int COMPRESS_BUFFER_SIZE = 64 * 1024; // 64K private static final LZ4Factory factory = LZ4Factory.fastestInstance(); - @Setup(Level.Invocation) + + // Newer servers default to ZSTD for native compression, while both clients decode only LZ4 blocks. + private static final String NATIVE_METHOD_SETTING = "network_compression_method"; + + @Setup(Level.Trial) public void setup() { LOGGER.info("Compressor type {}", factory.fastCompressor()); } @@ -51,4 +87,191 @@ public void CompressingOutputStreamV2(DataState dataState) { LOGGER.error("Error: ", e); } } + + @State(Scope.Benchmark) + public static class CompressionState { + @Param({"RowBinaryWithNamesAndTypes", "JSONEachRow"}) + String format; + + ClickHouseFormat clickHouseFormat; + byte[] payload; + ClickHouseClient clientV1; + Client clientV2Native; + Client clientV2Http; + + // Iteration level so it runs after the environment and tables are created at trial level. + @Setup(Level.Iteration) + public void setup(DataState dataState) throws Exception { + if (payload != null) { + return; + } + clickHouseFormat = ClickHouseFormat.valueOf(format); + try (Client client = getClientV2IncludeDb(false); + QueryResponse response = client.query(getSelectQuery(dataState.tableNameFilled), + new QuerySettings().setFormat(clickHouseFormat)).get()) { + payload = response.getInputStream().readAllBytes(); + } + LOGGER.info("Payload in {}: {} bytes", format, payload.length); + + clientV1 = getClientV1(false); + clientV2Native = getClientV2Builder(true, true) + .useHttpCompression(false) + .serverSetting(NATIVE_METHOD_SETTING, "lz4") + .build(); + clientV2Http = getClientV2Builder(true, true) + .useHttpCompression(true) + .build(); + } + + @TearDown(Level.Iteration) + public void truncate(DataState dataState) { + truncateTable(dataState.tableNameEmpty); + } + + @TearDown(Level.Trial) + public void tearDown() { + if (clientV1 != null) { + clientV1.close(); + } + if (clientV2Native != null) { + clientV2Native.close(); + } + if (clientV2Http != null) { + clientV2Http.close(); + } + } + } + + @Benchmark + public void queryV1NativeLz4(DataState dataState, CompressionState state, Blackhole blackhole) throws Exception { + queryV1(dataState, state, ClickHouseCompression.LZ4, blackhole); + } + + @Benchmark + public void queryV1HttpZstd(DataState dataState, CompressionState state, Blackhole blackhole) throws Exception { + queryV1(dataState, state, ClickHouseCompression.ZSTD, blackhole); + } + + @Benchmark + public void queryV1HttpBrotli(DataState dataState, CompressionState state, Blackhole blackhole) throws Exception { + queryV1(dataState, state, ClickHouseCompression.BROTLI, blackhole); + } + + @Benchmark + public void insertV1NativeLz4(DataState dataState, CompressionState state) throws Exception { + insertV1(dataState, state, ClickHouseCompression.LZ4); + } + + @Benchmark + public void insertV1HttpZstd(DataState dataState, CompressionState state) throws Exception { + insertV1(dataState, state, ClickHouseCompression.ZSTD); + } + + @Benchmark + public void insertV1HttpBrotli(DataState dataState, CompressionState state) throws Exception { + insertV1(dataState, state, ClickHouseCompression.BROTLI); + } + + @Benchmark + public void queryV2NativeLz4(DataState dataState, CompressionState state, Blackhole blackhole) throws Exception { + queryV2(dataState, state, state.clientV2Native, null, blackhole); + } + + @Benchmark + public void queryV2HttpLz4(DataState dataState, CompressionState state, Blackhole blackhole) throws Exception { + queryV2(dataState, state, state.clientV2Http, ClickHouseCompression.LZ4, blackhole); + } + + @Benchmark + public void queryV2HttpZstd(DataState dataState, CompressionState state, Blackhole blackhole) throws Exception { + queryV2(dataState, state, state.clientV2Http, ClickHouseCompression.ZSTD, blackhole); + } + + @Benchmark + public void queryV2HttpBrotli(DataState dataState, CompressionState state, Blackhole blackhole) throws Exception { + queryV2(dataState, state, state.clientV2Http, ClickHouseCompression.BROTLI, blackhole); + } + + @Benchmark + public void insertV2NativeLz4(DataState dataState, CompressionState state) throws Exception { + insertV2(dataState, state, state.clientV2Native, null); + } + + @Benchmark + public void insertV2HttpLz4(DataState dataState, CompressionState state) throws Exception { + insertV2(dataState, state, state.clientV2Http, ClickHouseCompression.LZ4); + } + + @Benchmark + public void insertV2HttpZstd(DataState dataState, CompressionState state) throws Exception { + insertV2(dataState, state, state.clientV2Http, ClickHouseCompression.ZSTD); + } + + private static void queryV1(DataState dataState, CompressionState state, ClickHouseCompression algorithm, + Blackhole blackhole) throws Exception { + try (ClickHouseResponse response = state.clientV1.read(getServer()) + .query(getSelectQuery(dataState.tableNameFilled)) + .format(state.clickHouseFormat) + .option(ClickHouseClientOption.ASYNC, false) + .option(ClickHouseClientOption.COMPRESS, true) + .option(ClickHouseClientOption.COMPRESS_ALGORITHM, algorithm) + .set(NATIVE_METHOD_SETTING, "lz4") + .executeAndWait()) { + blackhole.consume(drain(response.getInputStream())); + } + } + + private static void insertV1(DataState dataState, CompressionState state, ClickHouseCompression algorithm) + throws Exception { + try (ClickHouseResponse response = state.clientV1.read(getServer()) + .write() + .query(getInsertQuery(dataState.tableNameEmpty)) + .format(state.clickHouseFormat) + .option(ClickHouseClientOption.ASYNC, false) + .option(ClickHouseClientOption.DECOMPRESS, true) + .option(ClickHouseClientOption.DECOMPRESS_ALGORITHM, algorithm) + .data(new ByteArrayInputStream(state.payload)) + .executeAndWait()) { + response.getSummary(); + } + } + + /** + * @param httpAlgorithm encoding to request with {@code Accept-Encoding}, {@code null} for native compression + */ + private static void queryV2(DataState dataState, CompressionState state, Client client, + ClickHouseCompression httpAlgorithm, Blackhole blackhole) throws Exception { + QuerySettings settings = new QuerySettings().setFormat(state.clickHouseFormat); + if (httpAlgorithm != null) { + settings.httpHeader("Accept-Encoding", httpAlgorithm.encoding()); + } + try (QueryResponse response = client.query(getSelectQuery(dataState.tableNameFilled), settings).get()) { + blackhole.consume(drain(response.getInputStream())); + } + } + + /** + * @param httpAlgorithm encoding to send with {@code Content-Encoding}, {@code null} for native compression + */ + private static void insertV2(DataState dataState, CompressionState state, Client client, + ClickHouseCompression httpAlgorithm) throws Exception { + InsertSettings settings = new InsertSettings().compressClientRequest(true); + if (httpAlgorithm != null) { + settings.httpHeader("Content-Encoding", httpAlgorithm.encoding()); + } + try (InsertResponse response = client.insert(dataState.tableNameEmpty, + new ByteArrayInputStream(state.payload), state.clickHouseFormat, settings).get()) { + response.getWrittenRows(); + } + } + + private static long drain(InputStream in) throws Exception { + byte[] buffer = new byte[COMPRESS_BUFFER_SIZE]; + long total = 0; + int read; + while ((read = in.read(buffer)) != -1) { + total += read; + } + return total; + } }