diff --git a/clients/src/main/java/org/apache/kafka/common/record/KafkaLZ4BlockInputStream.java b/clients/src/main/java/org/apache/kafka/common/record/KafkaLZ4BlockInputStream.java index 850b1e96e55f2..9a378333e64db 100644 --- a/clients/src/main/java/org/apache/kafka/common/record/KafkaLZ4BlockInputStream.java +++ b/clients/src/main/java/org/apache/kafka/common/record/KafkaLZ4BlockInputStream.java @@ -77,11 +77,6 @@ public KafkaLZ4BlockInputStream(ByteBuffer in, BufferSupplier bufferSupplier, bo this.bufferSupplier = bufferSupplier; readHeader(); decompressionBuffer = bufferSupplier.get(maxBlockSize); - if (!decompressionBuffer.hasArray() || decompressionBuffer.arrayOffset() != 0) { - // require array backed decompression buffer with zero offset - // to simplify workaround for https://github.com/lz4/lz4-java/pull/65 - throw new RuntimeException("decompression buffer must have backing array with zero array offset"); - } finished = false; } @@ -131,10 +126,7 @@ private void readHeader() throws IOException { int len = in.position() - in.reset().position(); - int hash = in.hasArray() ? - // workaround for https://github.com/lz4/lz4-java/pull/65 - CHECKSUM.hash(in.array(), in.arrayOffset() + in.position(), len, 0) : - CHECKSUM.hash(in, in.position(), len, 0); + int hash = CHECKSUM.hash(in, in.position(), len, 0); in.position(in.position() + len); if (in.get() != (byte) ((hash >> 8) & 0xFF)) { throw new IOException(DESCRIPTOR_HASH_MISMATCH); @@ -172,22 +164,8 @@ private void readBlock() throws IOException { if (compressed) { try { - // workaround for https://github.com/lz4/lz4-java/pull/65 - final int bufferSize; - if (in.hasArray()) { - bufferSize = DECOMPRESSOR.decompress( - in.array(), - in.position() + in.arrayOffset(), - blockSize, - decompressionBuffer.array(), - 0, - maxBlockSize - ); - } else { - // decompressionBuffer has zero arrayOffset, so we don't need to worry about - // https://github.com/lz4/lz4-java/pull/65 - bufferSize = DECOMPRESSOR.decompress(in, in.position(), blockSize, decompressionBuffer, 0, maxBlockSize); - } + final int bufferSize = DECOMPRESSOR.decompress(in, in.position(), blockSize, decompressionBuffer, 0, + maxBlockSize); decompressionBuffer.position(0); decompressionBuffer.limit(bufferSize); decompressedBuffer = decompressionBuffer; @@ -201,10 +179,7 @@ private void readBlock() throws IOException { // verify checksum if (flg.isBlockChecksumSet()) { - // workaround for https://github.com/lz4/lz4-java/pull/65 - int hash = in.hasArray() ? - CHECKSUM.hash(in.array(), in.arrayOffset() + in.position(), blockSize, 0) : - CHECKSUM.hash(in, in.position(), blockSize, 0); + int hash = CHECKSUM.hash(in, in.position(), blockSize, 0); in.position(in.position() + blockSize); if (hash != in.getInt()) { throw new IOException(BLOCK_HASH_MISMATCH);