From cb73239d64b5444130a248b13c068d2412ae2557 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 29 Jul 2026 16:16:35 +0200 Subject: [PATCH 1/8] avoid over-reserving Vecs in ZSTDCodec --- parquet/src/compression.rs | 96 +++++++++++++++++++++++++------------- 1 file changed, 63 insertions(+), 33 deletions(-) diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index 6f3a18fe60b0..b90a45086add 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -503,29 +503,63 @@ pub use lz4_codec::*; #[cfg(any(feature = "zstd", test))] mod zstd_codec { + use zstd::zstd_safe; + use crate::compression::{Codec, ZstdLevel}; use crate::errors::Result; - use std::io::Cursor; + use std::io::Read; /// Codec for Zstandard compression algorithm. /// /// Uses `zstd::bulk` API with reusable compressor/decompressor contexts /// to avoid the overhead of reinitializing contexts for each operation. pub struct ZSTDCodec { - compressor: zstd::bulk::Compressor<'static>, - decompressor: zstd::bulk::Decompressor<'static>, + cctx: zstd_safe::CCtx<'static>, + dctx: zstd_safe::DCtx<'static>, } impl ZSTDCodec { /// Creates new Zstandard compression codec. pub(crate) fn new(level: ZstdLevel) -> Self { - Self { - compressor: zstd::bulk::Compressor::new(level.compression_level()) - .expect("valid zstd compression level"), - decompressor: zstd::bulk::Decompressor::new() - .expect("can create zstd decompressor"), + let mut cctx = zstd_safe::CCtx::create(); + cctx.set_parameter(zstd_safe::CParameter::CompressionLevel( + level.compression_level(), + )) + .expect("valid zstd compression level"); + + let dctx = zstd_safe::DCtx::create(); + + Self { cctx, dctx } + } + } + + /// Avoids zstd crate abstractions to minimize redundant copies; + /// [zstd::stream::Encoder] uses a buffered writer which is pointless for Vec. + fn compress_to_vec( + cctx: &mut zstd_safe::CCtx<'static>, + input_buf: &[u8], + output_buf: &mut Vec, + ) -> std::result::Result<(), zstd_safe::ErrorCode> { + cctx.reset(zstd_safe::ResetDirective::SessionOnly)?; + cctx.set_pledged_src_size(Some(input_buf.len() as u64))?; + + let mut input = zstd_safe::InBuffer::around(input_buf); + while input.pos < input.src.len() { + let mut output = zstd_safe::OutBuffer::around_pos(output_buf, output_buf.len()); + let end_op = zstd_safe::zstd_sys::ZSTD_EndDirective::ZSTD_e_continue; + let to_flush = cctx.compress_stream2(&mut output, &mut input, end_op)?; + output_buf.reserve(to_flush); + } + + loop { + let mut output = zstd_safe::OutBuffer::around_pos(output_buf, output_buf.len()); + let to_flush = cctx.end_stream(&mut output)?; + if to_flush == 0 { + break; } + output_buf.reserve_exact(to_flush); } + Ok(()) } impl Codec for ZSTDCodec { @@ -535,35 +569,31 @@ mod zstd_codec { output_buf: &mut Vec, uncompress_size: Option, ) -> Result { - let offset = output_buf.len(); - 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); - cursor.set_position(offset as u64); - let len = self - .decompressor - .decompress_to_buffer(input_buf, &mut cursor)?; + if let Some(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) + }) { + output_buf.reserve(len); + } + + let mut decoder = + zstd::stream::Decoder::with_context(input_buf, &mut self.dctx).single_frame(); + + // The default Read::read_to_end impl is acceptable; + // it reads directly into the Vec most of the time. + // Using raw DCtx here would be more annoying than the compress path. + let len = decoder.read_to_end(output_buf)?; Ok(len) } fn compress(&mut self, input_buf: &[u8], output_buf: &mut Vec) -> Result<()> { - let offset = output_buf.len(); - let len = zstd::zstd_safe::compress_bound(input_buf.len()); - output_buf.reserve(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(()) + compress_to_vec(&mut self.cctx, input_buf, output_buf).map_err(|code| { + let msg = zstd_safe::get_error_name(code); + std::io::Error::new(std::io::ErrorKind::Other, msg.to_string()).into() + }) } } } From 139281dc15f43c9bbcb52f43417c72f6a5e8d967 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 29 Jul 2026 16:38:03 +0200 Subject: [PATCH 2/8] simplify LZ4Codec writing and reading --- parquet/src/compression.rs | 23 ++--------------------- 1 file changed, 2 insertions(+), 21 deletions(-) diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index b90a45086add..5a89ba794c65 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -446,8 +446,6 @@ mod lz4_codec { use crate::compression::Codec; use crate::errors::{ParquetError, Result}; - const LZ4_BUFFER_SIZE: usize = 4096; - /// Codec for LZ4 compression algorithm. pub struct LZ4Codec {} @@ -466,30 +464,13 @@ mod lz4_codec { _uncompress_size: Option, ) -> Result { let mut decoder = lz4_flex::frame::FrameDecoder::new(input_buf); - let mut buffer: [u8; LZ4_BUFFER_SIZE] = [0; LZ4_BUFFER_SIZE]; - let mut total_len = 0; - loop { - let len = decoder.read(&mut buffer)?; - if len == 0 { - break; - } - total_len += len; - output_buf.write_all(&buffer[0..len])?; - } + let total_len = decoder.read_to_end(output_buf)?; Ok(total_len) } fn compress(&mut self, input_buf: &[u8], output_buf: &mut Vec) -> Result<()> { let mut encoder = lz4_flex::frame::FrameEncoder::new(output_buf); - let mut from = 0; - loop { - let to = std::cmp::min(from + LZ4_BUFFER_SIZE, input_buf.len()); - encoder.write_all(&input_buf[from..to])?; - from += LZ4_BUFFER_SIZE; - if from >= input_buf.len() { - break; - } - } + encoder.write_all(input_buf)?; match encoder.finish() { Ok(_) => Ok(()), Err(e) => Err(ParquetError::External(Box::new(e))), From 65223f0033c752dc603abd530e7bf3106ae5b92c Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 29 Jul 2026 16:53:27 +0200 Subject: [PATCH 3/8] fix clippy --- 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 5a89ba794c65..fbd765121f38 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -573,7 +573,7 @@ mod zstd_codec { fn compress(&mut self, input_buf: &[u8], output_buf: &mut Vec) -> Result<()> { compress_to_vec(&mut self.cctx, input_buf, output_buf).map_err(|code| { let msg = zstd_safe::get_error_name(code); - std::io::Error::new(std::io::ErrorKind::Other, msg.to_string()).into() + std::io::Error::other(msg.to_string()).into() }) } } From 0ee7ac3201faa07d9f188e96ce012515c95476b1 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 29 Jul 2026 18:07:19 +0200 Subject: [PATCH 4/8] reserve later when compressing zstd --- parquet/src/compression.rs | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index fbd765121f38..b0278145cf47 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -525,10 +525,13 @@ mod zstd_codec { cctx.set_pledged_src_size(Some(input_buf.len() as u64))?; let mut input = zstd_safe::InBuffer::around(input_buf); - while input.pos < input.src.len() { + loop { let mut output = zstd_safe::OutBuffer::around_pos(output_buf, output_buf.len()); let end_op = zstd_safe::zstd_sys::ZSTD_EndDirective::ZSTD_e_continue; let to_flush = cctx.compress_stream2(&mut output, &mut input, end_op)?; + if input.pos == input.src.len() { + break; // let the end_stream loop below call reserve_exact with the finalized amount + } output_buf.reserve(to_flush); } @@ -565,7 +568,8 @@ mod zstd_codec { // The default Read::read_to_end impl is acceptable; // it reads directly into the Vec most of the time. - // Using raw DCtx here would be more annoying than the compress path. + // Using raw DCtx here would be more annoying than the compress path, + // but could be done if better Vec reserve control is desired in the future. let len = decoder.read_to_end(output_buf)?; Ok(len) } From 26d5ad0b99cb83b6245a455045822ee8c72d1d69 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Wed, 29 Jul 2026 18:12:15 +0200 Subject: [PATCH 5/8] use uncompress_size in codecs --- parquet/src/compression.rs | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index b0278145cf47..58d018f6c5fe 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -281,8 +281,11 @@ mod gzip_codec { &mut self, input_buf: &[u8], output_buf: &mut Vec, - _uncompress_size: Option, + uncompress_size: Option, ) -> Result { + if let Some(len) = uncompress_size { + output_buf.reserve(len); + } let mut decoder = read::MultiGzDecoder::new(input_buf); decoder.read_to_end(output_buf).map_err(|e| e.into()) } @@ -389,8 +392,10 @@ mod brotli_codec { output_buf: &mut Vec, uncompress_size: Option, ) -> Result { - let buffer_size = uncompress_size.unwrap_or(BROTLI_DEFAULT_BUFFER_SIZE); - brotli::Decompressor::new(input_buf, buffer_size) + if let Some(len) = uncompress_size { + output_buf.reserve(len); + } + brotli::Decompressor::new(input_buf, BROTLI_DEFAULT_BUFFER_SIZE) .read_to_end(output_buf) .map_err(|e| e.into()) } @@ -461,8 +466,11 @@ mod lz4_codec { &mut self, input_buf: &[u8], output_buf: &mut Vec, - _uncompress_size: Option, + uncompress_size: Option, ) -> Result { + if let Some(len) = uncompress_size { + output_buf.reserve(len); + } let mut decoder = lz4_flex::frame::FrameDecoder::new(input_buf); let total_len = decoder.read_to_end(output_buf)?; Ok(total_len) From 9ccc3893680f119f6f36590717b8accc6640d982 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Thu, 30 Jul 2026 18:49:02 +0200 Subject: [PATCH 6/8] decode all frames in ZSTDCodec --- parquet/src/compression.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index 58d018f6c5fe..bf2f3719d9b0 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -571,8 +571,7 @@ mod zstd_codec { output_buf.reserve(len); } - let mut decoder = - zstd::stream::Decoder::with_context(input_buf, &mut self.dctx).single_frame(); + let mut decoder = zstd::stream::Decoder::with_context(input_buf, &mut self.dctx); // The default Read::read_to_end impl is acceptable; // it reads directly into the Vec most of the time. From 1dc6c79afb04925081c546c8a68dca9c4492b948 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Thu, 30 Jul 2026 19:52:20 +0200 Subject: [PATCH 7/8] enable experimental zstd, reserve generously and go back to zstd bulk decode --- parquet/Cargo.toml | 2 +- parquet/src/compression.rs | 48 ++++++++++++++++++++++---------------- 2 files changed, 29 insertions(+), 21 deletions(-) diff --git a/parquet/Cargo.toml b/parquet/Cargo.toml index 0ae0182263ed..303e7a0e730c 100644 --- a/parquet/Cargo.toml +++ b/parquet/Cargo.toml @@ -57,7 +57,7 @@ brotli = { version = "8.0", default-features = false, features = ["std"], option # To use `flate2` you must enable either the `flate2-zlib-rs` or `flate2-rust_backend` backends flate2 = { version = "1.1", default-features = false, optional = true } lz4_flex = { version = "0.14", default-features = false, features = ["std", "frame"], optional = true } -zstd = { version = "0.13", optional = true, default-features = false } +zstd = { version = "0.13", optional = true, default-features = false, features = ["experimental"] } chrono = { workspace = true } num-bigint = { version = "0.5", default-features = false } num-integer = { version = "0.1.46", default-features = false, features = ["std"] } diff --git a/parquet/src/compression.rs b/parquet/src/compression.rs index bf2f3719d9b0..c28521e9a65f 100644 --- a/parquet/src/compression.rs +++ b/parquet/src/compression.rs @@ -496,7 +496,7 @@ mod zstd_codec { use crate::compression::{Codec, ZstdLevel}; use crate::errors::Result; - use std::io::Read; + use std::io::Cursor; /// Codec for Zstandard compression algorithm. /// @@ -522,6 +522,11 @@ mod zstd_codec { } } + fn map_error_code(code: zstd_safe::ErrorCode) -> crate::errors::ParquetError { + let msg = zstd_safe::get_error_name(code); + std::io::Error::other(msg.to_string()).into() + } + /// Avoids zstd crate abstractions to minimize redundant copies; /// [zstd::stream::Encoder] uses a buffered writer which is pointless for Vec. fn compress_to_vec( @@ -530,7 +535,10 @@ mod zstd_codec { output_buf: &mut Vec, ) -> std::result::Result<(), zstd_safe::ErrorCode> { cctx.reset(zstd_safe::ResetDirective::SessionOnly)?; + // causes frame header to include size cctx.set_pledged_src_size(Some(input_buf.len() as u64))?; + // zstd can avoid allocating and copying into an input window + cctx.set_parameter(zstd_safe::CParameter::StableInBuffer(true))?; let mut input = zstd_safe::InBuffer::around(input_buf); loop { @@ -545,7 +553,8 @@ mod zstd_codec { loop { let mut output = zstd_safe::OutBuffer::around_pos(output_buf, output_buf.len()); - let to_flush = cctx.end_stream(&mut output)?; + let end_op = zstd_safe::zstd_sys::ZSTD_EndDirective::ZSTD_e_end; + let to_flush = cctx.compress_stream2(&mut output, &mut input, end_op)?; if to_flush == 0 { break; } @@ -561,31 +570,30 @@ mod zstd_codec { output_buf: &mut Vec, uncompress_size: Option, ) -> Result { - if let Some(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) - }) { - output_buf.reserve(len); + if let Some(len) = uncompress_size { + output_buf.reserve_exact(len); + } else if let Some(len) = zstd_safe::find_decompressed_size(input_buf) + .map_err(|err| std::io::Error::other(err.to_string()))? + { + output_buf.reserve_exact(len as usize); + } else { + let len = zstd_safe::decompress_bound(input_buf).map_err(map_error_code)?; + output_buf.reserve(len as usize); } - let mut decoder = zstd::stream::Decoder::with_context(input_buf, &mut self.dctx); + let offset = output_buf.len(); + let mut cursor = Cursor::new(output_buf); + cursor.set_position(offset as u64); - // The default Read::read_to_end impl is acceptable; - // it reads directly into the Vec most of the time. - // Using raw DCtx here would be more annoying than the compress path, - // but could be done if better Vec reserve control is desired in the future. - let len = decoder.read_to_end(output_buf)?; + let len = self + .dctx + .decompress(&mut cursor, input_buf) + .map_err(map_error_code)?; Ok(len) } fn compress(&mut self, input_buf: &[u8], output_buf: &mut Vec) -> Result<()> { - compress_to_vec(&mut self.cctx, input_buf, output_buf).map_err(|code| { - let msg = zstd_safe::get_error_name(code); - std::io::Error::other(msg.to_string()).into() - }) + compress_to_vec(&mut self.cctx, input_buf, output_buf).map_err(map_error_code) } } } From 61aaf89a8f8a65db8865a1485770d7a6edaca6b3 Mon Sep 17 00:00:00 2001 From: Michal Piatkowski <291740709+MassivePizza@users.noreply.github.com> Date: Fri, 31 Jul 2026 15:03:12 +0200 Subject: [PATCH 8/8] fix cargo check --- parquet/Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parquet/Cargo.toml b/parquet/Cargo.toml index 303e7a0e730c..7e8a22e629ea 100644 --- a/parquet/Cargo.toml +++ b/parquet/Cargo.toml @@ -85,7 +85,7 @@ insta = { workspace = true, default-features = true } brotli = { version = "8.0", default-features = false, features = ["std"] } flate2 = { version = "1.0", default-features = false, features = ["rust_backend"] } lz4_flex = { version = "0.14", default-features = false, features = ["std", "frame"] } -zstd = { version = "0.13", default-features = false } +zstd = { version = "0.13", default-features = false, features = ["experimental"] } serde_json = { version = "1.0", features = ["std"], default-features = false } arrow = { workspace = true, features = ["ipc", "test_utils", "prettyprint", "json"] } arrow-cast = { workspace = true }