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]