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]

Reply via email to