Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,10 @@

### New Features

- **[client-v2,jdbc-v2]** Added ZSTD compression support for Block compression stream. Previously only
LZ4 was supported in this case. Note: Added ZSTD library and native libraries to `-all` JDBC package because it is
now required to work with server. (https://github.com/ClickHouse/clickhouse-java/issues/3105).

- **[migration-helpers]** Added `migration-helpers` module containing `ConfigurationMigrationHelper` and
`ConfigPropertyCache` to convert configuration properties and connection URLs from v1 (0.7.1) format to v2 (0.9.8+)
format (automatically prefixing ClickHouse server settings with `clickhouse_setting_`, custom headers with
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -160,25 +160,99 @@
protected abstract Class<? extends ClickHouseClient> getClientClass();

protected Map<ClickHouseOption, Serializable> getClientOptions() {
return Collections.emptyMap();
return Collections.singletonMap(ClickHouseClientOption.CUSTOM_SETTINGS, "network_compression_method=lz4");
}

protected ClickHouseClientBuilder initClient(ClickHouseClientBuilder builder) {
return builder;
}

protected ClickHouseClient getClient(ClickHouseConfig... configs) {
return initClient(ClickHouseClient.builder()).config(new ClickHouseConfig(configs))
Map<ClickHouseOption, Serializable> defaultOptions = new HashMap<>(getClientOptions());
ClickHouseConfig baseConfig = new ClickHouseConfig(defaultOptions);
List<ClickHouseConfig> list = new ArrayList<>();
list.add(baseConfig);
if (configs != null) {
Collections.addAll(list, configs);
}
return initClient(ClickHouseClient.builder().options(defaultOptions)).config(new ClickHouseConfig(list))
.nodeSelector(ClickHouseNodeSelector.of(getProtocol())).build();
}

protected ClickHouseClient getSecureClient(ClickHouseConfig... configs) {

Check warning on line 182 in clickhouse-client/src/test/java/com/clickhouse/client/ClientIntegrationTest.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Update this method so that its implementation is not identical to "getClient" on line 170.

See more on https://sonarcloud.io/project/issues?id=ClickHouse_clickhouse-java&issues=AaDrnYxZxBtXXZoTJdo5&open=AaDrnYxZxBtXXZoTJdo5&pullRequest=3157
return initClient(ClickHouseClient.builder())
.config(new ClickHouseConfig(configs))
Map<ClickHouseOption, Serializable> defaultOptions = new HashMap<>(getClientOptions());
ClickHouseConfig baseConfig = new ClickHouseConfig(defaultOptions);
List<ClickHouseConfig> list = new ArrayList<>();
list.add(baseConfig);
if (configs != null) {
Collections.addAll(list, configs);
}
return initClient(ClickHouseClient.builder().options(defaultOptions))
.config(new ClickHouseConfig(list))
.nodeSelector(ClickHouseNodeSelector.of(getProtocol()))
.build();
}

private ClickHouseNode addCustomSettings(ClickHouseNode node) {
if (node == null) {
return null;
}
String key = ClickHouseClientOption.CUSTOM_SETTINGS.getKey();
String setting = "network_compression_method=lz4";
String existing = node.getOptions().get(key);
if (existing != null && !existing.isEmpty()) {
if (!existing.contains("network_compression_method")) {
setting = existing + "," + setting;
} else {
setting = existing;
}
}
String httpKey = "custom_http_params";
String httpSetting = "network_compression_method=lz4";
String httpExisting = node.getOptions().get(httpKey);
if (httpExisting != null && !httpExisting.isEmpty()) {
if (!httpExisting.contains("network_compression_method")) {
httpSetting = httpExisting + "," + httpSetting;
} else {
httpSetting = httpExisting;
}
}
return ClickHouseNode.builder(node)
.addOption(key, setting)
.addOption(httpKey, httpSetting)
.build();
}

@Override
protected ClickHouseNode getSecureServer(ClickHouseProtocol protocol) {
return addCustomSettings(super.getSecureServer(protocol));
}

@Override
protected ClickHouseNode getSecureServer(ClickHouseProtocol protocol, ClickHouseNode base) {
return addCustomSettings(super.getSecureServer(protocol, base));
}

@Override
protected ClickHouseNode getServer(ClickHouseProtocol protocol) {
return addCustomSettings(super.getServer(protocol));
}

@Override
protected ClickHouseNode getServer(ClickHouseProtocol protocol, ClickHouseNode base) {
return addCustomSettings(super.getServer(protocol, base));
}

@Override
protected ClickHouseNode getServer(ClickHouseProtocol protocol, int port) {
return addCustomSettings(super.getServer(protocol, port));
}

@Override
protected ClickHouseNode getServer(ClickHouseProtocol protocol, Map<String, String> options) {
return addCustomSettings(super.getServer(protocol, options));
}

protected ClickHouseNode getSecureServer(ClickHouseNode base) {
return getSecureServer(getProtocol(), base);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,8 +74,9 @@ public boolean supports(Class<?> clazz) {

@Override
protected Map<ClickHouseOption, Serializable> getClientOptions() {
return Collections.singletonMap(ClickHouseHttpOption.CONNECTION_PROVIDER,
HttpConnectionProvider.APACHE_HTTP_CLIENT);
Map<ClickHouseOption, Serializable> options = new HashMap<>(super.getClientOptions());
options.put(ClickHouseHttpOption.CONNECTION_PROVIDER, HttpConnectionProvider.APACHE_HTTP_CLIENT);
return options;
}

@Test(groups = { "unit" }, dataProvider = "replicaTags")
Expand Down Expand Up @@ -351,7 +352,7 @@ public void testConnectionTTL(Map<ClickHouseOption, Serializable> options, int o
proxy.addStubMapping(WireMock.post(WireMock.anyUrl())
.willReturn(WireMock.aResponse().proxiedFrom(targetURI.build().toString())).build());

Map<ClickHouseOption, Serializable> baseOptions = new HashMap<>();
Map<ClickHouseOption, Serializable> baseOptions = new HashMap<>(getClientOptions());
baseOptions.put(ClickHouseClientOption.PROXY_PORT, proxyPort);
baseOptions.put(ClickHouseClientOption.PROXY_HOST, "localhost");
baseOptions.put(ClickHouseClientOption.PROXY_TYPE, ClickHouseProxyType.HTTP);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -103,8 +103,9 @@ protected Class<? extends ClickHouseClient> getClientClass() {

@Override
protected Map<ClickHouseOption, Serializable> getClientOptions() {
return Collections.singletonMap(ClickHouseHttpOption.CONNECTION_PROVIDER,
HttpConnectionProvider.HTTP_URL_CONNECTION);
Map<ClickHouseOption, Serializable> options = new HashMap<>(super.getClientOptions());
options.put(ClickHouseHttpOption.CONNECTION_PROVIDER, HttpConnectionProvider.HTTP_URL_CONNECTION);
return options;
}

@Test(groups = { "integration" })
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ public void testSetRoleDifferentConnections(String[] roles, String setRoleExpr,
properties.setProperty(ClickHouseDefaults.PASSWORD.getKey(), ClickHouseServerForTest.getPassword());
properties.setProperty(ClickHouseHttpOption.REMEMBER_LAST_SET_ROLES.getKey(), "true");
properties.setProperty(ClickHouseHttpOption.CONNECTION_PROVIDER.getKey(), connectionProvider);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, properties);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, addCustomSettings(properties));
String serverVersion = getServerVersion(dataSource.getConnection());
if (ClickHouseVersion.of(serverVersion).check("(,24.3]")) {
System.out.println("Test is skipped: feature is supported since 24.4");
Expand Down Expand Up @@ -117,7 +117,7 @@ public void testSetRolesAccessingTableRows() throws SQLException {
Properties properties = new Properties();
properties.setProperty(ClickHouseDefaults.PASSWORD.getKey(), ClickHouseServerForTest.getPassword());
properties.setProperty(ClickHouseHttpOption.REMEMBER_LAST_SET_ROLES.getKey(), "true");
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, properties);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, addCustomSettings(properties));
String serverVersion = getServerVersion(dataSource.getConnection());
if (ClickHouseVersion.of(serverVersion).check("(,24.3]")) {
System.out.println("Test is skipped: feature is supported since 24.4");
Expand Down Expand Up @@ -191,7 +191,7 @@ public void testPasswordAuthentication(String identifyWith, String identifyBy) t
String url = String.format("jdbc:ch:%s", getEndpointString());
Properties properties = new Properties();
properties.setProperty(ClickHouseHttpOption.REMEMBER_LAST_SET_ROLES.getKey(), "true");
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, properties);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, addCustomSettings(properties));

try (Connection connection = dataSource.getConnection("access_dba", "123")) {
Statement st = connection.createStatement();
Expand Down Expand Up @@ -231,7 +231,7 @@ public void testSwitchingBasicAuthToClickHouseHeaders(String identifyWith, Strin
String url = String.format("jdbc:ch:%s", getEndpointString());
Properties properties = new Properties();
properties.put(ClickHouseHttpOption.USE_BASIC_AUTHENTICATION.getKey(), false);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, properties);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, addCustomSettings(properties));

try (Connection connection = dataSource.getConnection("access_dba", "123")) {
Statement st = connection.createStatement();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
import com.clickhouse.client.BaseIntegrationTest;
import com.clickhouse.client.ClickHouseNode;
import com.clickhouse.client.ClickHouseProtocol;
import com.clickhouse.client.config.ClickHouseClientOption;
import com.clickhouse.client.http.config.ClickHouseHttpOption;

import javax.sql.DataSource;
Expand All @@ -28,6 +29,10 @@ public abstract class JdbcIntegrationTest extends BaseIntegrationTest {

protected String buildJdbcUrl(ClickHouseProtocol protocol, String prefix, String url) {
if (url != null && url.startsWith("jdbc:")) {
if (protocol != ClickHouseProtocol.MYSQL && !url.contains("custom_settings")) {
char sep = url.indexOf('?') >= 0 ? '&' : '?';
return url + sep + "custom_settings=network_compression_method=lz4";
}
return url;
}

Expand Down Expand Up @@ -58,6 +63,14 @@ protected String buildJdbcUrl(ClickHouseProtocol protocol, String prefix, String
builder.append('?').append(ClickHouseHttpOption.CONNECTION_PROVIDER.getKey()).append('=')
.append(CUSTOM_PROTOCOL_NAME);
}

if (protocol != ClickHouseProtocol.MYSQL) {
String customSetting = "network_compression_method=lz4";
if (builder.indexOf("custom_settings") == -1) {
char sep = builder.indexOf("?") >= 0 ? '&' : '?';
builder.append(sep).append("custom_settings=").append(customSetting);
Comment thread
chernser marked this conversation as resolved.
}
}
return builder.toString();
}

Expand Down Expand Up @@ -101,10 +114,31 @@ public DataSource newDataSource(String url) throws SQLException {
return newDataSource(url, new Properties());
}

public DataSource newDataSource(String url, Properties properties) throws SQLException {
protected Properties addCustomSettings(Properties properties) {
if (properties == null) {
properties = new Properties();
}
String customSettingsKey = ClickHouseClientOption.CUSTOM_SETTINGS.getKey();
String customSetting = "network_compression_method=lz4";
String existingCustom = properties.getProperty(customSettingsKey);
if (existingCustom == null || existingCustom.isEmpty()) {
properties.setProperty(customSettingsKey, customSetting);
} else if (!existingCustom.contains("network_compression_method")) {
properties.setProperty(customSettingsKey, existingCustom + "," + customSetting);
}

String customHttpParamsKey = "custom_http_params";
String existingHttpParams = properties.getProperty(customHttpParamsKey);
if (existingHttpParams == null || existingHttpParams.isEmpty()) {
properties.setProperty(customHttpParamsKey, customSetting);
} else if (!existingHttpParams.contains("network_compression_method")) {
properties.setProperty(customHttpParamsKey, existingHttpParams + "," + customSetting);
}
return properties;
}

public DataSource newDataSource(String url, Properties properties) throws SQLException {
properties = addCustomSettings(properties);
if (!properties.containsKey("password")) {
properties.put("password", getPassword());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ public void test01Decompress() throws SQLException {
prop.setProperty("decompress", "true");
prop.setProperty("decompress_algorithm", "lz4");
String url = String.format("jdbc:ch:%s", getEndpointString(true));
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, prop);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, addCustomSettings(prop));
String columnNames = "event_id";
String columnValues = "('event_id String')";
String sql = String.format("INSERT INTO %s (%s) SELECT %s FROM input %s", TABLE_NAME, columnNames, columnNames, columnValues);
Expand Down Expand Up @@ -59,7 +59,7 @@ public void test02Decompress() throws SQLException {
prop.setProperty("decompress", "true");
prop.setProperty("decompress_algorithm", "lz4");
String url = String.format("jdbc:ch:%s", getEndpointString(true));
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, prop);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, addCustomSettings(prop));
String columnNames = "event_id";
String columnValues = "('event_id String')";
String sql = String.format("INSERT INTO %s (%s) SELECT %s FROM input %s", TABLE_NAME, columnNames, columnNames, columnValues);
Expand Down Expand Up @@ -93,7 +93,7 @@ public void test03Decompress() throws SQLException {
prop.setProperty("decompress", "true");
prop.setProperty("decompress_algorithm", "lz4");
String url = String.format("jdbc:ch:%s", getEndpointString(true));
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, prop);
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, addCustomSettings(prop));
String columnNames = "event_id, num01,event_id_01 ";
String columnValues = "('event_id String, num01 Int8, event_id_01 String')";
String sql = String.format("INSERT INTO %s (%s) SELECT %s FROM input %s", TABLE_NAME, columnNames, columnNames, columnValues);
Expand Down Expand Up @@ -126,7 +126,7 @@ public void test03Decompress() throws SQLException {
public void testIssue1373() throws SQLException {
String TABLE_NAME = "issue_1373";
String url = String.format("jdbc:ch:%s", getEndpointString(true));
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, new Properties());
ClickHouseDataSource dataSource = new ClickHouseDataSource(url, addCustomSettings(new Properties()));
String columnNames = "event_id, num01,event_id_01 ";
String columnValues = "('event_id String, num01 Int8, event_id_01 String')";
String sql = String.format("INSERT INTO %s (%s) SELECT %s FROM input %s", TABLE_NAME, columnNames, columnNames, columnValues);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ public Connection getJdbcConnectionV1(Properties properties) throws SQLException
info.putAll(properties);
}

return new ClickHouseConnectionImpl(getJDBCEndpointString(), info);
return new ClickHouseConnectionImpl(getJDBCEndpointString(), addCustomSettings(info));
}

public Connection getJdbcConnectionV2(Properties properties) throws SQLException {
Expand Down
13 changes: 7 additions & 6 deletions client-v2/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,13 @@
<version>${lz4.version}</version>
</dependency>

<dependency>
<groupId>com.github.luben</groupId>
<artifactId>zstd-jni</artifactId>
<version>1.5.7-20</version>
<classifier>cloud</classifier>
</dependency>
Comment thread
chernser marked this conversation as resolved.

<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-compress</artifactId>
Expand Down Expand Up @@ -188,12 +195,6 @@
<version>5.19.0</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>com.github.luben</groupId>
<artifactId>zstd-jni</artifactId>
<version>1.5.7-6</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.bouncycastle</groupId>
<artifactId>bcprov-jdk18on</artifactId>
Expand Down
10 changes: 10 additions & 0 deletions client-v2/src/main/java/com/clickhouse/client/api/Client.java
Original file line number Diff line number Diff line change
Expand Up @@ -1333,6 +1333,16 @@ public Builder queryFormat(String format) {
return this;
}

/**
* Compression method used by client when sending data to server.
* @param method - method to use for compression (ex.: LZ4, ZSTD)
* @return this instance of builder
*/
public Builder compressionMethod(CompressionMethod method) {
this.configuration.put(ClientConfigProperties.COMPRESSION_METHOD.getKey(), method.name());
return this;
}

public Client build() {
// check if endpoint are empty. so can not initiate client
if (this.endpoints.isEmpty()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
import com.clickhouse.client.api.data_formats.ClickHouseFormatReader;
import com.clickhouse.client.api.data_formats.internal.AbstractBinaryFormatReader;
import com.clickhouse.client.api.enums.SSLMode;
import com.clickhouse.client.api.internal.ClickHouseLZ4OutputStream;
import com.clickhouse.client.api.internal.CompressedBlockOutputStream;
import com.clickhouse.data.ClickHouseDataType;
import com.clickhouse.data.ClickHouseFormat;
import org.slf4j.Logger;
Expand Down Expand Up @@ -99,7 +99,9 @@ public enum ClientConfigProperties {

USE_HTTP_COMPRESSION("client.use_http_compression", Boolean.class, "false"),

COMPRESSION_LZ4_UNCOMPRESSED_BUF_SIZE("compression.lz4.uncompressed_buffer_size", Integer.class, String.valueOf(ClickHouseLZ4OutputStream.UNCOMPRESSED_BUFF_SIZE)),
COMPRESSION_LZ4_UNCOMPRESSED_BUF_SIZE("compression.lz4.uncompressed_buffer_size", Integer.class, String.valueOf(CompressedBlockOutputStream.UNCOMPRESSED_BUFF_SIZE)),

COMPRESSION_METHOD("compression.method", CompressionMethod.class, CompressionMethod.ZSTD.name()),

DISABLE_NATIVE_COMPRESSION("disable_native_compression", Boolean.class, "false"),

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package com.clickhouse.client.api;

public enum CompressionMethod {

LZ4,

ZSTD
}
Loading
Loading