This is an automated email from the ASF dual-hosted git repository.
hubgeter pushed a commit to branch orc
in repository https://gitbox.apache.org/repos/asf/doris-thirdparty.git
The following commit(s) were added to refs/heads/orc by this push:
new 5860f94768b [improvement] Allocate the block decompression input
buffer lazily (#418)
5860f94768b is described below
commit 5860f94768b94406f31698ac66f4cc3d2a781324
Author: daidai <[email protected]>
AuthorDate: Tue Sep 29 16:21:52 2026 +0800
[improvement] Allocate the block decompression input buffer lazily (#418)
BlockDecompressionStream allocated and zeroed a block-sized inputDataBuffer
for every stream, although NextDecompress only uses it to stitch a chunk
that
spans input buffers. Every ZLIB, Snappy, LZ4 and ZSTD stream therefore held
a
second block next to its output buffer; with 256 KiB blocks a wide scan over
100 nullable string columns (PRESENT, LENGTH and DATA streams) pays about
75 MiB per reader for buffers that are usually unused.
Start the buffer empty and let NextDecompress grow it to the chunk length on
the first split chunk. Add a test that counts pool allocations for
contiguous
and split chunks of each block codec.
---
c++/src/Compression.cc | 5 ++--
c++/test/TestCompression.cc | 70 +++++++++++++++++++++++++++++++++++++++++++++
2 files changed, 73 insertions(+), 2 deletions(-)
diff --git a/c++/src/Compression.cc b/c++/src/Compression.cc
index 4b7f7f8f2e5..b6d3168ad35 100644
--- a/c++/src/Compression.cc
+++ b/c++/src/Compression.cc
@@ -744,7 +744,8 @@ namespace orc {
private:
// may need to stitch together multiple input buffers;
- // to give snappy a contiguous block
+ // to give snappy a contiguous block. Allocated on the first chunk that
spans input
+ // buffers, so streams whose chunks are always contiguous do not pay a
block per stream.
DataBuffer<char> inputDataBuffer;
};
@@ -752,7 +753,7 @@ namespace orc {
size_t blockSize,
MemoryPool& _pool,
ReaderMetrics* _metrics)
: DecompressionStream(std::move(inStream), blockSize, _pool, _metrics),
- inputDataBuffer(pool, blockSize) {}
+ inputDataBuffer(pool, 0) {}
void BlockDecompressionStream::NextDecompress(const void** data, int* size,
size_t availableSize) {
diff --git a/c++/test/TestCompression.cc b/c++/test/TestCompression.cc
index d8ddc24042a..47200225289 100644
--- a/c++/test/TestCompression.cc
+++ b/c++/test/TestCompression.cc
@@ -25,6 +25,8 @@
#include "wrap/orc-proto-wrapper.hh"
#include <algorithm>
+#include <cstdlib>
+#include <map>
namespace orc {
const int DEFAULT_MEM_STREAM_SIZE = 1024 * 1024 * 2; // 2M
@@ -422,4 +424,72 @@ namespace orc {
corrupted[3] = 0x07;
EXPECT_THROW(decompressAll(corrupted, 1024), ParseError);
}
+
+ class CountingMemoryPool : public MemoryPool {
+ public:
+ char* malloc(uint64_t size) override {
+ char* p = static_cast<char*>(std::malloc(size));
+ sizes[p] = size;
+ current += size;
+ peak = std::max(peak, current);
+ return p;
+ }
+ void free(char* p) override {
+ auto it = sizes.find(p);
+ current -= it->second;
+ sizes.erase(it);
+ std::free(p);
+ }
+
+ uint64_t current = 0;
+ uint64_t peak = 0;
+
+ private:
+ std::map<char*, uint64_t> sizes;
+ };
+
+ // Block codecs copy a chunk into a scratch buffer only when it spans input
buffers, so a
+ // stream whose chunks are contiguous must not allocate a second block.
+ TEST(Compression, block_decompression_input_buffer_is_lazy) {
+ const uint64_t blockSize = 256 * 1024;
+ std::string testData;
+ for (int i = 0; i < 1000; ++i) {
+ testData.push_back(static_cast<char>('a' + i % 26));
+ }
+ for (CompressionKind kind : {CompressionKind_ZLIB, CompressionKind_ZSTD,
CompressionKind_LZ4,
+ CompressionKind_SNAPPY}) {
+ SCOPED_TRACE(kind);
+ MemoryOutputStream memStream(DEFAULT_MEM_STREAM_SIZE);
+ compressAndVerify(kind, &memStream, CompressionStrategy_COMPRESSION,
1024, 1024,
+ *getDefaultPool(), testData.data(), testData.size());
+ const std::string compressed(memStream.getData(), memStream.getLength());
+ ASSERT_EQ(0, compressed[0] & 1);
+
+ // 0 hands over the whole stream in one buffer; 7 splits the chunk
across many buffers.
+ for (uint64_t inputBlockSize : {uint64_t{0}, uint64_t{7}}) {
+ SCOPED_TRACE(inputBlockSize);
+ CountingMemoryPool pool;
+ {
+ auto decompressStream =
+ createDecompressor(kind,
+ std::make_unique<SeekableArrayInputStream>(
+ compressed.data(), compressed.size(),
inputBlockSize),
+ blockSize, pool, getDefaultReaderMetrics());
+ std::string result;
+ const void* data;
+ int size;
+ while (decompressStream->Next(&data, &size)) {
+ result.append(static_cast<const char*>(data),
static_cast<size_t>(size));
+ }
+ EXPECT_EQ(testData, result);
+ if (inputBlockSize == 0) {
+ EXPECT_EQ(blockSize, pool.peak);
+ } else {
+ EXPECT_EQ(blockSize + compressed.size() - 3, pool.peak);
+ }
+ }
+ EXPECT_EQ(0, pool.current);
+ }
+ }
+ }
} // namespace orc
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]