This is an automated email from the ASF dual-hosted git repository.
Fokko pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/parquet-java.git
The following commit(s) were added to refs/heads/master by this push:
new 8e903700c MINOR: Reuse RLE decoder buffers for packed runs (#3749)
8e903700c is described below
commit 8e903700c18f55c9ff744e556a2b58da231afae6
Author: Divjot Arora <[email protected]>
AuthorDate: Tue Sep 8 09:51:27 2026 -0400
MINOR: Reuse RLE decoder buffers for packed runs (#3749)
* Reuse RLE packed buffers to avoid per-run allocation
* Fix Spotless formatting in RLE decoder tests
* Fix RLE decoder benchmark formatting
---
...RunLengthBitPackingHybridDecodingBenchmark.java | 84 ++++++++++++++++++++++
.../rle/RunLengthBitPackingHybridDecoder.java | 26 +++++--
.../rle/TestRunLengthBitPackingHybridEncoder.java | 54 ++++++++++++++
3 files changed, 157 insertions(+), 7 deletions(-)
diff --git
a/parquet-benchmarks/src/main/java/org/apache/parquet/benchmarks/RunLengthBitPackingHybridDecodingBenchmark.java
b/parquet-benchmarks/src/main/java/org/apache/parquet/benchmarks/RunLengthBitPackingHybridDecodingBenchmark.java
new file mode 100644
index 000000000..1948d1425
--- /dev/null
+++
b/parquet-benchmarks/src/main/java/org/apache/parquet/benchmarks/RunLengthBitPackingHybridDecodingBenchmark.java
@@ -0,0 +1,84 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.parquet.benchmarks;
+
+import java.io.ByteArrayInputStream;
+import java.io.IOException;
+import java.util.concurrent.TimeUnit;
+import org.apache.parquet.bytes.DirectByteBufferAllocator;
+import org.apache.parquet.column.values.rle.RunLengthBitPackingHybridDecoder;
+import org.apache.parquet.column.values.rle.RunLengthBitPackingHybridEncoder;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Level;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OperationsPerInvocation;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.Warmup;
+import org.openjdk.jmh.infra.Blackhole;
+
+@BenchmarkMode(Mode.Throughput)
+@OutputTimeUnit(TimeUnit.SECONDS)
+@Fork(1)
+@Warmup(iterations = 3, time = 1)
+@Measurement(iterations = 5, time = 1)
+@State(Scope.Thread)
+public class RunLengthBitPackingHybridDecodingBenchmark {
+
+ private static final int VALUE_COUNT = 100_000;
+
+ @Param({"1", "8", "16"})
+ public int bitWidth;
+
+ private byte[] encoded;
+
+ @Setup(Level.Trial)
+ public void setup() throws IOException {
+ RunLengthBitPackingHybridEncoder encoder = new
RunLengthBitPackingHybridEncoder(
+ bitWidth, 64 * 1024, 4 * 1024 * 1024, new DirectByteBufferAllocator());
+ int mask = (1 << bitWidth) - 1;
+ for (int i = 0; i < VALUE_COUNT; i += 32) {
+ for (int j = 0; j < 16 && i + j < VALUE_COUNT; j++) {
+ encoder.writeInt((i + j) & mask);
+ }
+ int repeated = (i / 32) & mask;
+ for (int j = 16; j < 32 && i + j < VALUE_COUNT; j++) {
+ encoder.writeInt(repeated);
+ }
+ }
+ encoded = encoder.toBytes().toByteArray();
+ encoder.close();
+ }
+
+ @Benchmark
+ @OperationsPerInvocation(VALUE_COUNT)
+ public void decode(Blackhole bh) throws IOException {
+ RunLengthBitPackingHybridDecoder decoder =
+ new RunLengthBitPackingHybridDecoder(bitWidth, new
ByteArrayInputStream(encoded));
+ for (int i = 0; i < VALUE_COUNT; i++) {
+ bh.consume(decoder.readInt());
+ }
+ }
+}
diff --git
a/parquet-column/src/main/java/org/apache/parquet/column/values/rle/RunLengthBitPackingHybridDecoder.java
b/parquet-column/src/main/java/org/apache/parquet/column/values/rle/RunLengthBitPackingHybridDecoder.java
index e55b276b2..0a798ef67 100644
---
a/parquet-column/src/main/java/org/apache/parquet/column/values/rle/RunLengthBitPackingHybridDecoder.java
+++
b/parquet-column/src/main/java/org/apache/parquet/column/values/rle/RunLengthBitPackingHybridDecoder.java
@@ -18,9 +18,10 @@
*/
package org.apache.parquet.column.values.rle;
-import java.io.DataInputStream;
+import java.io.EOFException;
import java.io.IOException;
import java.io.InputStream;
+import java.util.Arrays;
import org.apache.parquet.Preconditions;
import org.apache.parquet.bytes.BytesUtils;
import org.apache.parquet.column.values.bitpacking.BytePacker;
@@ -47,7 +48,9 @@ public class RunLengthBitPackingHybridDecoder {
private MODE mode;
private int currentCount;
private int currentValue;
- private int[] currentBuffer;
+ private int currentBufferPosition;
+ private int[] currentBuffer = new int[0];
+ private byte[] packedBytes = new byte[0];
public RunLengthBitPackingHybridDecoder(int bitWidth, InputStream in) {
LOG.debug("decoding bitWidth {}", bitWidth);
@@ -69,7 +72,7 @@ public class RunLengthBitPackingHybridDecoder {
result = currentValue;
break;
case PACKED:
- result = currentBuffer[currentBuffer.length - 1 - currentCount];
+ result = currentBuffer[currentBufferPosition++];
break;
default:
throw new ParquetDecodingException("not a valid mode " + mode);
@@ -91,16 +94,25 @@ public class RunLengthBitPackingHybridDecoder {
int numGroups = header >>> 1;
currentCount = numGroups * 8;
LOG.debug("reading {} values BIT PACKED", currentCount);
- currentBuffer = new int[currentCount]; // TODO: reuse a buffer
- byte[] bytes = new byte[numGroups * bitWidth];
+ if (currentBuffer.length < currentCount) {
+ currentBuffer = new int[currentCount];
+ }
+ currentBufferPosition = 0;
+ int bytesRequired = numGroups * bitWidth;
+ if (packedBytes.length < bytesRequired) {
+ packedBytes = new byte[bytesRequired];
+ }
// At the end of the file RLE data though, there might not be that
many bytes left.
int bytesToRead = (int) Math.ceil(currentCount * bitWidth / 8.0);
bytesToRead = Math.min(bytesToRead, in.available());
- new DataInputStream(in).readFully(bytes, 0, bytesToRead);
+ if (in.readNBytes(packedBytes, 0, bytesToRead) != bytesToRead) {
+ throw new EOFException();
+ }
+ Arrays.fill(packedBytes, bytesToRead, bytesRequired, (byte) 0);
for (int valueIndex = 0, byteIndex = 0;
valueIndex < currentCount;
valueIndex += 8, byteIndex += bitWidth) {
- packer.unpack8Values(bytes, byteIndex, currentBuffer, valueIndex);
+ packer.unpack8Values(packedBytes, byteIndex, currentBuffer,
valueIndex);
}
break;
default:
diff --git
a/parquet-column/src/test/java/org/apache/parquet/column/values/rle/TestRunLengthBitPackingHybridEncoder.java
b/parquet-column/src/test/java/org/apache/parquet/column/values/rle/TestRunLengthBitPackingHybridEncoder.java
index b4f484755..d73061316 100644
---
a/parquet-column/src/test/java/org/apache/parquet/column/values/rle/TestRunLengthBitPackingHybridEncoder.java
+++
b/parquet-column/src/test/java/org/apache/parquet/column/values/rle/TestRunLengthBitPackingHybridEncoder.java
@@ -298,6 +298,60 @@ public class TestRunLengthBitPackingHybridEncoder {
assertThat(stream.available()).isEqualTo(0);
}
+ @Test
+ public void testSmallerPackedRunAfterLargerPackedRun() throws Exception {
+ byte[] bytes = {
+ (byte) 5,
+ (byte) 0xe4,
+ (byte) 0xe4,
+ (byte) 0xe4,
+ (byte) 0xe4,
+ (byte) 16,
+ (byte) 2,
+ (byte) 3,
+ (byte) 0x39,
+ (byte) 0x39
+ };
+ RunLengthBitPackingHybridDecoder decoder =
+ new RunLengthBitPackingHybridDecoder(2, new
ByteArrayInputStream(bytes));
+
+ for (int i = 0; i < 16; i++) {
+ assertThat(decoder.readInt()).isEqualTo(i % 4);
+ }
+ for (int i = 0; i < 8; i++) {
+ assertThat(decoder.readInt()).isEqualTo(2);
+ }
+ assertThat(decoder.readInt()).isEqualTo(1);
+ assertThat(decoder.readInt()).isEqualTo(2);
+ assertThat(decoder.readInt()).isEqualTo(3);
+ assertThat(decoder.readInt()).isEqualTo(0);
+ assertThat(decoder.readInt()).isEqualTo(1);
+ assertThat(decoder.readInt()).isEqualTo(2);
+ assertThat(decoder.readInt()).isEqualTo(3);
+ assertThat(decoder.readInt()).isEqualTo(0);
+ }
+
+ @Test
+ public void testTruncatedPackedRunClearsReusedBytes() throws Exception {
+ byte[] bytes = {
+ (byte) 5, (byte) 0xe4, (byte) 0xe4, (byte) 0xe4, (byte) 0xe4, (byte) 16,
(byte) 2, (byte) 3, (byte) 0x39
+ };
+ RunLengthBitPackingHybridDecoder decoder =
+ new RunLengthBitPackingHybridDecoder(2, new
ByteArrayInputStream(bytes));
+
+ for (int i = 0; i < 24; i++) {
+ decoder.readInt();
+ }
+ assertThat(decoder.readInt()).isEqualTo(1);
+ assertThat(decoder.readInt()).isEqualTo(2);
+ assertThat(decoder.readInt()).isEqualTo(3);
+ assertThat(decoder.readInt()).isEqualTo(0);
+ assertThat(decoder.readInt()).isEqualTo(0);
+ assertThat(decoder.readInt()).isEqualTo(0);
+ assertThat(decoder.readInt()).isEqualTo(0);
+ assertThat(decoder.readInt()).isEqualTo(0);
+ }
+
private static List<Integer> unpack(int bitWidth, int numValues,
ByteArrayInputStream is) throws Exception {
BytePacker packer = Packer.LITTLE_ENDIAN.newBytePacker(bitWidth);