This is an automated email from the ASF dual-hosted git repository.

Gabriel39 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 78ba6c80934 [fix](parquet) Fix concatenated GZIP decoding and 
large-map regression memory (#68819)
78ba6c80934 is described below

commit 78ba6c80934ec870f0c8a5a84087b546c209edaf
Author: Gabriel <[email protected]>
AuthorDate: Fri Oct 9 21:10:48 2026 +0800

    [fix](parquet) Fix concatenated GZIP decoding and large-map regression 
memory (#68819)
    
    ### What problem does this PR solve?
    
    Issue Number: N/A
    
    Related PRs: #68808, #68593
    
    Two failures in the Parquet TVF regression suite prevent the full
    fixture set
    from completing:
    
    - The large-string-map fixture has two 1 GiB keys in a column chunk
    larger
    than 2 GiB. Keep the change from #68808: use a statement-local
    `batch_size=1`
    hint and disable aggregate pushdown. Both rows still undergo full
    Map/string
    decoding, including the dictionary-to-PLAIN transition, while each
    output
    batch contains one key. Oversized column-chunk coverage remains;
    accumulating
      both keys in one output ColumnString is no longer exercised.
    - The concatenated GZIP fixture has a 1416-byte compressed page
    containing two
    members that expand to 4096 and 8 bytes. `libdeflate_gzip_decompress`
    decodes
    only the first member, returning SHORT_OUTPUT against the 4104-byte page
    size.
    Use `libdeflate_gzip_decompress_ex` to consume every member, then
    validate the
    exact combined output size. Corrupt or truncated later members, trailing
      garbage, and size mismatches still fail.
    
    The HDFS query and expected results for the concatenated GZIP fixture
    remain
    unchanged. Unit coverage includes multiple members, empty members at
    every
    position, all-empty members, corrupt/truncated trailers, trailing
    garbage,
    and undersized/oversized output buffers.
    
    ### Release note
    
    Fix reading Parquet GZIP pages containing concatenated members.
    
    ### Check List (For Author)
    
    - [x] Local ASAN BE build.
    - [x] Complete `test_hdfs_parquet_group0` suite against HDFS: 55 query
    checks,
    1 successful suite, 0 failed suites, 0 fatal scripts, 0 skipped scripts.
    This includes both the large-map query and the concatenated GZIP query.
    - [x] `run-be-ut.sh --run --filter=BlockCompressionTest.*`: 7 tests
    passed.
    - [x] clang-format 16, clang-tidy on changed lines, BE build hygiene,
    and
      `git diff --check` passed.
    - Behavior changed: Parquet GZIP decoding supports all concatenated
    members;
    the large-map regression uses bounded output batches. Expected results
    are
      unchanged.
    - Does this need documentation: No.
    
    ### Check List (For Reviewer who merge this PR)
    
    - [ ] Confirm the release note
    - [ ] Confirm test cases
    - [ ] Add branch pick label
---
 be/src/util/block_compression.cpp                  | 30 ++++++++++---
 be/test/util/block_compression_test.cpp            | 51 ++++++++++++++++++++++
 .../tvf/test_hdfs_parquet_group0.groovy            |  5 ++-
 3 files changed, 78 insertions(+), 8 deletions(-)

diff --git a/be/src/util/block_compression.cpp 
b/be/src/util/block_compression.cpp
index 40ca6d17ad0..96e77b270df 100644
--- a/be/src/util/block_compression.cpp
+++ b/be/src/util/block_compression.cpp
@@ -1447,8 +1447,6 @@ public:
     ~GzipBlockCompressionByLibdeflate() override = default;
 
     Status decompress(const Slice& input, Slice* output) override {
-        // Parquet page headers give the exact uncompressed size, so the page 
must fill the
-        // output buffer. Without actual_out_nbytes_ret, libdeflate rejects 
shorter output.
         if (input.empty() && output->size == 0) {
             return Status::OK();
         }
@@ -1457,12 +1455,30 @@ public:
         if (!decompressor) {
             return Status::InternalError("libdeflate_alloc_decompressor 
error.");
         }
-        auto result = libdeflate_gzip_decompress(decompressor.get(), 
input.data, input.size,
-                                                 output->data, output->size, 
nullptr);
-        if (result != LIBDEFLATE_SUCCESS) {
+        // A Parquet GZIP page may contain concatenated members. libdeflate 
decodes only
+        // one member per call; the page header's exact size applies to their 
combined output.
+        size_t input_offset = 0;
+        size_t output_offset = 0;
+        while (input_offset < input.size) {
+            size_t consumed = 0;
+            size_t produced = 0;
+            auto result = libdeflate_gzip_decompress_ex(
+                    decompressor.get(), input.data + input_offset, input.size 
- input_offset,
+                    output->data + output_offset, output->size - 
output_offset, &consumed,
+                    &produced);
+            if (result != LIBDEFLATE_SUCCESS) {
+                return Status::InternalError(
+                        "libdeflate_gzip_decompress_ex error, res={}, input 
size={}, output "
+                        "size={}",
+                        result, input.size, output->size);
+            }
+            input_offset += consumed;
+            output_offset += produced;
+        }
+        if (output_offset != output->size) {
             return Status::InternalError(
-                    "libdeflate_gzip_decompress error, res={}, input size={}, 
output size={}",
-                    result, input.size, output->size);
+                    "GZIP page decompressed size mismatch, actual={}, 
expected={}", output_offset,
+                    output->size);
         }
         return Status::OK();
     }
diff --git a/be/test/util/block_compression_test.cpp 
b/be/test/util/block_compression_test.cpp
index c7364a59ec9..e333daf9b21 100644
--- a/be/test/util/block_compression_test.cpp
+++ b/be/test/util/block_compression_test.cpp
@@ -211,6 +211,57 @@ TEST_F(BlockCompressionTest, parquet_gzip) {
     EXPECT_FALSE(codec->decompress(Slice("not a gzip stream"), &output).ok());
 }
 
+static void check_concatenated_gzip_decompression(BlockCompressionCodec* codec,
+                                                  const std::string& 
compressed,
+                                                  const std::string& original) 
{
+    std::string restored(original.size(), '\0');
+    Slice output(restored);
+    auto status = codec->decompress(Slice(compressed), &output);
+    EXPECT_TRUE(status.ok()) << status.to_string();
+    EXPECT_EQ(original, restored);
+    EXPECT_EQ(original.size(), output.size);
+
+    if (!original.empty()) {
+        Slice short_output(restored.data(), original.size() - 1);
+        EXPECT_FALSE(codec->decompress(Slice(compressed), &short_output).ok());
+    }
+    std::string larger(original.size() + 1, '\0');
+    Slice long_output(larger);
+    EXPECT_FALSE(codec->decompress(Slice(compressed), &long_output).ok());
+
+    // Validate the final member even when previous members already filled the 
output.
+    EXPECT_FALSE(codec->decompress(Slice(compressed.data(), compressed.size() 
- 1), &output).ok());
+    std::string corrupted = compressed;
+    corrupted[corrupted.size() - 8] ^= 1;
+    EXPECT_FALSE(codec->decompress(Slice(corrupted), &output).ok());
+    std::string trailing = compressed + "trailing junk";
+    EXPECT_FALSE(codec->decompress(Slice(trailing), &output).ok());
+}
+
+TEST_F(BlockCompressionTest, parquet_gzip_concatenated_members) {
+    BlockCompressionCodec* codec = nullptr;
+    ASSERT_TRUE(get_block_compression_codec(tparquet::CompressionCodec::GZIP, 
&codec).ok());
+    BlockCompressionCodec* zlib_codec = nullptr;
+    ASSERT_TRUE(get_block_compression_codec(TFileCompressType::GZ, 
&zlib_codec).ok());
+
+    // A Parquet page's uncompressed size covers all members, including empty 
members.
+    for (const std::vector<std::string>& members :
+         {std::vector<std::string> {"first", std::string(4099, 'x')},
+          std::vector<std::string> {"", "first", "", std::string(4099, 'x'), 
""},
+          std::vector<std::string> {"", ""}}) {
+        std::string compressed;
+        std::string original;
+        for (const auto& member : members) {
+            faststring encoded;
+            ASSERT_TRUE(zlib_codec->compress(Slice(member), &encoded).ok());
+            compressed.append(reinterpret_cast<const char*>(encoded.data()), 
encoded.size());
+            original.append(member);
+        }
+        SCOPED_TRACE(original.size());
+        check_concatenated_gzip_decompression(codec, compressed, original);
+    }
+}
+
 static void check_snappy_decompression(BlockCompressionCodec* codec, const 
faststring& compressed,
                                        const std::string& original) {
     std::string restored(original.size(), '\0');
diff --git 
a/regression-test/suites/external_table_p0/tvf/test_hdfs_parquet_group0.groovy 
b/regression-test/suites/external_table_p0/tvf/test_hdfs_parquet_group0.groovy
index 017c0c09e92..37643276187 100644
--- 
a/regression-test/suites/external_table_p0/tvf/test_hdfs_parquet_group0.groovy
+++ 
b/regression-test/suites/external_table_p0/tvf/test_hdfs_parquet_group0.groovy
@@ -105,7 +105,10 @@ suite("test_hdfs_parquet_group0", "p0,external") {
 
 
             uri = "${defaultFS}" + 
"/user/doris/tvf_data/test_hdfs_parquet/group0/large_string_map.brotli.parquet"
-            order_qt_test_11 """ select count(arr) from HDFS(
+            // Read both 1 GiB keys one row per batch to avoid a 4 GiB output 
buffer allocation.
+            // Disable aggregate pushdown to retain full decoding of the >2 
GiB column chunk.
+            order_qt_test_11 """ select /*+ SET_VAR(batch_size=1, 
enable_push_down_no_group_agg=false) */
+                        count(arr) from HDFS(
                         "uri" = "${uri}",
                         "hadoop.username" = "${hdfsUserName}",
                         "format" = "parquet"); """


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to