From 9dc6a951f27137988f98054d6018394eae0d8ceb Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Mon, 28 Sep 2026 11:16:13 -0700 Subject: [PATCH 1/9] Added ZSTD block compression support --- CHANGELOG.md | 4 + client-v2/pom.xml | 13 +- .../internal/ClickHouseLZ4OutputStream.java | 6 +- ...Entity.java => CompressedBlockEntity.java} | 274 ++++++------- ...m.java => CompressedBlockInputStream.java} | 388 ++++++++++-------- .../api/internal/HttpAPIClientHelper.java | 12 +- ...va => CompressedBlockInputStreamTest.java} | 4 +- 7 files changed, 366 insertions(+), 335 deletions(-) rename client-v2/src/main/java/com/clickhouse/client/api/internal/{LZ4Entity.java => CompressedBlockEntity.java} (88%) rename client-v2/src/main/java/com/clickhouse/client/api/internal/{ClickHouseLZ4InputStream.java => CompressedBlockInputStream.java} (78%) rename client-v2/src/test/java/com/clickhouse/client/api/internal/{ClickHouseLZ4InputStreamTest.java => CompressedBlockInputStreamTest.java} (85%) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6af66dc6f..1ed89aa6c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/client-v2/pom.xml b/client-v2/pom.xml index fb6abe539..ad1f767cf 100644 --- a/client-v2/pom.xml +++ b/client-v2/pom.xml @@ -58,6 +58,13 @@ ${lz4.version} + + com.github.luben + zstd-jni + 1.5.7-20 + cloud + + org.apache.commons commons-compress @@ -188,12 +195,6 @@ 5.19.0 test - - com.github.luben - zstd-jni - 1.5.7-6 - test - org.bouncycastle bcprov-jdk18on diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4OutputStream.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4OutputStream.java index 14304e008..d7ff98671 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4OutputStream.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4OutputStream.java @@ -82,14 +82,14 @@ public void write( byte[] b, int off, int len) throws IOException { public void flush() throws IOException { if (inBuffer.position() > 0) { compressedBuffer.clear(); - compressedBuffer.put(16, ClickHouseLZ4InputStream.MAGIC); + compressedBuffer.put(16, CompressedBlockInputStream.MAGIC); int uncompressedLen = inBuffer.position(); inBuffer.flip(); int compressed = compressor.compress(inBuffer, 0, uncompressedLen, compressedBuffer, 25, compressedBuffer.remaining() - 25); int compressedSizeWithHeader = compressed + 9; - ClickHouseLZ4InputStream.setInt32(compressedBuffer.array(), 17, compressedSizeWithHeader); // compressed size with header - ClickHouseLZ4InputStream.setInt32(compressedBuffer.array(), 21, uncompressedLen); // uncompressed size + CompressedBlockInputStream.setInt32(compressedBuffer.array(), 17, compressedSizeWithHeader); // compressed size with header + CompressedBlockInputStream.setInt32(compressedBuffer.array(), 21, uncompressedLen); // uncompressed size long[] hash = ClickHouseCityHash.cityHash128(compressedBuffer.array(), 16, compressedSizeWithHeader); setInt64(compressedBuffer.array(), 0, hash[0]); setInt64(compressedBuffer.array(), 8, hash[1]); diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/LZ4Entity.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockEntity.java similarity index 88% rename from client-v2/src/main/java/com/clickhouse/client/api/internal/LZ4Entity.java rename to client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockEntity.java index 1cd75c7e3..538d2c67e 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/LZ4Entity.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockEntity.java @@ -1,137 +1,137 @@ -package com.clickhouse.client.api.internal; - -import net.jpountz.lz4.LZ4Factory; -import org.apache.commons.compress.compressors.lz4.FramedLZ4CompressorInputStream; -import org.apache.commons.compress.compressors.lz4.FramedLZ4CompressorOutputStream; -import org.apache.hc.core5.function.Supplier; -import org.apache.hc.core5.http.Header; -import org.apache.hc.core5.http.HttpEntity; - -import java.io.IOException; -import java.io.InputStream; -import java.io.OutputStream; -import java.util.List; -import java.util.Set; - -class LZ4Entity implements HttpEntity { - - private final HttpEntity httpEntity; - - private final boolean useHttpCompression; - - private final int bufferSize; - - private final boolean isResponse; - - private boolean serverCompression; - - private boolean clientCompression; - - private LZ4Factory lz4Factory = null; - - LZ4Entity(HttpEntity httpEntity, boolean useHttpCompression, boolean serverCompression, boolean clientCompression, - int bufferSize, boolean isResponse, LZ4Factory lz4Factory) { - this.httpEntity = httpEntity; - this.useHttpCompression = useHttpCompression; - this.bufferSize = bufferSize; - this.serverCompression = serverCompression; - this.clientCompression = clientCompression; - this.isResponse = isResponse; - this.lz4Factory = lz4Factory; - } - - @Override - public boolean isRepeatable() { - return httpEntity.isRepeatable(); - } - - @Override - public InputStream getContent() throws IOException, UnsupportedOperationException { - if (!isResponse && serverCompression) { - throw new UnsupportedOperationException("Unsupported: getting compressed content of request"); - } else if (serverCompression) { - if (useHttpCompression) { - InputStream content = httpEntity.getContent(); - try { - return new FramedLZ4CompressorInputStream(content); - } catch (IOException e) { - // This is the easiest way to handle empty content because - // - streams at this point wrapped with something else and we can't check content length - // - exception is thrown with no details - // So we just return original content and if there is a real data in it we will get error later - return content; - } - } else { - return new ClickHouseLZ4InputStream(httpEntity.getContent(), lz4Factory.fastDecompressor(), - bufferSize); - } - } else { - return httpEntity.getContent(); - } - } - - @Override - public void writeTo(OutputStream outStream) throws IOException { - if (isResponse && serverCompression) { - // called by us to get compressed response - throw new UnsupportedOperationException("Unsupported: writing compressed response to elsewhere"); - } else if (clientCompression) { - // called by client to send data - OutputStream compressingStream; - if (useHttpCompression) { - compressingStream = new FramedLZ4CompressorOutputStream(outStream); - } else { - compressingStream = new ClickHouseLZ4OutputStream(outStream, lz4Factory.fastCompressor(), bufferSize); - } - - try { - httpEntity.writeTo(compressingStream); - } finally { - compressingStream.close(); - } - } else { - httpEntity.writeTo(outStream); - } - } - - @Override - public boolean isStreaming() { - return httpEntity.isStreaming(); - } - - @Override - public Supplier> getTrailers() { - return httpEntity.getTrailers(); - } - - @Override - public void close() throws IOException { - httpEntity.close(); - } - - @Override - public long getContentLength() { - // compressed request length is unknown event if it is a byte[] - return isResponse ? httpEntity.getContentLength() : -1; - } - - @Override - public String getContentType() { - return httpEntity.getContentType(); - } - - @Override - public String getContentEncoding() { - return httpEntity.getContentEncoding(); - } - - @Override - public boolean isChunked() { - return httpEntity.isChunked(); - } - - @Override - public Set getTrailerNames() { - return httpEntity.getTrailerNames(); - } -} +package com.clickhouse.client.api.internal; + +import net.jpountz.lz4.LZ4Factory; +import org.apache.commons.compress.compressors.lz4.FramedLZ4CompressorInputStream; +import org.apache.commons.compress.compressors.lz4.FramedLZ4CompressorOutputStream; +import org.apache.hc.core5.function.Supplier; +import org.apache.hc.core5.http.Header; +import org.apache.hc.core5.http.HttpEntity; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.util.List; +import java.util.Set; + +class CompressedBlockEntity implements HttpEntity { + + private final HttpEntity httpEntity; + + private final boolean useHttpCompression; + + private final int bufferSize; + + private final boolean isResponse; + + private final boolean serverCompression; + + private final boolean clientCompression; + + private LZ4Factory lz4Factory = null; + + CompressedBlockEntity(HttpEntity httpEntity, boolean useHttpCompression, boolean serverCompression, boolean clientCompression, + int bufferSize, boolean isResponse, LZ4Factory lz4Factory) { + this.httpEntity = httpEntity; + this.useHttpCompression = useHttpCompression; + this.bufferSize = bufferSize; + this.serverCompression = serverCompression; + this.clientCompression = clientCompression; + this.isResponse = isResponse; + this.lz4Factory = lz4Factory; + } + + @Override + public boolean isRepeatable() { + return httpEntity.isRepeatable(); + } + + @Override + public InputStream getContent() throws IOException, UnsupportedOperationException { + if (!isResponse && serverCompression) { + throw new UnsupportedOperationException("Unsupported: getting compressed content of request"); + } else if (serverCompression) { + if (useHttpCompression) { + InputStream content = httpEntity.getContent(); + try { + return new FramedLZ4CompressorInputStream(content); + } catch (IOException e) { + // This is the easiest way to handle empty content because + // - streams at this point wrapped with something else and we can't check content length + // - exception is thrown with no details + // So we just return original content and if there is a real data in it we will get error later + return content; + } + } else { + return new CompressedBlockInputStream(httpEntity.getContent(), lz4Factory.fastDecompressor(), + bufferSize); + } + } else { + return httpEntity.getContent(); + } + } + + @Override + public void writeTo(OutputStream outStream) throws IOException { + if (isResponse && serverCompression) { + // called by us to get compressed response + throw new UnsupportedOperationException("Unsupported: writing compressed response to elsewhere"); + } else if (clientCompression) { + // called by client to send data + OutputStream compressingStream; + if (useHttpCompression) { + compressingStream = new FramedLZ4CompressorOutputStream(outStream); + } else { + compressingStream = new ClickHouseLZ4OutputStream(outStream, lz4Factory.fastCompressor(), bufferSize); + } + + try { + httpEntity.writeTo(compressingStream); + } finally { + compressingStream.close(); + } + } else { + httpEntity.writeTo(outStream); + } + } + + @Override + public boolean isStreaming() { + return httpEntity.isStreaming(); + } + + @Override + public Supplier> getTrailers() { + return httpEntity.getTrailers(); + } + + @Override + public void close() throws IOException { + httpEntity.close(); + } + + @Override + public long getContentLength() { + // compressed request length is unknown event if it is a byte[] + return isResponse ? httpEntity.getContentLength() : -1; + } + + @Override + public String getContentType() { + return httpEntity.getContentType(); + } + + @Override + public String getContentEncoding() { + return httpEntity.getContentEncoding(); + } + + @Override + public boolean isChunked() { + return httpEntity.isChunked(); + } + + @Override + public Set getTrailerNames() { + return httpEntity.getTrailerNames(); + } +} diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4InputStream.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java similarity index 78% rename from client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4InputStream.java rename to client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java index 980b06cf9..4890ac73a 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4InputStream.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java @@ -1,182 +1,208 @@ -package com.clickhouse.client.api.internal; - -import com.clickhouse.client.api.ClientException; -import com.clickhouse.data.ClickHouseByteUtils; -import com.clickhouse.data.ClickHouseCityHash; -import com.clickhouse.data.ClickHouseUtils; -import net.jpountz.lz4.LZ4FastDecompressor; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.io.EOFException; -import java.io.IOException; -import java.io.InputStream; -import java.nio.ByteBuffer; - -public class ClickHouseLZ4InputStream extends InputStream { - - private static Logger LOG = LoggerFactory.getLogger(ClickHouseLZ4InputStream.class); - private final LZ4FastDecompressor decompressor; - - private final InputStream in; - - private ByteBuffer buffer; - - private byte[] tmpBuffer = new byte[1]; - - - public ClickHouseLZ4InputStream(InputStream in, LZ4FastDecompressor decompressor, int bufferSize) { - super(); - LOG.debug("Using LZ4 decompressor with buffer size {}", bufferSize); - this.decompressor = decompressor; - this.in = in; - this.buffer = ByteBuffer.allocate(bufferSize); - this.buffer.limit(0); - } - - @Override - public int read() throws IOException { - int n = read(tmpBuffer, 0, 1); - return n == -1 ? -1 : tmpBuffer[0] & 0xFF; - } - - @Override - public int read(byte[] b, int off, int len) throws IOException { - if (b == null) { - throw new NullPointerException("b is null"); - } else if (off < 0) { - throw new IndexOutOfBoundsException("off is negative"); - } else if (len < 0) { - throw new IndexOutOfBoundsException("len is negative"); - } else if (off + len > b.length) { - throw new IndexOutOfBoundsException("off + len is greater than b.length"); - } else if (len == 0) { - return 0; - } - - int readBytes = 0; - do { - int remaining = Math.min(len - readBytes, buffer.remaining()); - buffer.get(b, off + readBytes, remaining); - readBytes += remaining; - } while (readBytes < len && refill() != -1); - - return readBytes == 0 ? -1 : readBytes; - } - - - static final byte MAGIC = (byte) 0x82; - static final int HEADER_LENGTH = 25; - - final byte[] headerBuff = new byte[HEADER_LENGTH]; - - /** - * Method ensures to read all bytes from the input stream. - * In case of network connection it may be a case when not all bytes are read at once. - * @throws IOException - */ - private boolean readFully(byte[] b, int off, int len) throws IOException { - int n = 0; - while (n < len) { - int count = in.read(b, off + n, len - n); - if (count < 0) { - if (n == 0) { - return false; - } +package com.clickhouse.client.api.internal; + +import com.clickhouse.client.api.ClientException; +import com.clickhouse.data.ClickHouseByteUtils; +import com.clickhouse.data.ClickHouseCityHash; +import com.clickhouse.data.ClickHouseUtils; +import com.github.luben.zstd.Zstd; +import net.jpountz.lz4.LZ4FastDecompressor; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.EOFException; +import java.io.IOException; +import java.io.InputStream; +import java.nio.ByteBuffer; + +public class CompressedBlockInputStream extends InputStream { + + private static Logger LOG = LoggerFactory.getLogger(CompressedBlockInputStream.class); + private final LZ4FastDecompressor decompressor; + + private final InputStream in; + + private ByteBuffer buffer; + + private byte[] tmpBuffer = new byte[1]; + + + public CompressedBlockInputStream(InputStream in, LZ4FastDecompressor decompressor, int bufferSize) { + super(); + LOG.debug("Using LZ4 decompressor with buffer size {}", bufferSize); + this.decompressor = decompressor; + this.in = in; + this.buffer = ByteBuffer.allocate(bufferSize); + this.buffer.limit(0); + } + + @Override + public int read() throws IOException { + int n = read(tmpBuffer, 0, 1); + return n == -1 ? -1 : tmpBuffer[0] & 0xFF; + } + + @Override + public int read(byte[] b, int off, int len) throws IOException { + if (b == null) { + throw new NullPointerException("b is null"); + } else if (off < 0) { + throw new IndexOutOfBoundsException("off is negative"); + } else if (len < 0) { + throw new IndexOutOfBoundsException("len is negative"); + } else if (off + len > b.length) { + throw new IndexOutOfBoundsException("off + len is greater than b.length"); + } else if (len == 0) { + return 0; + } + + int readBytes = 0; + do { + int remaining = Math.min(len - readBytes, buffer.remaining()); + buffer.get(b, off + readBytes, remaining); + readBytes += remaining; + } while (readBytes < len && refill() != -1); + + return readBytes == 0 ? -1 : readBytes; + } + + + static final byte MAGIC = (byte) 0x82; + static final byte MAGIC_LZ4 = (byte) 0x82; + static final byte MAGIC_ZSTD_3 = (byte) 0x90; + static final byte MAGIC_NONE = (byte) 0x02; + static final int HEADER_LENGTH = 25; + + final byte[] headerBuff = new byte[HEADER_LENGTH]; + + /** + * Method ensures to read all bytes from the input stream. + * In case of network connection it may be a case when not all bytes are read at once. + * @throws IOException + */ + private boolean readFully(byte[] b, int off, int len) throws IOException { + int n = 0; + while (n < len) { + int count = in.read(b, off + n, len - n); + if (count < 0) { + if (n == 0) { + return false; + } throw new IOException(ClickHouseUtils.format("Incomplete read: %s of %s", n, len)); - } - n += count; - } - - return true; - } - - public byte[] getHeaderBuffer() { - return headerBuff; - } - - public InputStream getInputStream() { - return in; - } - - private int refill() throws IOException { - - // read header - boolean readFully = readFully(headerBuff, 0, HEADER_LENGTH); - if (!readFully) { - return -1; - } - - if (headerBuff[16] != MAGIC) { - // 1 byte - 0x82 (shows this is LZ4) - throw new ClientException("Invalid LZ4 magic byte: '" + headerBuff[16] + "'"); - } - - // 4 bytes - size of the compressed data including 9 bytes of the header - int compressedSizeWithHeader = getInt32(headerBuff, 17); - // 4 bytes - size of uncompressed data - int uncompressedSize = getInt32(headerBuff, 21); - - int offset = 9; - final byte[] block = new byte[compressedSizeWithHeader]; - block[0] = MAGIC; - setInt32(block, 1, compressedSizeWithHeader); - setInt32(block, 5, uncompressedSize); - // compressed data: compressed_size - 9 bytes - int remaining = compressedSizeWithHeader - offset; - - readFully = readFully(block, offset, remaining); - if (!readFully) { - throw new EOFException("Unexpected end of stream"); - } - - long[] real = ClickHouseCityHash.cityHash128(block, 0, compressedSizeWithHeader); - if (real[0] != getInt64(headerBuff, 0) || real[1] != ClickHouseByteUtils.getInt64(headerBuff, 8)) { - throw new ClientException("Corrupted stream: checksum mismatch"); - } - - if (buffer.capacity() < uncompressedSize) { - buffer = ByteBuffer.allocate(uncompressedSize); - } - decompressor.decompress(ByteBuffer.wrap(block), offset, buffer, 0, uncompressedSize); - buffer.position(0); - buffer.limit(uncompressedSize); - return uncompressedSize; - } - - /** - * Read int32 Little Endian - * @param bytes - * @param offset - * @return - */ - static int getInt32(byte[] bytes, int offset) { - return (0xFF & bytes[offset]) | ((0xFF & bytes[offset + 1]) << 8) | ((0xFF & bytes[offset + 2]) << 16) - | ((0xFF & bytes[offset + 3]) << 24); - } - - /** - * Read int64 Little Endian - * @param bytes - * @param offset - * @return - */ - static long getInt64(byte[] bytes, int offset) { - return (0xFFL & bytes[offset]) | ((0xFFL & bytes[offset + 1]) << 8) | ((0xFFL & bytes[offset + 2]) << 16) - | ((0xFFL & bytes[offset + 3]) << 24) | ((0xFFL & bytes[offset + 4]) << 32) - | ((0xFFL & bytes[offset + 5]) << 40) | ((0xFFL & bytes[offset + 6]) << 48) - | ((0xFFL & bytes[offset + 7]) << 56); - } - - static void setInt32(byte[] bytes, int offset, int value) { - bytes[offset] = (byte) (0xFF & value); - bytes[offset + 1] = (byte) (0xFF & (value >> 8)); - bytes[offset + 2] = (byte) (0xFF & (value >> 16)); - bytes[offset + 3] = (byte) (0xFF & (value >> 24)); - } - - @Override - public void close() throws IOException { - in.close(); - } -} + } + n += count; + } + + return true; + } + + public byte[] getHeaderBuffer() { + return headerBuff; + } + + public InputStream getInputStream() { + return in; + } + + private int refill() throws IOException { + + // read header + boolean readFully = readFully(headerBuff, 0, HEADER_LENGTH); + if (!readFully) { + return -1; + } + + byte magicNumber = headerBuff[16]; + switch (magicNumber) { + case MAGIC_LZ4: + case MAGIC_ZSTD_3: + case MAGIC_NONE: + break; + default: + throw new ClientException("Invalid LZ4 magic byte: '" + magicNumber + "'"); + } + + // 4 bytes - size of the compressed data including 9 bytes of the header + int compressedSizeWithHeader = getInt32(headerBuff, 17); + // 4 bytes - size of uncompressed data + int uncompressedSize = getInt32(headerBuff, 21); + + int offset = 9; + final byte[] block = new byte[compressedSizeWithHeader]; + block[0] = magicNumber; + setInt32(block, 1, compressedSizeWithHeader); + setInt32(block, 5, uncompressedSize); + // compressed data: compressed_size - 9 bytes + int remaining = compressedSizeWithHeader - offset; + + readFully = readFully(block, offset, remaining); + if (!readFully) { + throw new EOFException("Unexpected end of stream"); + } + + long[] real = ClickHouseCityHash.cityHash128(block, 0, compressedSizeWithHeader); + if (real[0] != getInt64(headerBuff, 0) || real[1] != ClickHouseByteUtils.getInt64(headerBuff, 8)) { + throw new ClientException("Corrupted stream: checksum mismatch"); + } + + if (buffer.capacity() < uncompressedSize) { + buffer = ByteBuffer.allocate(uncompressedSize); + } + + switch (magicNumber) { + case MAGIC_LZ4: { + ByteBuffer blockBuff = ByteBuffer.wrap(block, offset, remaining); + decompressor.decompress(blockBuff, offset, buffer, 0, uncompressedSize); + break; + } + case MAGIC_ZSTD_3: { + Zstd.decompressByteArray(buffer.array(), 0, uncompressedSize, block, offset, remaining); + break; + } + case MAGIC_NONE: + // block is not compressed - just put it. + buffer.put(block, 0, uncompressedSize); + break; + default: + throw new ClientException("bug: Should not reach here"); + } + buffer.position(0); + buffer.limit(uncompressedSize); + return uncompressedSize; + } + + /** + * Read int32 Little Endian + * @param bytes + * @param offset + * @return + */ + static int getInt32(byte[] bytes, int offset) { + return (0xFF & bytes[offset]) | ((0xFF & bytes[offset + 1]) << 8) | ((0xFF & bytes[offset + 2]) << 16) + | ((0xFF & bytes[offset + 3]) << 24); + } + + /** + * Read int64 Little Endian + * @param bytes + * @param offset + * @return + */ + static long getInt64(byte[] bytes, int offset) { + return (0xFFL & bytes[offset]) | ((0xFFL & bytes[offset + 1]) << 8) | ((0xFFL & bytes[offset + 2]) << 16) + | ((0xFFL & bytes[offset + 3]) << 24) | ((0xFFL & bytes[offset + 4]) << 32) + | ((0xFFL & bytes[offset + 5]) << 40) | ((0xFFL & bytes[offset + 6]) << 48) + | ((0xFFL & bytes[offset + 7]) << 56); + } + + static void setInt32(byte[] bytes, int offset, int value) { + bytes[offset] = (byte) (0xFF & value); + bytes[offset + 1] = (byte) (0xFF & (value >> 8)); + bytes[offset + 2] = (byte) (0xFF & (value >> 16)); + bytes[offset + 3] = (byte) (0xFF & (value >> 24)); + } + + @Override + public void close() throws IOException { + in.close(); + } +} diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/HttpAPIClientHelper.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/HttpAPIClientHelper.java index 9dab47a39..b3295f83c 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/HttpAPIClientHelper.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/HttpAPIClientHelper.java @@ -448,8 +448,8 @@ private ServerException readNotClickHouseError(HttpEntity httpEntity, String que break; } catch (ClientException e) { // Invalid LZ4 Magic - if (body instanceof ClickHouseLZ4InputStream) { - ClickHouseLZ4InputStream stream = (ClickHouseLZ4InputStream) body; + if (body instanceof CompressedBlockInputStream) { + CompressedBlockInputStream stream = (CompressedBlockInputStream) body; body = stream.getInputStream(); byte[] lzHeader = stream.getHeaderBuffer(); // Here is read part of original body offset = Math.min(lzHeader.length, buffer.length); @@ -479,8 +479,8 @@ private static ServerException readClickHouseError(HttpEntity httpEntity, int se rBytes = body.read(buffer); } catch (ClientException e) { // Invalid LZ4 Magic - if (body instanceof ClickHouseLZ4InputStream) { - ClickHouseLZ4InputStream stream = (ClickHouseLZ4InputStream) body; + if (body instanceof CompressedBlockInputStream) { + CompressedBlockInputStream stream = (CompressedBlockInputStream) body; body = stream.getInputStream(); byte[] headerBuffer = stream.getHeaderBuffer(); System.arraycopy(headerBuffer, 0, buffer, 0, headerBuffer.length); @@ -1029,7 +1029,7 @@ private HttpEntity wrapRequestEntity(HttpEntity httpEntity, Map return new CompressedEntity(httpEntity, false, CompressorStreamFactory.getSingleton()); } else if (clientCompression && !appCompressedData) { int buffSize = ClientConfigProperties.COMPRESSION_LZ4_UNCOMPRESSED_BUF_SIZE.getOrDefault(requestConfig); - return new LZ4Entity(httpEntity, useHttpCompression, false, true, + return new CompressedBlockEntity(httpEntity, useHttpCompression, false, true, buffSize, false, lz4Factory); } else { return httpEntity; @@ -1048,7 +1048,7 @@ private HttpEntity wrapResponseEntity(HttpEntity httpEntity, int httpStatus, Map // data compression if (serverCompression && !(httpStatus == HttpStatus.SC_FORBIDDEN || httpStatus == HttpStatus.SC_UNAUTHORIZED)) { int buffSize = ClientConfigProperties.COMPRESSION_LZ4_UNCOMPRESSED_BUF_SIZE.getOrDefault(requestConfig); - return new LZ4Entity(httpEntity, useHttpCompression, true, false, buffSize, true, lz4Factory); + return new CompressedBlockEntity(httpEntity, useHttpCompression, true, false, buffSize, true, lz4Factory); } return httpEntity; diff --git a/client-v2/src/test/java/com/clickhouse/client/api/internal/ClickHouseLZ4InputStreamTest.java b/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java similarity index 85% rename from client-v2/src/test/java/com/clickhouse/client/api/internal/ClickHouseLZ4InputStreamTest.java rename to client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java index e6ab846e5..42890eb39 100644 --- a/client-v2/src/test/java/com/clickhouse/client/api/internal/ClickHouseLZ4InputStreamTest.java +++ b/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java @@ -7,12 +7,12 @@ import org.testng.Assert; import org.testng.annotations.Test; -public class ClickHouseLZ4InputStreamTest { +public class CompressedBlockInputStreamTest { @Test public void reportsActualByteCountsForTruncatedHeader() { byte[] truncatedHeader = new byte[10]; - ClickHouseLZ4InputStream input = new ClickHouseLZ4InputStream( + CompressedBlockInputStream input = new CompressedBlockInputStream( new ByteArrayInputStream(truncatedHeader), LZ4Factory.fastestJavaInstance().fastDecompressor(), 8192); From 0342bd98b0800ffbff93981748d79096a970114e Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Mon, 28 Sep 2026 12:02:00 -0700 Subject: [PATCH 2/9] Added option to set compression method for outgoing requests --- .../com/clickhouse/client/api/Client.java | 10 + .../client/api/ClientConfigProperties.java | 6 +- .../client/api/CompressionMethod.java | 8 + .../api/internal/CompressedBlockEntity.java | 20 +- .../internal/CompressedBlockInputStream.java | 6 +- ....java => CompressedBlockOutputStream.java} | 273 ++++++++++-------- .../api/internal/HttpAPIClientHelper.java | 7 +- .../com/clickhouse/client/ClientTests.java | 53 +--- 8 files changed, 208 insertions(+), 175 deletions(-) create mode 100644 client-v2/src/main/java/com/clickhouse/client/api/CompressionMethod.java rename client-v2/src/main/java/com/clickhouse/client/api/internal/{ClickHouseLZ4OutputStream.java => CompressedBlockOutputStream.java} (61%) diff --git a/client-v2/src/main/java/com/clickhouse/client/api/Client.java b/client-v2/src/main/java/com/clickhouse/client/api/Client.java index d0eebf9d0..da6ab010b 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/Client.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/Client.java @@ -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()) { diff --git a/client-v2/src/main/java/com/clickhouse/client/api/ClientConfigProperties.java b/client-v2/src/main/java/com/clickhouse/client/api/ClientConfigProperties.java index 07c11e9bf..63173c5a2 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/ClientConfigProperties.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/ClientConfigProperties.java @@ -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; @@ -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"), diff --git a/client-v2/src/main/java/com/clickhouse/client/api/CompressionMethod.java b/client-v2/src/main/java/com/clickhouse/client/api/CompressionMethod.java new file mode 100644 index 000000000..5ffd91e70 --- /dev/null +++ b/client-v2/src/main/java/com/clickhouse/client/api/CompressionMethod.java @@ -0,0 +1,8 @@ +package com.clickhouse.client.api; + +public enum CompressionMethod { + + LZ4, + + ZSTD +} diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockEntity.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockEntity.java index 538d2c67e..10990dbb8 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockEntity.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockEntity.java @@ -1,8 +1,10 @@ package com.clickhouse.client.api.internal; +import com.clickhouse.client.api.CompressionMethod; import net.jpountz.lz4.LZ4Factory; import org.apache.commons.compress.compressors.lz4.FramedLZ4CompressorInputStream; import org.apache.commons.compress.compressors.lz4.FramedLZ4CompressorOutputStream; +import org.apache.hc.client5.http.ClientProtocolException; import org.apache.hc.core5.function.Supplier; import org.apache.hc.core5.http.Header; import org.apache.hc.core5.http.HttpEntity; @@ -27,10 +29,12 @@ class CompressedBlockEntity implements HttpEntity { private final boolean clientCompression; - private LZ4Factory lz4Factory = null; + private final LZ4Factory lz4Factory; + + private final CompressionMethod compressionMethod; CompressedBlockEntity(HttpEntity httpEntity, boolean useHttpCompression, boolean serverCompression, boolean clientCompression, - int bufferSize, boolean isResponse, LZ4Factory lz4Factory) { + int bufferSize, boolean isResponse, LZ4Factory lz4Factory, CompressionMethod compressMethod) { this.httpEntity = httpEntity; this.useHttpCompression = useHttpCompression; this.bufferSize = bufferSize; @@ -38,6 +42,7 @@ class CompressedBlockEntity implements HttpEntity { this.clientCompression = clientCompression; this.isResponse = isResponse; this.lz4Factory = lz4Factory; + this.compressionMethod = compressMethod; } @Override @@ -81,7 +86,16 @@ public void writeTo(OutputStream outStream) throws IOException { if (useHttpCompression) { compressingStream = new FramedLZ4CompressorOutputStream(outStream); } else { - compressingStream = new ClickHouseLZ4OutputStream(outStream, lz4Factory.fastCompressor(), bufferSize); + switch (compressionMethod) { + case LZ4: + compressingStream = new CompressedBlockOutputStream.LZ4OutputStream(outStream, lz4Factory.fastCompressor(), bufferSize); + break; + case ZSTD: + compressingStream = new CompressedBlockOutputStream.ZSTDOutputStream(outStream, bufferSize); + break; + default: + throw new ClientProtocolException("bug: unsupported compression method " + compressionMethod); + } } try { diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java index 4890ac73a..bfdb095f3 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java @@ -16,19 +16,19 @@ public class CompressedBlockInputStream extends InputStream { - private static Logger LOG = LoggerFactory.getLogger(CompressedBlockInputStream.class); + private static final Logger LOG = LoggerFactory.getLogger(CompressedBlockInputStream.class); private final LZ4FastDecompressor decompressor; private final InputStream in; private ByteBuffer buffer; - private byte[] tmpBuffer = new byte[1]; + private final byte[] tmpBuffer = new byte[1]; public CompressedBlockInputStream(InputStream in, LZ4FastDecompressor decompressor, int bufferSize) { super(); - LOG.debug("Using LZ4 decompressor with buffer size {}", bufferSize); + LOG.debug("Using decompressor with buffer size {}", bufferSize); this.decompressor = decompressor; this.in = in; this.buffer = ByteBuffer.allocate(bufferSize); diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4OutputStream.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java similarity index 61% rename from client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4OutputStream.java rename to client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java index d7ff98671..157cf6fb2 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/ClickHouseLZ4OutputStream.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java @@ -1,118 +1,155 @@ -package com.clickhouse.client.api.internal; - -import com.clickhouse.data.ClickHouseCityHash; -import net.jpountz.lz4.LZ4Compressor; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.io.IOException; -import java.io.OutputStream; -import java.nio.ByteBuffer; - -public class ClickHouseLZ4OutputStream extends OutputStream { - - private static Logger LOG = LoggerFactory.getLogger(ClickHouseLZ4OutputStream.class); - public static final int UNCOMPRESSED_BUFF_SIZE = 64 * 1024; // 64K is most optimal for LZ4 compression - - private final ByteBuffer inBuffer; - - private final OutputStream out; - - private final LZ4Compressor compressor; - - private byte tmpBuffer[] = new byte[1]; - - private final ByteBuffer compressedBuffer; - - private static int HEADER_LEN = 15; // 9 bytes for header, 6 bytes for checksum - - - public ClickHouseLZ4OutputStream(OutputStream out, LZ4Compressor compressor, int bufferSize) { - super(); - LOG.debug("Using LZ4 compressor with buffer size {}", bufferSize); - this.inBuffer = ByteBuffer.allocate(bufferSize); - this.out = out; - this.compressor = compressor; - this.compressedBuffer = ByteBuffer.allocate(compressor.maxCompressedLength(inBuffer.capacity()) + HEADER_LEN); - } - - @Override - public void write(int b) throws IOException { - if (inBuffer.remaining() == 0) { - flush(); - } - inBuffer.put((byte) b); - } - - @Override - public void write(byte[] b) throws IOException { - if (b.length == 1) { - write(b[0]); - } else { - write(b, 0, b.length); - } - } - - @Override - public void write( byte[] b, int off, int len) throws IOException { - if (b == null) { - throw new NullPointerException("b is null"); - } else if (off < 0) { - throw new IndexOutOfBoundsException("off is negative"); - } else if (len < 0) { - throw new IndexOutOfBoundsException("len is negative"); - } else if (off + len > b.length) { - throw new IndexOutOfBoundsException("off + len is greater than b.length"); - } else if (len == 0) { - return; - } - - int writtenBytes = 0; - do { - if (inBuffer.remaining() == 0) { - flush(); // flush will make inBuffer clear - } - int remaining = Math.min(len - writtenBytes, inBuffer.remaining()); - inBuffer.put(b, off + writtenBytes, remaining); - writtenBytes += remaining; - } while (writtenBytes < len); - } - - @Override - public void flush() throws IOException { - if (inBuffer.position() > 0) { - compressedBuffer.clear(); - compressedBuffer.put(16, CompressedBlockInputStream.MAGIC); - int uncompressedLen = inBuffer.position(); - inBuffer.flip(); - int compressed = compressor.compress(inBuffer, 0, uncompressedLen, compressedBuffer, 25, - compressedBuffer.remaining() - 25); - int compressedSizeWithHeader = compressed + 9; - CompressedBlockInputStream.setInt32(compressedBuffer.array(), 17, compressedSizeWithHeader); // compressed size with header - CompressedBlockInputStream.setInt32(compressedBuffer.array(), 21, uncompressedLen); // uncompressed size - long[] hash = ClickHouseCityHash.cityHash128(compressedBuffer.array(), 16, compressedSizeWithHeader); - setInt64(compressedBuffer.array(), 0, hash[0]); - setInt64(compressedBuffer.array(), 8, hash[1]); - compressedBuffer.flip(); - out.write(compressedBuffer.array(), 0, compressed + 25); - inBuffer.clear(); - } - } - - - static void setInt64(byte[] bytes, int offset, long value) { - bytes[offset] = (byte) (0xFF & value); - bytes[offset + 1] = (byte) (0xFF & (value >> 8)); - bytes[offset + 2] = (byte) (0xFF & (value >> 16)); - bytes[offset + 3] = (byte) (0xFF & (value >> 24)); - bytes[offset + 4] = (byte) (0xFF & (value >> 32)); - bytes[offset + 5] = (byte) (0xFF & (value >> 40)); - bytes[offset + 6] = (byte) (0xFF & (value >> 48)); - bytes[offset + 7] = (byte) (0xFF & (value >> 56)); - } - @Override - public void close() throws IOException { - flush(); - out.close(); - } -} +package com.clickhouse.client.api.internal; + +import com.clickhouse.data.ClickHouseCityHash; +import com.github.luben.zstd.Zstd; +import net.jpountz.lz4.LZ4Compressor; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.io.OutputStream; +import java.nio.ByteBuffer; + +public abstract class CompressedBlockOutputStream extends OutputStream { + + private static final Logger LOG = LoggerFactory.getLogger(CompressedBlockOutputStream.class); + public static final int UNCOMPRESSED_BUFF_SIZE = 64 * 1024; // 64K is most optimal for LZ4 compression + + private final ByteBuffer inBuffer; + + private final OutputStream out; + + private final ByteBuffer compressedBuffer; + + private static final int HEADER_LEN = 15; // 9 bytes for header, 6 bytes for checksum + + + public CompressedBlockOutputStream(OutputStream out, int bufferSize) { + super(); + LOG.debug("Using compressor with buffer size {}", bufferSize); + this.inBuffer = ByteBuffer.allocate(bufferSize); + this.out = out; + this.compressedBuffer = ByteBuffer.allocate(bufferSize + HEADER_LEN); + } + + @Override + public void write(int b) throws IOException { + if (inBuffer.remaining() == 0) { + flush(); + } + inBuffer.put((byte) b); + } + + @Override + public void write(byte[] b) throws IOException { + if (b.length == 1) { + write(b[0]); + } else { + write(b, 0, b.length); + } + } + + @Override + public void write( byte[] b, int off, int len) throws IOException { + if (b == null) { + throw new NullPointerException("b is null"); + } else if (off < 0) { + throw new IndexOutOfBoundsException("off is negative"); + } else if (len < 0) { + throw new IndexOutOfBoundsException("len is negative"); + } else if (off + len > b.length) { + throw new IndexOutOfBoundsException("off + len is greater than b.length"); + } else if (len == 0) { + return; + } + + int writtenBytes = 0; + do { + if (inBuffer.remaining() == 0) { + flush(); // flush will make inBuffer clear + } + int remaining = Math.min(len - writtenBytes, inBuffer.remaining()); + inBuffer.put(b, off + writtenBytes, remaining); + writtenBytes += remaining; + } while (writtenBytes < len); + } + + @Override + public void flush() throws IOException { + if (inBuffer.position() > 0) { + compressedBuffer.clear(); + compressedBuffer.put(16, getMagicNumber()); + int uncompressedLen = inBuffer.position(); + inBuffer.flip(); + int compressed = compressData(inBuffer, uncompressedLen, compressedBuffer, compressedBuffer); + int compressedSizeWithHeader = compressed + 9; + CompressedBlockInputStream.setInt32(compressedBuffer.array(), 17, compressedSizeWithHeader); // compressed size with header + CompressedBlockInputStream.setInt32(compressedBuffer.array(), 21, uncompressedLen); // uncompressed size + long[] hash = ClickHouseCityHash.cityHash128(compressedBuffer.array(), 16, compressedSizeWithHeader); + setInt64(compressedBuffer.array(), 0, hash[0]); + setInt64(compressedBuffer.array(), 8, hash[1]); + compressedBuffer.flip(); + out.write(compressedBuffer.array(), 0, compressed + 25); + inBuffer.clear(); + } + } + + protected abstract byte getMagicNumber(); + + protected abstract int compressData(ByteBuffer inBuffer, int uncompressedLen, ByteBuffer compressedBuffer, ByteBuffer buffer); + + static void setInt64(byte[] bytes, int offset, long value) { + bytes[offset] = (byte) (0xFF & value); + bytes[offset + 1] = (byte) (0xFF & (value >> 8)); + bytes[offset + 2] = (byte) (0xFF & (value >> 16)); + bytes[offset + 3] = (byte) (0xFF & (value >> 24)); + bytes[offset + 4] = (byte) (0xFF & (value >> 32)); + bytes[offset + 5] = (byte) (0xFF & (value >> 40)); + bytes[offset + 6] = (byte) (0xFF & (value >> 48)); + bytes[offset + 7] = (byte) (0xFF & (value >> 56)); + } + @Override + public void close() throws IOException { + flush(); + out.close(); + } + + public static class LZ4OutputStream extends CompressedBlockOutputStream { + + private final LZ4Compressor compressor; + + public LZ4OutputStream(OutputStream out, LZ4Compressor compressor, int bufferSize) { + super(out, bufferSize); + this.compressor = compressor; + } + + @Override + protected int compressData(ByteBuffer inBuffer, int uncompressedLen, ByteBuffer compressedBuffer, ByteBuffer buffer) { + return compressor.compress(inBuffer, 0, uncompressedLen, compressedBuffer, 25, + compressedBuffer.remaining() - 25); + } + + @Override + protected byte getMagicNumber() { + return CompressedBlockInputStream.MAGIC_LZ4; + } + } + + public static class ZSTDOutputStream extends CompressedBlockOutputStream { + + public ZSTDOutputStream(OutputStream out, int bufferSize) { + super(out, bufferSize); + } + + @Override + protected int compressData(ByteBuffer inBuffer, int uncompressedLen, ByteBuffer compressedBuffer, ByteBuffer buffer) { + return (int) Zstd.compressByteArray(compressedBuffer.array(), 25, compressedBuffer.remaining() - 25, + inBuffer.array(), 0, uncompressedLen, 3); + } + + @Override + protected byte getMagicNumber() { + return CompressedBlockInputStream.MAGIC_ZSTD_3; + } + } +} diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/HttpAPIClientHelper.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/HttpAPIClientHelper.java index b3295f83c..ff8c22d53 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/HttpAPIClientHelper.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/HttpAPIClientHelper.java @@ -6,6 +6,7 @@ import com.clickhouse.client.api.ClientException; import com.clickhouse.client.api.ClientFaultCause; import com.clickhouse.client.api.ClientMisconfigurationException; +import com.clickhouse.client.api.CompressionMethod; import com.clickhouse.client.api.ConnectionInitiationException; import com.clickhouse.client.api.ConnectionReuseStrategy; import com.clickhouse.client.api.DataTransferException; @@ -1029,8 +1030,9 @@ private HttpEntity wrapRequestEntity(HttpEntity httpEntity, Map return new CompressedEntity(httpEntity, false, CompressorStreamFactory.getSingleton()); } else if (clientCompression && !appCompressedData) { int buffSize = ClientConfigProperties.COMPRESSION_LZ4_UNCOMPRESSED_BUF_SIZE.getOrDefault(requestConfig); + CompressionMethod compressMethod = ClientConfigProperties.COMPRESSION_METHOD.getOrDefault(requestConfig); return new CompressedBlockEntity(httpEntity, useHttpCompression, false, true, - buffSize, false, lz4Factory); + buffSize, false, lz4Factory, compressMethod); } else { return httpEntity; } @@ -1048,7 +1050,8 @@ private HttpEntity wrapResponseEntity(HttpEntity httpEntity, int httpStatus, Map // data compression if (serverCompression && !(httpStatus == HttpStatus.SC_FORBIDDEN || httpStatus == HttpStatus.SC_UNAUTHORIZED)) { int buffSize = ClientConfigProperties.COMPRESSION_LZ4_UNCOMPRESSED_BUF_SIZE.getOrDefault(requestConfig); - return new CompressedBlockEntity(httpEntity, useHttpCompression, true, false, buffSize, true, lz4Factory); + return new CompressedBlockEntity(httpEntity, useHttpCompression, true, false, buffSize, true, lz4Factory, + ClientConfigProperties.COMPRESSION_METHOD.getDefObjVal()); } return httpEntity; diff --git a/client-v2/src/test/java/com/clickhouse/client/ClientTests.java b/client-v2/src/test/java/com/clickhouse/client/ClientTests.java index 50a80f589..6c3a73ba1 100644 --- a/client-v2/src/test/java/com/clickhouse/client/ClientTests.java +++ b/client-v2/src/test/java/com/clickhouse/client/ClientTests.java @@ -5,13 +5,14 @@ import com.clickhouse.client.api.ClientException; import com.clickhouse.client.api.ClientFaultCause; import com.clickhouse.client.api.ClientMisconfigurationException; +import com.clickhouse.client.api.CompressionMethod; import com.clickhouse.client.api.ConnectionInitiationException; import com.clickhouse.client.api.ConnectionReuseStrategy; import com.clickhouse.client.api.ServerException; import com.clickhouse.client.api.command.CommandSettings; import com.clickhouse.client.api.enums.Protocol; import com.clickhouse.client.api.insert.InsertSettings; -import com.clickhouse.client.api.internal.ClickHouseLZ4OutputStream; +import com.clickhouse.client.api.internal.CompressedBlockOutputStream; import com.clickhouse.client.api.internal.CredentialsManager; import com.clickhouse.client.api.internal.ServerSettings; import com.clickhouse.client.api.internal.ValidationUtils; @@ -334,7 +335,7 @@ public void testDefaultSettings() { Assert.assertEquals(config.get(p.getKey()), p.getDefaultValue(), "Default value doesn't match"); } } - Assert.assertEquals(config.size(), 36); // to check everything is set. Increment when new added. + Assert.assertEquals(config.size(), 37); // to check everything is set. Increment when new added. } try (Client client = new Client.Builder() @@ -369,7 +370,7 @@ public void testDefaultSettings() { .queryFormat(ClickHouseFormat.CSV.name()) .build()) { Map config = client.getConfiguration(); - Assert.assertEquals(config.size(), 39); // to check everything is set. Increment when new added. + Assert.assertEquals(config.size(), 40); // to check everything is set. Increment when new added. Assert.assertEquals(config.get(ClientConfigProperties.DATABASE.getKey()), "mydb"); Assert.assertEquals(config.get(ClientConfigProperties.MAX_EXECUTION_TIME.getKey()), "10"); Assert.assertEquals(config.get(ClientConfigProperties.COMPRESSION_LZ4_UNCOMPRESSED_BUF_SIZE.getKey()), "300000"); @@ -396,6 +397,8 @@ public void testDefaultSettings() { Assert.assertEquals(config.get(ClientConfigProperties.SSL_MODE.getKey()), "STRICT"); Assert.assertEquals(config.get(ClientConfigProperties.BINARY_STRING_SUPPORT.getKey()), "true"); Assert.assertEquals(config.get(ClientConfigProperties.INPUT_OUTPUT_FORMAT.getKey()), "CSV"); + Assert.assertEquals(config.get(ClientConfigProperties.COMPRESSION_METHOD.getKey()), CompressionMethod.ZSTD.name()); + // Add new assertion for default value } } @@ -419,50 +422,6 @@ public void testSocketBufferIsUnsetUnlessConfigured(ClientConfigProperties optio } } - @Test(groups = {"integration"}) - public void testWithOldDefaults() { - try (Client client = new Client.Builder() - .setUsername("default") - .setPassword("seceret") - .addEndpoint("http://localhost:8123") - .setDefaultDatabase("default") - .setExecutionTimeout(0, MILLIS) - .setLZ4UncompressedBufferSize(ClickHouseLZ4OutputStream.UNCOMPRESSED_BUFF_SIZE) - .disableNativeCompression(false) - .useServerTimeZone(true) - .setServerTimeZone("UTC") - .useAsyncRequests(false) - .setMaxConnections(10) - .setConnectionRequestTimeout(10, SECONDS) - .setConnectionReuseStrategy(ConnectionReuseStrategy.FIFO) - .enableConnectionPool(true) - .setConnectionTTL(-1, MILLIS) - .retryOnFailures(ClientFaultCause.NoHttpResponse, ClientFaultCause.ConnectTimeout, - ClientFaultCause.ConnectionRequestTimeout, ClientFaultCause.ServerRetryable) - .setClientNetworkBufferSize(300_000) - .setMaxRetries(3) - .allowBinaryReaderToReuseBuffers(false) - .columnToMethodMatchingStrategy(DefaultColumnToMethodMatchingStrategy.INSTANCE) - .useHTTPBasicAuth(true) - .compressClientRequest(false) - .compressServerResponse(true) - .useHttpCompression(false) - .appCompressedData(false) - .setSocketTimeout(0, SECONDS) - .setSocketRcvbuf(804800) - .setSocketSndbuf(804800) - .build()) { - Map config = client.getConfiguration(); - for (ClientConfigProperties p : ClientConfigProperties.values()) { - if (p.getDefaultValue() != null) { - Assert.assertTrue(config.containsKey(p.getKey()), "Default value should be set for " + p.getKey()); - Assert.assertEquals(config.get(p.getKey()), p.getDefaultValue(), "Default value doesn't match"); - } - } - Assert.assertEquals(config.size(), 38); // to check everything is set. Increment when new added. - } - } - @DataProvider(name = "sessionRoles") private static Object[][] sessionRoles() { return new Object[][]{ From ba0b7ca26a0e7030011f282fbff3d83b9c7216db Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Mon, 28 Sep 2026 14:09:54 -0700 Subject: [PATCH 3/9] Updated tests to run them with 2 main compression methods: LZ4 and ZSTD --- ...sertClientContentLZ4CompressionTests.java} | 169 +++++++++--------- ...sertClientContentZSTDCompressionTests.java | 85 +++++++++ .../InsertClientHttpCompressionTests.java | 3 +- .../clickhouse/client/insert/InsertTests.java | 6 +- .../query/BinaryReadyReusesBuffersTests.java | 4 +- .../QueryServerContentCompressionTests.java | 8 - ...QueryServerContentLZ4CompressionTests.java | 10 ++ ...ueryServerContentZSTDCompressionTests.java | 10 ++ .../QueryServerHttpCompressionTests.java | 3 +- .../clickhouse/client/query/QueryTests.java | 10 +- 10 files changed, 209 insertions(+), 99 deletions(-) rename client-v2/src/test/java/com/clickhouse/client/insert/{InsertClientContentCompressionTests.java => InsertClientContentLZ4CompressionTests.java} (94%) create mode 100644 client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentZSTDCompressionTests.java delete mode 100644 client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentCompressionTests.java create mode 100644 client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentLZ4CompressionTests.java create mode 100644 client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentZSTDCompressionTests.java diff --git a/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentCompressionTests.java b/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentLZ4CompressionTests.java similarity index 94% rename from client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentCompressionTests.java rename to client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentLZ4CompressionTests.java index 67d7f23ef..772ffb66c 100644 --- a/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentCompressionTests.java +++ b/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentLZ4CompressionTests.java @@ -1,84 +1,85 @@ -package com.clickhouse.client.insert; - -import com.clickhouse.client.ClickHouseNode; -import com.clickhouse.client.ClickHouseProtocol; -import com.clickhouse.client.ClickHouseServerForTest; -import com.clickhouse.client.api.Client; -import com.clickhouse.client.api.data_formats.ClickHouseBinaryFormatReader; -import com.clickhouse.client.api.insert.InsertResponse; -import com.clickhouse.client.api.insert.InsertSettings; -import com.clickhouse.client.api.internal.ServerSettings; -import com.clickhouse.client.api.query.GenericRecord; -import com.clickhouse.client.api.query.QueryResponse; -import org.apache.commons.lang3.RandomStringUtils; -import org.testng.Assert; -import org.testng.annotations.Test; - -import java.util.Collections; -import java.util.List; -import java.util.UUID; -import java.util.concurrent.TimeUnit; - -public class InsertClientContentCompressionTests extends InsertTests { - public InsertClientContentCompressionTests() { - super(true, false); - } - - - @Test(groups = { "integration" }) - public void testInsertAndReadBackWithSecureConnection() { - if (isCloud()) { - return; - } - ClickHouseNode secureServer = getSecureServer(ClickHouseProtocol.HTTP); - - try (Client client = new Client.Builder() - .addEndpoint("https://localhost:" + secureServer.getPort()) - .setUsername("default") - .setPassword(ClickHouseServerForTest.getPassword()) - .setRootCertificate("containers/clickhouse-server/certs/localhost.crt") - .setDefaultDatabase(ClickHouseServerForTest.getDatabase()) - .compressClientRequest(true) - .build()) { - final String tableName = "single_pojo_table"; - final String createSQL = SamplePOJO.generateTableCreateSQL(tableName); - final SamplePOJO pojo = new SamplePOJO(); - - initTable(tableName, createSQL); - - client.register(SamplePOJO.class, client.getTableSchema(tableName)); - InsertSettings settings = new InsertSettings() - .setDeduplicationToken(RandomStringUtils.randomAlphabetic(36)) - .setQueryId(String.valueOf(UUID.randomUUID())) - .serverSetting(ServerSettings.ASYNC_INSERT, "0"); - System.out.println("Inserting POJO: " + pojo); - try (InsertResponse response = client.insert(tableName, Collections.singletonList(pojo), settings).get(10, TimeUnit.SECONDS)) { - Assert.assertEquals(response.getWrittenRows(), 1); - } - - try (QueryResponse queryResponse = - client.query("SELECT * FROM " + tableName + " LIMIT 1").get(10, TimeUnit.SECONDS)) { - - ClickHouseBinaryFormatReader reader = client.newBinaryFormatReader(queryResponse); - Assert.assertNotNull(reader.next()); - - Assert.assertEquals(reader.getByte("byteValue"), pojo.getByteValue()); - Assert.assertEquals(reader.getByte("int8"), pojo.getInt8()); - Assert.assertEquals(reader.getShort("uint8"), pojo.getUint8()); - Assert.assertEquals(reader.getShort("int16"), pojo.getInt16()); - Assert.assertEquals(reader.getInteger("int32"), pojo.getInt32()); - Assert.assertEquals(reader.getLong("int64"), pojo.getInt64()); - Assert.assertEquals(reader.getFloat("float32"), pojo.getFloat32()); - Assert.assertEquals(reader.getDouble("float64"), pojo.getFloat64()); - Assert.assertEquals(reader.getString("string"), pojo.getString()); - Assert.assertEquals(reader.getString("fixedString"), pojo.getFixedString()); - } - List records = client.queryAll("SELECT timezone()"); - Assert.assertTrue(records.size() > 0); - Assert.assertEquals(records.get(0).getString(1), "UTC"); - } catch (Exception e) { - e.printStackTrace(); - Assert.fail(e.getMessage()); - } - } -} +package com.clickhouse.client.insert; + +import com.clickhouse.client.ClickHouseNode; +import com.clickhouse.client.ClickHouseProtocol; +import com.clickhouse.client.ClickHouseServerForTest; +import com.clickhouse.client.api.Client; +import com.clickhouse.client.api.CompressionMethod; +import com.clickhouse.client.api.data_formats.ClickHouseBinaryFormatReader; +import com.clickhouse.client.api.insert.InsertResponse; +import com.clickhouse.client.api.insert.InsertSettings; +import com.clickhouse.client.api.internal.ServerSettings; +import com.clickhouse.client.api.query.GenericRecord; +import com.clickhouse.client.api.query.QueryResponse; +import org.apache.commons.lang3.RandomStringUtils; +import org.testng.Assert; +import org.testng.annotations.Test; + +import java.util.Collections; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.TimeUnit; + +public class InsertClientContentLZ4CompressionTests extends InsertTests { + public InsertClientContentLZ4CompressionTests() { + super(true, false, CompressionMethod.LZ4); + } + + + @Test(groups = { "integration" }) + public void testInsertAndReadBackWithSecureConnection() { + if (isCloud()) { + return; + } + ClickHouseNode secureServer = getSecureServer(ClickHouseProtocol.HTTP); + + try (Client client = new Client.Builder() + .addEndpoint("https://localhost:" + secureServer.getPort()) + .setUsername("default") + .setPassword(ClickHouseServerForTest.getPassword()) + .setRootCertificate("containers/clickhouse-server/certs/localhost.crt") + .setDefaultDatabase(ClickHouseServerForTest.getDatabase()) + .compressClientRequest(true) + .build()) { + final String tableName = "single_pojo_table"; + final String createSQL = SamplePOJO.generateTableCreateSQL(tableName); + final SamplePOJO pojo = new SamplePOJO(); + + initTable(tableName, createSQL); + + client.register(SamplePOJO.class, client.getTableSchema(tableName)); + InsertSettings settings = new InsertSettings() + .setDeduplicationToken(RandomStringUtils.randomAlphabetic(36)) + .setQueryId(String.valueOf(UUID.randomUUID())) + .serverSetting(ServerSettings.ASYNC_INSERT, "0"); + System.out.println("Inserting POJO: " + pojo); + try (InsertResponse response = client.insert(tableName, Collections.singletonList(pojo), settings).get(10, TimeUnit.SECONDS)) { + Assert.assertEquals(response.getWrittenRows(), 1); + } + + try (QueryResponse queryResponse = + client.query("SELECT * FROM " + tableName + " LIMIT 1").get(10, TimeUnit.SECONDS)) { + + ClickHouseBinaryFormatReader reader = client.newBinaryFormatReader(queryResponse); + Assert.assertNotNull(reader.next()); + + Assert.assertEquals(reader.getByte("byteValue"), pojo.getByteValue()); + Assert.assertEquals(reader.getByte("int8"), pojo.getInt8()); + Assert.assertEquals(reader.getShort("uint8"), pojo.getUint8()); + Assert.assertEquals(reader.getShort("int16"), pojo.getInt16()); + Assert.assertEquals(reader.getInteger("int32"), pojo.getInt32()); + Assert.assertEquals(reader.getLong("int64"), pojo.getInt64()); + Assert.assertEquals(reader.getFloat("float32"), pojo.getFloat32()); + Assert.assertEquals(reader.getDouble("float64"), pojo.getFloat64()); + Assert.assertEquals(reader.getString("string"), pojo.getString()); + Assert.assertEquals(reader.getString("fixedString"), pojo.getFixedString()); + } + List records = client.queryAll("SELECT timezone()"); + Assert.assertTrue(records.size() > 0); + Assert.assertEquals(records.get(0).getString(1), "UTC"); + } catch (Exception e) { + e.printStackTrace(); + Assert.fail(e.getMessage()); + } + } +} diff --git a/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentZSTDCompressionTests.java b/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentZSTDCompressionTests.java new file mode 100644 index 000000000..9501340da --- /dev/null +++ b/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientContentZSTDCompressionTests.java @@ -0,0 +1,85 @@ +package com.clickhouse.client.insert; + +import com.clickhouse.client.ClickHouseNode; +import com.clickhouse.client.ClickHouseProtocol; +import com.clickhouse.client.ClickHouseServerForTest; +import com.clickhouse.client.api.Client; +import com.clickhouse.client.api.CompressionMethod; +import com.clickhouse.client.api.data_formats.ClickHouseBinaryFormatReader; +import com.clickhouse.client.api.insert.InsertResponse; +import com.clickhouse.client.api.insert.InsertSettings; +import com.clickhouse.client.api.internal.ServerSettings; +import com.clickhouse.client.api.query.GenericRecord; +import com.clickhouse.client.api.query.QueryResponse; +import org.apache.commons.lang3.RandomStringUtils; +import org.testng.Assert; +import org.testng.annotations.Test; + +import java.util.Collections; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.TimeUnit; + +public class InsertClientContentZSTDCompressionTests extends InsertTests { + public InsertClientContentZSTDCompressionTests() { + super(true, false, CompressionMethod.ZSTD); + } + + + @Test(groups = { "integration" }) + public void testInsertAndReadBackWithSecureConnection() { + if (isCloud()) { + return; + } + ClickHouseNode secureServer = getSecureServer(ClickHouseProtocol.HTTP); + + try (Client client = new Client.Builder() + .addEndpoint("https://localhost:" + secureServer.getPort()) + .setUsername("default") + .setPassword(ClickHouseServerForTest.getPassword()) + .setRootCertificate("containers/clickhouse-server/certs/localhost.crt") + .setDefaultDatabase(ClickHouseServerForTest.getDatabase()) + .compressClientRequest(true) + .build()) { + final String tableName = "single_pojo_table"; + final String createSQL = SamplePOJO.generateTableCreateSQL(tableName); + final SamplePOJO pojo = new SamplePOJO(); + + initTable(tableName, createSQL); + + client.register(SamplePOJO.class, client.getTableSchema(tableName)); + InsertSettings settings = new InsertSettings() + .setDeduplicationToken(RandomStringUtils.randomAlphabetic(36)) + .setQueryId(String.valueOf(UUID.randomUUID())) + .serverSetting(ServerSettings.ASYNC_INSERT, "0"); + System.out.println("Inserting POJO: " + pojo); + try (InsertResponse response = client.insert(tableName, Collections.singletonList(pojo), settings).get(10, TimeUnit.SECONDS)) { + Assert.assertEquals(response.getWrittenRows(), 1); + } + + try (QueryResponse queryResponse = + client.query("SELECT * FROM " + tableName + " LIMIT 1").get(10, TimeUnit.SECONDS)) { + + ClickHouseBinaryFormatReader reader = client.newBinaryFormatReader(queryResponse); + Assert.assertNotNull(reader.next()); + + Assert.assertEquals(reader.getByte("byteValue"), pojo.getByteValue()); + Assert.assertEquals(reader.getByte("int8"), pojo.getInt8()); + Assert.assertEquals(reader.getShort("uint8"), pojo.getUint8()); + Assert.assertEquals(reader.getShort("int16"), pojo.getInt16()); + Assert.assertEquals(reader.getInteger("int32"), pojo.getInt32()); + Assert.assertEquals(reader.getLong("int64"), pojo.getInt64()); + Assert.assertEquals(reader.getFloat("float32"), pojo.getFloat32()); + Assert.assertEquals(reader.getDouble("float64"), pojo.getFloat64()); + Assert.assertEquals(reader.getString("string"), pojo.getString()); + Assert.assertEquals(reader.getString("fixedString"), pojo.getFixedString()); + } + List records = client.queryAll("SELECT timezone()"); + Assert.assertTrue(records.size() > 0); + Assert.assertEquals(records.get(0).getString(1), "UTC"); + } catch (Exception e) { + e.printStackTrace(); + Assert.fail(e.getMessage()); + } + } +} diff --git a/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientHttpCompressionTests.java b/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientHttpCompressionTests.java index 18a8f5d3f..f272b5ee4 100644 --- a/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientHttpCompressionTests.java +++ b/client-v2/src/test/java/com/clickhouse/client/insert/InsertClientHttpCompressionTests.java @@ -1,5 +1,6 @@ package com.clickhouse.client.insert; +import com.clickhouse.client.api.CompressionMethod; import com.clickhouse.client.api.insert.InsertResponse; import com.clickhouse.client.api.insert.InsertSettings; import com.clickhouse.client.api.metrics.OperationMetrics; @@ -20,7 +21,7 @@ public class InsertClientHttpCompressionTests extends InsertTests { public InsertClientHttpCompressionTests() { - super(true, true); + super(true, true, CompressionMethod.ZSTD); } diff --git a/client-v2/src/test/java/com/clickhouse/client/insert/InsertTests.java b/client-v2/src/test/java/com/clickhouse/client/insert/InsertTests.java index cf9eb057b..b08ef5587 100644 --- a/client-v2/src/test/java/com/clickhouse/client/insert/InsertTests.java +++ b/client-v2/src/test/java/com/clickhouse/client/insert/InsertTests.java @@ -70,14 +70,17 @@ public class InsertTests extends BaseIntegrationTest { private boolean useHttpCompression = false; + private CompressionMethod compressionMethod = CompressionMethod.LZ4; + static final int EXECUTE_CMD_TIMEOUT = 10; // seconds InsertTests() { } - public InsertTests(boolean useClientCompression, boolean useHttpCompression) { + public InsertTests(boolean useClientCompression, boolean useHttpCompression, CompressionMethod compressionMethod) { this.useClientCompression = useClientCompression; this.useHttpCompression = useHttpCompression; + this.compressionMethod = compressionMethod; } @BeforeMethod(groups = { "integration" }) @@ -104,6 +107,7 @@ protected Client.Builder newClient() { .setPassword(ClickHouseServerForTest.getPassword()) .compressClientRequest(useClientCompression) .useHttpCompression(useHttpCompression) + .compressionMethod(compressionMethod) .setDefaultDatabase(ClickHouseServerForTest.getDatabase()) .serverSetting(ServerSettings.ASYNC_INSERT, "0") .serverSetting(ServerSettings.WAIT_END_OF_QUERY, "1") diff --git a/client-v2/src/test/java/com/clickhouse/client/query/BinaryReadyReusesBuffersTests.java b/client-v2/src/test/java/com/clickhouse/client/query/BinaryReadyReusesBuffersTests.java index 714275a12..3e493a6c3 100644 --- a/client-v2/src/test/java/com/clickhouse/client/query/BinaryReadyReusesBuffersTests.java +++ b/client-v2/src/test/java/com/clickhouse/client/query/BinaryReadyReusesBuffersTests.java @@ -1,8 +1,10 @@ package com.clickhouse.client.query; +import com.clickhouse.client.api.CompressionMethod; + public class BinaryReadyReusesBuffersTests extends QueryTests { public BinaryReadyReusesBuffersTests() { - super(false, false, true); + super(false, false, true, CompressionMethod.LZ4); } } diff --git a/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentCompressionTests.java b/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentCompressionTests.java deleted file mode 100644 index 1001b08be..000000000 --- a/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentCompressionTests.java +++ /dev/null @@ -1,8 +0,0 @@ -package com.clickhouse.client.query; - -public class QueryServerContentCompressionTests extends QueryTests { - - QueryServerContentCompressionTests() { - super(true, false); - } -} diff --git a/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentLZ4CompressionTests.java b/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentLZ4CompressionTests.java new file mode 100644 index 000000000..4b0968dc4 --- /dev/null +++ b/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentLZ4CompressionTests.java @@ -0,0 +1,10 @@ +package com.clickhouse.client.query; + +import com.clickhouse.client.api.CompressionMethod; + +public class QueryServerContentLZ4CompressionTests extends QueryTests { + + QueryServerContentLZ4CompressionTests() { + super(true, false, CompressionMethod.LZ4); + } +} diff --git a/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentZSTDCompressionTests.java b/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentZSTDCompressionTests.java new file mode 100644 index 000000000..5274806a7 --- /dev/null +++ b/client-v2/src/test/java/com/clickhouse/client/query/QueryServerContentZSTDCompressionTests.java @@ -0,0 +1,10 @@ +package com.clickhouse.client.query; + +import com.clickhouse.client.api.CompressionMethod; + +public class QueryServerContentZSTDCompressionTests extends QueryTests { + + QueryServerContentZSTDCompressionTests() { + super(true, false, CompressionMethod.ZSTD); + } +} diff --git a/client-v2/src/test/java/com/clickhouse/client/query/QueryServerHttpCompressionTests.java b/client-v2/src/test/java/com/clickhouse/client/query/QueryServerHttpCompressionTests.java index 57e9cad16..0c049b5bd 100644 --- a/client-v2/src/test/java/com/clickhouse/client/query/QueryServerHttpCompressionTests.java +++ b/client-v2/src/test/java/com/clickhouse/client/query/QueryServerHttpCompressionTests.java @@ -1,5 +1,6 @@ package com.clickhouse.client.query; +import com.clickhouse.client.api.CompressionMethod; import com.clickhouse.client.api.data_formats.internal.BinaryStreamReader; import com.clickhouse.client.api.query.GenericRecord; import com.clickhouse.client.api.query.QuerySettings; @@ -14,7 +15,7 @@ public class QueryServerHttpCompressionTests extends QueryTests { QueryServerHttpCompressionTests() { - super(true, true); + super(true, true, CompressionMethod.ZSTD); } diff --git a/client-v2/src/test/java/com/clickhouse/client/query/QueryTests.java b/client-v2/src/test/java/com/clickhouse/client/query/QueryTests.java index d26f59764..0b01d0f3c 100644 --- a/client-v2/src/test/java/com/clickhouse/client/query/QueryTests.java +++ b/client-v2/src/test/java/com/clickhouse/client/query/QueryTests.java @@ -9,6 +9,7 @@ import com.clickhouse.client.api.Client; import com.clickhouse.client.api.ClientConfigProperties; import com.clickhouse.client.api.ClientException; +import com.clickhouse.client.api.CompressionMethod; import com.clickhouse.client.api.ServerException; import com.clickhouse.client.api.command.CommandSettings; import com.clickhouse.client.api.data_formats.ClickHouseBinaryFormatReader; @@ -104,19 +105,22 @@ public class QueryTests extends BaseIntegrationTest { private boolean useHttpCompression = false; + private CompressionMethod compressionMethod; + private boolean usePreallocatedBuffers = false; QueryTests(){ } - public QueryTests(boolean useServerCompression, boolean useHttpCompression) { - this(useServerCompression, useHttpCompression, false); + public QueryTests(boolean useServerCompression, boolean useHttpCompression, CompressionMethod compressionMethod) { + this(useServerCompression, useHttpCompression, false, compressionMethod); } - public QueryTests(boolean useServerCompression, boolean useHttpCompression, boolean usePreallocatedBuffers) { + public QueryTests(boolean useServerCompression, boolean useHttpCompression, boolean usePreallocatedBuffers, CompressionMethod compressionMethod) { this.useServerCompression = useServerCompression; this.useHttpCompression = useHttpCompression; this.usePreallocatedBuffers = usePreallocatedBuffers; + this.compressionMethod = compressionMethod; } @BeforeMethod(groups = {"integration"}) From 3069a432682366f3de8120b72dbd217497c1ed01 Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Mon, 28 Sep 2026 22:08:09 -0700 Subject: [PATCH 4/9] cleaned up constants and fixed v1 tests --- .../client/ClientIntegrationTest.java | 82 ++++++++++++++++++- .../http/ApacheHttpConnectionImplTest.java | 7 +- .../client/http/ClickHouseHttpClientTest.java | 5 +- .../internal/CompressedBlockInputStream.java | 5 +- .../internal/CompressedBlockOutputStream.java | 11 ++- 5 files changed, 94 insertions(+), 16 deletions(-) diff --git a/clickhouse-client/src/test/java/com/clickhouse/client/ClientIntegrationTest.java b/clickhouse-client/src/test/java/com/clickhouse/client/ClientIntegrationTest.java index f1b9931da..07333646d 100644 --- a/clickhouse-client/src/test/java/com/clickhouse/client/ClientIntegrationTest.java +++ b/clickhouse-client/src/test/java/com/clickhouse/client/ClientIntegrationTest.java @@ -160,7 +160,7 @@ private void setClientOptions(ClickHouseRequest request) { protected abstract Class getClientClass(); protected Map getClientOptions() { - return Collections.emptyMap(); + return Collections.singletonMap(ClickHouseClientOption.CUSTOM_SETTINGS, "network_compression_method=lz4"); } protected ClickHouseClientBuilder initClient(ClickHouseClientBuilder builder) { @@ -168,17 +168,91 @@ protected ClickHouseClientBuilder initClient(ClickHouseClientBuilder builder) { } protected ClickHouseClient getClient(ClickHouseConfig... configs) { - return initClient(ClickHouseClient.builder()).config(new ClickHouseConfig(configs)) + Map defaultOptions = new HashMap<>(getClientOptions()); + ClickHouseConfig baseConfig = new ClickHouseConfig(defaultOptions); + List 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) { - return initClient(ClickHouseClient.builder()) - .config(new ClickHouseConfig(configs)) + Map defaultOptions = new HashMap<>(getClientOptions()); + ClickHouseConfig baseConfig = new ClickHouseConfig(defaultOptions); + List 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 options) { + return addCustomSettings(super.getServer(protocol, options)); + } + protected ClickHouseNode getSecureServer(ClickHouseNode base) { return getSecureServer(getProtocol(), base); } diff --git a/clickhouse-http-client/src/test/java/com/clickhouse/client/http/ApacheHttpConnectionImplTest.java b/clickhouse-http-client/src/test/java/com/clickhouse/client/http/ApacheHttpConnectionImplTest.java index 9774c745d..33d922e09 100644 --- a/clickhouse-http-client/src/test/java/com/clickhouse/client/http/ApacheHttpConnectionImplTest.java +++ b/clickhouse-http-client/src/test/java/com/clickhouse/client/http/ApacheHttpConnectionImplTest.java @@ -74,8 +74,9 @@ public boolean supports(Class clazz) { @Override protected Map getClientOptions() { - return Collections.singletonMap(ClickHouseHttpOption.CONNECTION_PROVIDER, - HttpConnectionProvider.APACHE_HTTP_CLIENT); + Map options = new HashMap<>(super.getClientOptions()); + options.put(ClickHouseHttpOption.CONNECTION_PROVIDER, HttpConnectionProvider.APACHE_HTTP_CLIENT); + return options; } @Test(groups = { "unit" }, dataProvider = "replicaTags") @@ -351,7 +352,7 @@ public void testConnectionTTL(Map options, int o proxy.addStubMapping(WireMock.post(WireMock.anyUrl()) .willReturn(WireMock.aResponse().proxiedFrom(targetURI.build().toString())).build()); - Map baseOptions = new HashMap<>(); + Map baseOptions = new HashMap<>(getClientOptions()); baseOptions.put(ClickHouseClientOption.PROXY_PORT, proxyPort); baseOptions.put(ClickHouseClientOption.PROXY_HOST, "localhost"); baseOptions.put(ClickHouseClientOption.PROXY_TYPE, ClickHouseProxyType.HTTP); diff --git a/clickhouse-http-client/src/test/java/com/clickhouse/client/http/ClickHouseHttpClientTest.java b/clickhouse-http-client/src/test/java/com/clickhouse/client/http/ClickHouseHttpClientTest.java index dd26a1ecd..3c7e59f8b 100644 --- a/clickhouse-http-client/src/test/java/com/clickhouse/client/http/ClickHouseHttpClientTest.java +++ b/clickhouse-http-client/src/test/java/com/clickhouse/client/http/ClickHouseHttpClientTest.java @@ -103,8 +103,9 @@ protected Class getClientClass() { @Override protected Map getClientOptions() { - return Collections.singletonMap(ClickHouseHttpOption.CONNECTION_PROVIDER, - HttpConnectionProvider.HTTP_URL_CONNECTION); + Map options = new HashMap<>(super.getClientOptions()); + options.put(ClickHouseHttpOption.CONNECTION_PROVIDER, HttpConnectionProvider.HTTP_URL_CONNECTION); + return options; } @Test(groups = { "integration" }) diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java index bfdb095f3..6a7d7726e 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java @@ -65,12 +65,11 @@ public int read(byte[] b, int off, int len) throws IOException { return readBytes == 0 ? -1 : readBytes; } - - static final byte MAGIC = (byte) 0x82; static final byte MAGIC_LZ4 = (byte) 0x82; static final byte MAGIC_ZSTD_3 = (byte) 0x90; static final byte MAGIC_NONE = (byte) 0x02; static final int HEADER_LENGTH = 25; + static final int MAGIC_NUM_POS = 16; final byte[] headerBuff = new byte[HEADER_LENGTH]; @@ -111,7 +110,7 @@ private int refill() throws IOException { return -1; } - byte magicNumber = headerBuff[16]; + byte magicNumber = headerBuff[MAGIC_NUM_POS]; switch (magicNumber) { case MAGIC_LZ4: case MAGIC_ZSTD_3: diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java index 157cf6fb2..3791d0019 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java @@ -125,8 +125,8 @@ public LZ4OutputStream(OutputStream out, LZ4Compressor compressor, int bufferSiz @Override protected int compressData(ByteBuffer inBuffer, int uncompressedLen, ByteBuffer compressedBuffer, ByteBuffer buffer) { - return compressor.compress(inBuffer, 0, uncompressedLen, compressedBuffer, 25, - compressedBuffer.remaining() - 25); + return compressor.compress(inBuffer, 0, uncompressedLen, compressedBuffer, CompressedBlockInputStream.HEADER_LENGTH, + compressedBuffer.remaining() - CompressedBlockInputStream.HEADER_LENGTH); } @Override @@ -137,14 +137,17 @@ protected byte getMagicNumber() { public static class ZSTDOutputStream extends CompressedBlockOutputStream { + private static final int COMPRESSION_LEVEL = 3; + public ZSTDOutputStream(OutputStream out, int bufferSize) { super(out, bufferSize); } @Override protected int compressData(ByteBuffer inBuffer, int uncompressedLen, ByteBuffer compressedBuffer, ByteBuffer buffer) { - return (int) Zstd.compressByteArray(compressedBuffer.array(), 25, compressedBuffer.remaining() - 25, - inBuffer.array(), 0, uncompressedLen, 3); + return (int) Zstd.compressByteArray(compressedBuffer.array(), CompressedBlockInputStream.HEADER_LENGTH, + compressedBuffer.remaining() - CompressedBlockInputStream.HEADER_LENGTH, + inBuffer.array(), 0, uncompressedLen, COMPRESSION_LEVEL); } @Override From fce6ed39c35b94efa9b71f7ab9bc1123bca3d870 Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Tue, 29 Sep 2026 06:42:01 -0700 Subject: [PATCH 5/9] Fixed the offset issues in LZ4 decompression code --- .../internal/CompressedBlockInputStream.java | 5 +- .../CompressedBlockInputStreamTest.java | 240 +++++++++++++++++- 2 files changed, 241 insertions(+), 4 deletions(-) diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java index 6a7d7726e..5f749fb64 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockInputStream.java @@ -149,8 +149,7 @@ private int refill() throws IOException { switch (magicNumber) { case MAGIC_LZ4: { - ByteBuffer blockBuff = ByteBuffer.wrap(block, offset, remaining); - decompressor.decompress(blockBuff, offset, buffer, 0, uncompressedSize); + decompressor.decompress(block, offset, buffer.array(), 0, uncompressedSize); break; } case MAGIC_ZSTD_3: { @@ -159,7 +158,7 @@ private int refill() throws IOException { } case MAGIC_NONE: // block is not compressed - just put it. - buffer.put(block, 0, uncompressedSize); + System.arraycopy(block, offset, buffer.array(), 0, uncompressedSize); break; default: throw new ClientException("bug: Should not reach here"); diff --git a/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java b/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java index 42890eb39..1be4b797e 100644 --- a/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java +++ b/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java @@ -1,20 +1,31 @@ package com.clickhouse.client.api.internal; import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.EOFException; import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.util.Arrays; +import com.clickhouse.client.api.ClientException; +import com.clickhouse.data.ClickHouseByteUtils; +import com.clickhouse.data.ClickHouseCityHash; +import com.github.luben.zstd.Zstd; import net.jpountz.lz4.LZ4Factory; import org.testng.Assert; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; public class CompressedBlockInputStreamTest { + private static final LZ4Factory LZ4_FACTORY = LZ4Factory.fastestInstance(); + @Test public void reportsActualByteCountsForTruncatedHeader() { byte[] truncatedHeader = new byte[10]; CompressedBlockInputStream input = new CompressedBlockInputStream( new ByteArrayInputStream(truncatedHeader), - LZ4Factory.fastestJavaInstance().fastDecompressor(), + LZ4_FACTORY.fastDecompressor(), 8192); IOException exception = Assert.expectThrows(IOException.class, @@ -22,4 +33,231 @@ public void reportsActualByteCountsForTruncatedHeader() { Assert.assertEquals(exception.getMessage(), "Incomplete read: 10 of 25"); } + + @DataProvider(name = "compressionAlgorithms") + public static Object[][] compressionAlgorithms() { + return new Object[][] { + { CompressedBlockInputStream.MAGIC_LZ4 }, + { CompressedBlockInputStream.MAGIC_ZSTD_3 }, + { CompressedBlockInputStream.MAGIC_NONE } + }; + } + + @Test(dataProvider = "compressionAlgorithms") + public void testBlockDecompression(byte magic) throws IOException { + byte[] payload = "Hello ClickHouse, testing LZ4, ZSTD and uncompressed blocks!".getBytes(StandardCharsets.UTF_8); + byte[] block = createBlock(magic, payload); + + CompressedBlockInputStream input = new CompressedBlockInputStream( + new ByteArrayInputStream(block), + LZ4_FACTORY.fastDecompressor(), + 1024); + + ByteArrayOutputStream out = new ByteArrayOutputStream(); + byte[] buf = new byte[11]; + int n; + while ((n = input.read(buf, 0, buf.length)) != -1) { + out.write(buf, 0, n); + } + + Assert.assertEquals(out.toByteArray(), payload); + } + + @Test(dataProvider = "compressionAlgorithms") + public void testReadByteByByte(byte magic) throws IOException { + byte[] payload = generateData(256); + byte[] block = createBlock(magic, payload); + + CompressedBlockInputStream input = new CompressedBlockInputStream( + new ByteArrayInputStream(block), + LZ4_FACTORY.fastDecompressor(), + 1024); + + ByteArrayOutputStream out = new ByteArrayOutputStream(); + int b; + while ((b = input.read()) != -1) { + out.write(b); + } + + Assert.assertEquals(out.toByteArray(), payload); + } + + @Test(dataProvider = "compressionAlgorithms") + public void testMultiBlockStream(byte magic) throws IOException { + byte[] part1 = generateData(500); + byte[] part2 = generateData(700); + + ByteArrayOutputStream combined = new ByteArrayOutputStream(); + combined.write(createBlock(magic, part1)); + combined.write(createBlock(magic, part2)); + + CompressedBlockInputStream input = new CompressedBlockInputStream( + new ByteArrayInputStream(combined.toByteArray()), + LZ4_FACTORY.fastDecompressor(), + 256); + + ByteArrayOutputStream out = new ByteArrayOutputStream(); + byte[] buf = new byte[64]; + int n; + while ((n = input.read(buf, 0, buf.length)) != -1) { + out.write(buf, 0, n); + } + + byte[] expected = new byte[part1.length + part2.length]; + System.arraycopy(part1, 0, expected, 0, part1.length); + System.arraycopy(part2, 0, expected, part1.length, part2.length); + + Assert.assertEquals(out.toByteArray(), expected); + } + + @Test + public void testRoundTripWithOutputStream() throws IOException { + byte[] payload = generateData(4096); + + // LZ4 + ByteArrayOutputStream lz4Out = new ByteArrayOutputStream(); + try (CompressedBlockOutputStream out = new CompressedBlockOutputStream.LZ4OutputStream( + lz4Out, LZ4_FACTORY.fastCompressor(), 1024)) { + out.write(payload); + } + try (CompressedBlockInputStream in = new CompressedBlockInputStream( + new ByteArrayInputStream(lz4Out.toByteArray()), LZ4_FACTORY.fastDecompressor(), 512)) { + ByteArrayOutputStream decompressed = new ByteArrayOutputStream(); + byte[] buf = new byte[128]; + int n; + while ((n = in.read(buf)) != -1) { + decompressed.write(buf, 0, n); + } + Assert.assertEquals(decompressed.toByteArray(), payload); + } + + // ZSTD + ByteArrayOutputStream zstdOut = new ByteArrayOutputStream(); + try (CompressedBlockOutputStream out = new CompressedBlockOutputStream.ZSTDOutputStream(zstdOut, 1024)) { + out.write(payload); + } + try (CompressedBlockInputStream in = new CompressedBlockInputStream( + new ByteArrayInputStream(zstdOut.toByteArray()), LZ4_FACTORY.fastDecompressor(), 512)) { + ByteArrayOutputStream decompressed = new ByteArrayOutputStream(); + byte[] buf = new byte[128]; + int n; + while ((n = in.read(buf)) != -1) { + decompressed.write(buf, 0, n); + } + Assert.assertEquals(decompressed.toByteArray(), payload); + } + } + + @Test + public void testInvalidMagicByte() { + byte[] payload = generateData(32); + byte[] block = createBlock((byte) 0x7F, payload); + + CompressedBlockInputStream input = new CompressedBlockInputStream( + new ByteArrayInputStream(block), + LZ4_FACTORY.fastDecompressor(), + 1024); + + ClientException ex = Assert.expectThrows(ClientException.class, + () -> input.read(new byte[1], 0, 1)); + Assert.assertTrue(ex.getMessage().contains("Invalid LZ4 magic byte")); + } + + @Test + public void testChecksumMismatch() { + byte[] payload = generateData(64); + byte[] block = createBlock(CompressedBlockInputStream.MAGIC_LZ4, payload); + block[block.length - 1] ^= 0xFF; + + CompressedBlockInputStream input = new CompressedBlockInputStream( + new ByteArrayInputStream(block), + LZ4_FACTORY.fastDecompressor(), + 1024); + + ClientException ex = Assert.expectThrows(ClientException.class, + () -> input.read(new byte[1], 0, 1)); + Assert.assertTrue(ex.getMessage().contains("checksum mismatch")); + } + + @Test + public void testUnexpectedEndOfStreamInBlock() { + byte[] payload = generateData(64); + byte[] block = createBlock(CompressedBlockInputStream.MAGIC_LZ4, payload); + byte[] truncated = Arrays.copyOf(block, 25); // header only, 0 bytes of payload + + CompressedBlockInputStream input = new CompressedBlockInputStream( + new ByteArrayInputStream(truncated), + LZ4_FACTORY.fastDecompressor(), + 1024); + + EOFException ex = Assert.expectThrows(EOFException.class, () -> input.read(new byte[1], 0, 1)); + Assert.assertEquals(ex.getMessage(), "Unexpected end of stream"); + } + + @Test + public void testIncompleteReadInBlockPayload() { + byte[] payload = generateData(64); + byte[] block = createBlock(CompressedBlockInputStream.MAGIC_LZ4, payload); + byte[] truncated = Arrays.copyOf(block, block.length - 5); + + CompressedBlockInputStream input = new CompressedBlockInputStream( + new ByteArrayInputStream(truncated), + LZ4_FACTORY.fastDecompressor(), + 1024); + + IOException ex = Assert.expectThrows(IOException.class, () -> input.read(new byte[1], 0, 1)); + Assert.assertTrue(ex.getMessage().startsWith("Incomplete read:")); + } + + @Test + public void testMissingLZ4Decompressor() { + byte[] payload = generateData(32); + byte[] block = createBlock(CompressedBlockInputStream.MAGIC_LZ4, payload); + + CompressedBlockInputStream input = new CompressedBlockInputStream( + new ByteArrayInputStream(block), + null, + 1024); + + ClientException ex = Assert.expectThrows(ClientException.class, + () -> input.read(new byte[1], 0, 1)); + Assert.assertTrue(ex.getMessage().contains("LZ4 decompressor is not configured")); + } + + private static byte[] generateData(int length) { + byte[] data = new byte[length]; + for (int i = 0; i < length; i++) { + data[i] = (byte) ('A' + (i % 26)); + } + return data; + } + + private static byte[] createBlock(byte magic, byte[] payload) { + int uncompressedSize = payload.length; + byte[] compressedPayload; + if (magic == CompressedBlockInputStream.MAGIC_LZ4) { + compressedPayload = LZ4_FACTORY.fastCompressor().compress(payload); + } else if (magic == CompressedBlockInputStream.MAGIC_ZSTD_3) { + compressedPayload = Zstd.compress(payload, 3); + } else if (magic == CompressedBlockInputStream.MAGIC_NONE) { + compressedPayload = payload; + } else { + compressedPayload = payload; + } + + int compressedSizeWithHeader = compressedPayload.length + 9; + byte[] block = new byte[compressedSizeWithHeader]; + block[0] = magic; + CompressedBlockInputStream.setInt32(block, 1, compressedSizeWithHeader); + CompressedBlockInputStream.setInt32(block, 5, uncompressedSize); + System.arraycopy(compressedPayload, 0, block, 9, compressedPayload.length); + + long[] checksum = ClickHouseCityHash.cityHash128(block, 0, compressedSizeWithHeader); + + byte[] rawBlock = new byte[16 + compressedSizeWithHeader]; + ClickHouseByteUtils.setInt64(rawBlock, 0, checksum[0]); + ClickHouseByteUtils.setInt64(rawBlock, 8, checksum[1]); + System.arraycopy(block, 0, rawBlock, 16, compressedSizeWithHeader); + return rawBlock; + } } From a8efd35ece908e6f994bad9d7995c87e0d549bd3 Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Tue, 29 Sep 2026 06:46:42 -0700 Subject: [PATCH 6/9] Fixed compress LZ4 buffer allocation --- .../client/api/internal/CompressedBlockOutputStream.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java index 3791d0019..750bdf80c 100644 --- a/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java +++ b/client-v2/src/main/java/com/clickhouse/client/api/internal/CompressedBlockOutputStream.java @@ -21,15 +21,12 @@ public abstract class CompressedBlockOutputStream extends OutputStream { private final ByteBuffer compressedBuffer; - private static final int HEADER_LEN = 15; // 9 bytes for header, 6 bytes for checksum - - public CompressedBlockOutputStream(OutputStream out, int bufferSize) { super(); LOG.debug("Using compressor with buffer size {}", bufferSize); this.inBuffer = ByteBuffer.allocate(bufferSize); this.out = out; - this.compressedBuffer = ByteBuffer.allocate(bufferSize + HEADER_LEN); + this.compressedBuffer = ByteBuffer.allocate(bufferSize + CompressedBlockInputStream.HEADER_LENGTH); } @Override From aa0271399aaf2f7e57d54e45dc56f559ed4cb3d3 Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Tue, 29 Sep 2026 06:54:30 -0700 Subject: [PATCH 7/9] removed useless test for null compressor --- .../internal/CompressedBlockInputStreamTest.java | 15 --------------- 1 file changed, 15 deletions(-) diff --git a/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java b/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java index 1be4b797e..e83537dcf 100644 --- a/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java +++ b/client-v2/src/test/java/com/clickhouse/client/api/internal/CompressedBlockInputStreamTest.java @@ -209,21 +209,6 @@ public void testIncompleteReadInBlockPayload() { Assert.assertTrue(ex.getMessage().startsWith("Incomplete read:")); } - @Test - public void testMissingLZ4Decompressor() { - byte[] payload = generateData(32); - byte[] block = createBlock(CompressedBlockInputStream.MAGIC_LZ4, payload); - - CompressedBlockInputStream input = new CompressedBlockInputStream( - new ByteArrayInputStream(block), - null, - 1024); - - ClientException ex = Assert.expectThrows(ClientException.class, - () -> input.read(new byte[1], 0, 1)); - Assert.assertTrue(ex.getMessage().contains("LZ4 decompressor is not configured")); - } - private static byte[] generateData(int length) { byte[] data = new byte[length]; for (int i = 0; i < length; i++) { From f7350ea54326d3bcb0d5c2f92b0761045d320e94 Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Tue, 29 Sep 2026 13:08:50 -0700 Subject: [PATCH 8/9] Fixed jdbc-v1 tests --- .../clickhouse/jdbc/AccessManagementTest.java | 8 ++--- .../clickhouse/jdbc/JdbcIntegrationTest.java | 36 ++++++++++++++++++- .../com/clickhouse/jdbc/JdbcIssuesTest.java | 8 ++--- .../comparison/DateTimeComparisonTest.java | 2 +- 4 files changed, 44 insertions(+), 10 deletions(-) diff --git a/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/AccessManagementTest.java b/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/AccessManagementTest.java index bba0cbb68..2ca8d89fa 100644 --- a/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/AccessManagementTest.java +++ b/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/AccessManagementTest.java @@ -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"); @@ -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"); @@ -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(); @@ -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(); diff --git a/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/JdbcIntegrationTest.java b/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/JdbcIntegrationTest.java index 97eea693a..1dd6cc1c7 100644 --- a/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/JdbcIntegrationTest.java +++ b/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/JdbcIntegrationTest.java @@ -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; @@ -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; } @@ -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); + } + } return builder.toString(); } @@ -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()); } diff --git a/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/JdbcIssuesTest.java b/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/JdbcIssuesTest.java index 9dc97a494..7a7873e45 100644 --- a/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/JdbcIssuesTest.java +++ b/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/JdbcIssuesTest.java @@ -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); @@ -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); @@ -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); @@ -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); diff --git a/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/comparison/DateTimeComparisonTest.java b/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/comparison/DateTimeComparisonTest.java index c941d39b6..dd2173716 100644 --- a/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/comparison/DateTimeComparisonTest.java +++ b/clickhouse-jdbc/src/test/java/com/clickhouse/jdbc/comparison/DateTimeComparisonTest.java @@ -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 { From 21f895cc23c06c3ac9b485725a48f351575b1eba Mon Sep 17 00:00:00 2001 From: Sergey Chernov Date: Tue, 29 Sep 2026 19:10:30 -0700 Subject: [PATCH 9/9] Fixed tests for r2dbc --- .../com/clickhouse/r2dbc/BaseR2dbcTest.java | 79 ++++++++++++++++++- .../r2dbc/spi/test/R2DBCTestKitImplTest.java | 6 +- 2 files changed, 79 insertions(+), 6 deletions(-) diff --git a/clickhouse-r2dbc/src/test/java/com/clickhouse/r2dbc/BaseR2dbcTest.java b/clickhouse-r2dbc/src/test/java/com/clickhouse/r2dbc/BaseR2dbcTest.java index b0c32e51c..d03589848 100644 --- a/clickhouse-r2dbc/src/test/java/com/clickhouse/r2dbc/BaseR2dbcTest.java +++ b/clickhouse-r2dbc/src/test/java/com/clickhouse/r2dbc/BaseR2dbcTest.java @@ -1,11 +1,15 @@ package com.clickhouse.r2dbc; +import java.util.Map; + import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import com.clickhouse.client.BaseIntegrationTest; +import com.clickhouse.client.ClickHouseNode; import com.clickhouse.client.ClickHouseProtocol; import com.clickhouse.client.ClickHouseServerForTest; +import com.clickhouse.client.config.ClickHouseClientOption; import io.r2dbc.spi.ConnectionFactories; import io.r2dbc.spi.ConnectionFactory; @@ -27,16 +31,85 @@ public static void afterSuite() throws Exception { ClickHouseServerForTest.afterSuite(); } + 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 options) { + return addCustomSettings(super.getServer(protocol, options)); + } + protected ConnectionFactory getConnectionFactory(ClickHouseProtocol protocol, String... parameters) { StringBuilder builder = new StringBuilder(getServer(protocol).toUri("r2dbc:ch:").toString()); for (String queryString : parameters) { if (queryString != null && !queryString.isEmpty()) { - if (queryString.charAt(0) != '&') { - builder.append('&'); + char sep = builder.indexOf("?") >= 0 ? '&' : '?'; + if (queryString.charAt(0) == '&' || queryString.charAt(0) == '?') { + queryString = queryString.substring(1); } - builder.append(queryString); + builder.append(sep).append(queryString); } } + if (builder.indexOf("custom_settings") == -1) { + char sep = builder.indexOf("?") >= 0 ? '&' : '?'; + builder.append(sep).append("custom_settings=network_compression_method=lz4"); + } + if (builder.indexOf("custom_http_params") == -1) { + char sep = builder.indexOf("?") >= 0 ? '&' : '?'; + builder.append(sep).append("custom_http_params=network_compression_method=lz4"); + } ConnectionFactory connectionFactory = ConnectionFactories.get(builder.toString()); return connectionFactory; } diff --git a/clickhouse-r2dbc/src/test/java/com/clickhouse/r2dbc/spi/test/R2DBCTestKitImplTest.java b/clickhouse-r2dbc/src/test/java/com/clickhouse/r2dbc/spi/test/R2DBCTestKitImplTest.java index d082d6a37..e547bce09 100644 --- a/clickhouse-r2dbc/src/test/java/com/clickhouse/r2dbc/spi/test/R2DBCTestKitImplTest.java +++ b/clickhouse-r2dbc/src/test/java/com/clickhouse/r2dbc/spi/test/R2DBCTestKitImplTest.java @@ -54,7 +54,7 @@ public static void setup() throws Exception { ClickHouseServerForTest.beforeSuite(); connectionFactory = ConnectionFactories.get( - format("r2dbc:clickhouse:%s://%s:%s@%s/%s?falan=filan&custom_http_params=async_insert=0&%s#tag1", DEFAULT_PROTOCOL, USER, PASSWORD, + format("r2dbc:clickhouse:%s://%s:%s@%s/%s?falan=filan&custom_http_params=async_insert=0,network_compression_method=lz4&custom_settings=network_compression_method=lz4&%s#tag1", DEFAULT_PROTOCOL, USER, PASSWORD, getClickHouseAddress(DEFAULT_PROTOCOL, false), DATABASE, EXTRA_PARAM)); jdbcTemplate = jdbcTemplate(null); } @@ -85,10 +85,10 @@ private static JdbcTemplate jdbcTemplate(String database) throws SQLException { Driver driver = new ClickHouseDriver(); DriverManager.registerDriver(driver); if (database == null) { - source.setJdbcUrl(format("jdbc:clickhouse:%s://%s?custom_http_params=async_insert=0%s", DEFAULT_PROTOCOL, + source.setJdbcUrl(format("jdbc:clickhouse:%s://%s?custom_http_params=async_insert=0,network_compression_method=lz4&custom_settings=network_compression_method=lz4&%s", DEFAULT_PROTOCOL, getClickHouseAddress(DEFAULT_PROTOCOL, false), EXTRA_PARAM)); } else { - source.setJdbcUrl(format("jdbc:clickhouse:%s://%s/%s?custom_http_params=async_insert=0&%s", DEFAULT_PROTOCOL, + source.setJdbcUrl(format("jdbc:clickhouse:%s://%s/%s?custom_http_params=async_insert=0,network_compression_method=lz4&custom_settings=network_compression_method=lz4&%s", DEFAULT_PROTOCOL, getClickHouseAddress(DEFAULT_PROTOCOL, false), DATABASE, EXTRA_PARAM)); }