From 02523a9f930bf43faa621ce900c5ef185a6045b2 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 15 Jul 2026 11:48:15 +0200 Subject: [PATCH 1/5] perf: use Cursor in ZSTDCodec to avoid Vec alloc and copy --- Cargo.lock | 1 + parquet/Cargo.toml | 3 ++- parquet/src/compression.rs | 28 +++++++++++++++++++--------- 3 files changed, 22 insertions(+), 10 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index d781954e8564..6e7262fcca47 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2370,6 +2370,7 @@ dependencies = [ "tokio", "twox-hash", "zstd", + "zstd-safe", ] [[package]] diff --git a/parquet/Cargo.toml b/parquet/Cargo.toml index f3342f03fc23..d36f86bc5cdc 100644 --- a/parquet/Cargo.toml +++ b/parquet/Cargo.toml @@ -58,6 +58,7 @@ brotli = { version = "8.0", default-features = false, features = ["std"], option flate2 = { version = "1.1", default-features = false, optional = true } lz4_flex = { version = "0.13", default-features = false, features = ["std", "frame"], optional = true } zstd = { version = "0.13", optional = true, default-features = false } +zstd-safe = { version = "7.2", optional = true, default-features = false } chrono = { workspace = true } num-bigint = { version = "0.5", default-features = false } num-integer = { version = "0.1.46", default-features = false, features = ["std"] } @@ -118,7 +119,7 @@ async = ["futures", "tokio"] # Enable object_store integration object_store = ["dep:object_store", "async"] # Group Zstd dependencies -zstd = ["dep:zstd"] +zstd = ["dep:zstd", "dep:zstd-safe"] # Verify 32-bit CRC checksum when decoding parquet pages crc = ["dep:crc32fast"] # Enable SIMD UTF-8 validation diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index fe2fb59c5b8c..51d0e9668da0 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -505,6 +505,7 @@ pub use lz4_codec::*; mod zstd_codec { use crate::compression::{Codec, ZstdLevel}; use crate::errors::Result; + use std::io::Cursor; /// Codec for Zstandard compression algorithm. /// @@ -534,22 +535,31 @@ mod zstd_codec { output_buf: &mut Vec, uncompress_size: Option, ) -> Result { + let offset = output_buf.len(); let capacity = uncompress_size.unwrap_or_else(|| { // Get the decompressed size from the zstd frame header - zstd::zstd_safe::get_frame_content_size(input_buf) - .ok() - .flatten() - .unwrap_or(input_buf.len() as u64 * 4) as usize + // See doc of upper_bound about "experimental" feature. + zstd::bulk::Decompressor::upper_bound(input_buf) + .unwrap_or(input_buf.len().saturating_mul(4)) }); - let decompressed = self.decompressor.decompress(input_buf, capacity)?; - let len = decompressed.len(); - output_buf.extend_from_slice(&decompressed); + output_buf.reserve(capacity); + + let mut cursor = Cursor::new(output_buf); + cursor.set_position(offset as u64); + let len = self + .decompressor + .decompress_to_buffer(input_buf, &mut cursor)?; Ok(len) } fn compress(&mut self, input_buf: &[u8], output_buf: &mut Vec) -> Result<()> { - let compressed = self.compressor.compress(input_buf)?; - output_buf.extend_from_slice(&compressed); + let offset = output_buf.len(); + let buffer_len = zstd_safe::compress_bound(input_buf.len()); + output_buf.reserve(buffer_len); + + let mut cursor = Cursor::new(output_buf); + cursor.set_position(offset as u64); + let _written = self.compressor.compress_to_buffer(input_buf, &mut cursor)?; Ok(()) } } From 0ecd05b704bc56556647f9e9d567fc6e0e498ece Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 15 Jul 2026 12:11:12 +0200 Subject: [PATCH 2/5] nit: rename vars --- parquet/src/compression.rs | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index 51d0e9668da0..2ab9e9249c8e 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -536,13 +536,13 @@ mod zstd_codec { uncompress_size: Option, ) -> Result { let offset = output_buf.len(); - let capacity = uncompress_size.unwrap_or_else(|| { - // Get the decompressed size from the zstd frame header + let len = uncompress_size.unwrap_or_else(|| { + // Get the decompressed size from the zstd frame header. // See doc of upper_bound about "experimental" feature. zstd::bulk::Decompressor::upper_bound(input_buf) .unwrap_or(input_buf.len().saturating_mul(4)) }); - output_buf.reserve(capacity); + output_buf.reserve(len); let mut cursor = Cursor::new(output_buf); cursor.set_position(offset as u64); @@ -554,8 +554,8 @@ mod zstd_codec { fn compress(&mut self, input_buf: &[u8], output_buf: &mut Vec) -> Result<()> { let offset = output_buf.len(); - let buffer_len = zstd_safe::compress_bound(input_buf.len()); - output_buf.reserve(buffer_len); + let len = zstd_safe::compress_bound(input_buf.len()); + output_buf.reserve(len); let mut cursor = Cursor::new(output_buf); cursor.set_position(offset as u64); From d0f851338f16f4f51ab67a1d6472a3131c6da10b Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Fri, 17 Jul 2026 12:41:26 +0200 Subject: [PATCH 3/5] use zstd::zstd_safe re-export --- Cargo.lock | 1 - parquet/Cargo.toml | 3 +-- parquet/src/compression.rs | 2 +- 3 files changed, 2 insertions(+), 4 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 6e7262fcca47..d781954e8564 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2370,7 +2370,6 @@ dependencies = [ "tokio", "twox-hash", "zstd", - "zstd-safe", ] [[package]] diff --git a/parquet/Cargo.toml b/parquet/Cargo.toml index d36f86bc5cdc..f3342f03fc23 100644 --- a/parquet/Cargo.toml +++ b/parquet/Cargo.toml @@ -58,7 +58,6 @@ brotli = { version = "8.0", default-features = false, features = ["std"], option flate2 = { version = "1.1", default-features = false, optional = true } lz4_flex = { version = "0.13", default-features = false, features = ["std", "frame"], optional = true } zstd = { version = "0.13", optional = true, default-features = false } -zstd-safe = { version = "7.2", optional = true, default-features = false } chrono = { workspace = true } num-bigint = { version = "0.5", default-features = false } num-integer = { version = "0.1.46", default-features = false, features = ["std"] } @@ -119,7 +118,7 @@ async = ["futures", "tokio"] # Enable object_store integration object_store = ["dep:object_store", "async"] # Group Zstd dependencies -zstd = ["dep:zstd", "dep:zstd-safe"] +zstd = ["dep:zstd"] # Verify 32-bit CRC checksum when decoding parquet pages crc = ["dep:crc32fast"] # Enable SIMD UTF-8 validation diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index 2ab9e9249c8e..2170e6db15c8 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -554,7 +554,7 @@ mod zstd_codec { fn compress(&mut self, input_buf: &[u8], output_buf: &mut Vec) -> Result<()> { let offset = output_buf.len(); - let len = zstd_safe::compress_bound(input_buf.len()); + let len = zstd::zstd_safe::compress_bound(input_buf.len()); output_buf.reserve(len); let mut cursor = Cursor::new(output_buf); From 02bd80c8a6b849b1c6166f512dc8f70e45d7a5d9 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Mon, 20 Jul 2026 14:52:18 +0200 Subject: [PATCH 4/5] clarify upper_bound comment --- parquet/src/compression.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index 2170e6db15c8..6c4ca79c4d22 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -537,7 +537,7 @@ mod zstd_codec { ) -> Result { let offset = output_buf.len(); let len = uncompress_size.unwrap_or_else(|| { - // Get the decompressed size from the zstd frame header. + // Get the decompressed size of all zstd frames in the input. // See doc of upper_bound about "experimental" feature. zstd::bulk::Decompressor::upper_bound(input_buf) .unwrap_or(input_buf.len().saturating_mul(4)) From d0d52ddc55660f31ac4812fef545f57caebd9344 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Mon, 20 Jul 2026 15:35:35 +0200 Subject: [PATCH 5/5] restore get_frame_content_size and remove upper_bound --- parquet/src/compression.rs | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index 6c4ca79c4d22..6f3a18fe60b0 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -536,12 +536,15 @@ mod zstd_codec { uncompress_size: Option, ) -> Result { let offset = output_buf.len(); - let len = uncompress_size.unwrap_or_else(|| { - // Get the decompressed size of all zstd frames in the input. - // See doc of upper_bound about "experimental" feature. - zstd::bulk::Decompressor::upper_bound(input_buf) - .unwrap_or(input_buf.len().saturating_mul(4)) - }); + let len = uncompress_size + .or_else(|| { + // Get the decompressed size from the zstd frame header + zstd::zstd_safe::get_frame_content_size(input_buf) + .ok() + .flatten() + .map(|size| size as usize) + }) + .unwrap_or(input_buf.len().saturating_mul(4)); output_buf.reserve(len); let mut cursor = Cursor::new(output_buf);