From 8492ad980ae9df00beda1383e63a6dc7f4cc0d77 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Mon, 14 Sep 2026 15:27:33 +0800 Subject: [PATCH 1/2] fix(fm_index): store blocks in Java's codec envelope, not the SST one --- crates/paimon/src/btree/block.rs | 192 ++++++++++++++++++++------- crates/paimon/src/btree/mod.rs | 2 +- crates/paimon/src/fm_index/format.rs | 70 +++++++++- 3 files changed, 211 insertions(+), 53 deletions(-) diff --git a/crates/paimon/src/btree/block.rs b/crates/paimon/src/btree/block.rs index 18e4ae35d..9e41cbcef 100644 --- a/crates/paimon/src/btree/block.rs +++ b/crates/paimon/src/btree/block.rs @@ -71,6 +71,11 @@ impl BlockCompressionType { /// Compress a Java Paimon block, including the outer uncompressed-size varint /// and the LZ4/LZO codec envelope when applicable. Compression is retained only /// when it saves at least 12.5%, matching Java's block writer. +/// +/// The varint belongs to the SST and bitmap-index *block formats*, not to the +/// codec: Java writes it in `SstFileWriter#writeBlock` and +/// `BitmapGlobalIndexFormat#encodeBlock`. Formats that store `BlockCompressor` +/// output verbatim want [`compress_codec_block`] instead. pub(crate) fn compress_block( data: &[u8], compression_type: BlockCompressionType, @@ -80,21 +85,52 @@ pub(crate) fn compress_block( return Ok((Cow::Borrowed(data), BlockCompressionType::None)); } + let mut encoded = Vec::with_capacity(13 + data.len()); + encode_var_int(&mut encoded, block_length(data.len())?)?; + append_codec_envelope(&mut encoded, data, compression_type, compression_level)?; + Ok(retain_if_smaller(encoded, data, compression_type)) +} + +/// Compress a block the way a format that stores `BlockCompressor` output +/// verbatim expects: the LZ4/LZO codec header when applicable, and no outer +/// uncompressed-size varint. Java's `FMIndexFile#writeBlock` stores the +/// compressor's bytes unchanged and records `storedLength = compressedLength`. +pub(crate) fn compress_codec_block( + data: &[u8], + compression_type: BlockCompressionType, + compression_level: i32, +) -> io::Result<(Cow<'_, [u8]>, BlockCompressionType)> { + if compression_type == BlockCompressionType::None { + return Ok((Cow::Borrowed(data), BlockCompressionType::None)); + } + + let mut encoded = Vec::with_capacity(8 + data.len()); + append_codec_envelope(&mut encoded, data, compression_type, compression_level)?; + Ok(retain_if_smaller(encoded, data, compression_type)) +} + +fn block_length(len: usize) -> io::Result { + i32::try_from(len) + .map_err(|_| io::Error::new(io::ErrorKind::InvalidData, "Block is larger than i32::MAX")) +} + +/// Append what Java's `BlockCompressor#compress` writes: for LZ4 and LZO the +/// `[le32 compressedLength][le32 srcLen]` header and then the payload, for ZSTD +/// the bare frame. +fn append_codec_envelope( + encoded: &mut Vec, + data: &[u8], + compression_type: BlockCompressionType, + compression_level: i32, +) -> io::Result<()> { let payload = match compression_type { - BlockCompressionType::None => unreachable!("handled above"), + BlockCompressionType::None => unreachable!("callers handle None"), BlockCompressionType::Zstd => zstd::bulk::compress(data, compression_level) .map_err(|error| io::Error::new(io::ErrorKind::InvalidInput, error))?, BlockCompressionType::Lz4 => lz4_flex::block::compress(data), BlockCompressionType::Lzo => lzokay_native::compress(data) .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?, }; - let mut encoded = Vec::with_capacity(13 + payload.len()); - encode_var_int( - &mut encoded, - i32::try_from(data.len()).map_err(|_| { - io::Error::new(io::ErrorKind::InvalidData, "Block is larger than i32::MAX") - })?, - )?; if matches!( compression_type, BlockCompressionType::Lz4 | BlockCompressionType::Lzo @@ -109,46 +145,62 @@ pub(crate) fn compress_block( })? .to_le_bytes(), ); - encoded.extend_from_slice( - &i32::try_from(data.len()) - .map_err(|_| { - io::Error::new(io::ErrorKind::InvalidData, "Block is larger than i32::MAX") - })? - .to_le_bytes(), - ); + encoded.extend_from_slice(&block_length(data.len())?.to_le_bytes()); } encoded.extend_from_slice(&payload); + Ok(()) +} + +/// Java keeps the compressed form only when it saves at least 12.5%; both +/// `SstFileWriter#writeBlock` and `FMIndexFile#writeBlock` compare against the +/// bytes they are about to store, so each caller applies this to its own output. +fn retain_if_smaller<'a>( + encoded: Vec, + data: &'a [u8], + compression_type: BlockCompressionType, +) -> (Cow<'a, [u8]>, BlockCompressionType) { if encoded.len() < data.len() - (data.len() / 8) { - Ok((Cow::Owned(encoded), compression_type)) + (Cow::Owned(encoded), compression_type) } else { - Ok((Cow::Borrowed(data), BlockCompressionType::None)) + (Cow::Borrowed(data), BlockCompressionType::None) } } /// Decode the payload written by [`compress_block`] or Java's corresponding -/// block compressors. +/// block compressors: an uncompressed-size varint followed by the codec block. pub(crate) fn decompress_block( data: &[u8], compression_type: BlockCompressionType, ) -> io::Result> { - decompress_block_inner(data, compression_type, None) -} + if compression_type == BlockCompressionType::None { + return Ok(data.to_vec()); + } -pub(crate) fn decompress_block_with_expected_size( - data: &[u8], - compression_type: BlockCompressionType, - expected_size: usize, -) -> io::Result> { - decompress_block_inner(data, compression_type, Some(expected_size)) + let mut cursor = Cursor::new(data); + let uncompressed_size = usize::try_from(decode_var_int(&mut cursor)?).map_err(|_| { + io::Error::new( + io::ErrorKind::InvalidData, + "Invalid negative uncompressed block size", + ) + })?; + let compressed_start = cursor.position() as usize; + decompress_codec_payload( + &data[compressed_start..], + compression_type, + uncompressed_size, + ) } -fn decompress_block_inner( +/// Decode a codec block stored without the outer uncompressed-size varint, the +/// way Java's `FMIndexFile#decodeStoredBlock` reads it. The decoded size comes +/// from the caller's own metadata, so it also bounds the allocation. +pub(crate) fn decompress_codec_block( data: &[u8], compression_type: BlockCompressionType, - expected_size: Option, + uncompressed_size: usize, ) -> io::Result> { if compression_type == BlockCompressionType::None { - if expected_size.is_some_and(|expected_size| expected_size != data.len()) { + if uncompressed_size != data.len() { return Err(io::Error::new( io::ErrorKind::InvalidData, "Uncompressed block size does not match its metadata", @@ -156,24 +208,16 @@ fn decompress_block_inner( } return Ok(data.to_vec()); } + decompress_codec_payload(data, compression_type, uncompressed_size) +} - let mut cursor = Cursor::new(data); - let uncompressed_size = usize::try_from(decode_var_int(&mut cursor)?).map_err(|_| { - io::Error::new( - io::ErrorKind::InvalidData, - "Invalid negative uncompressed block size", - ) - })?; - if expected_size.is_some_and(|expected_size| expected_size != uncompressed_size) { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "Compressed block size does not match its metadata", - )); - } - let compressed_start = cursor.position() as usize; - let compressed = &data[compressed_start..]; +fn decompress_codec_payload( + compressed: &[u8], + compression_type: BlockCompressionType, + uncompressed_size: usize, +) -> io::Result> { match compression_type { - BlockCompressionType::None => unreachable!("handled above"), + BlockCompressionType::None => unreachable!("callers handle None"), BlockCompressionType::Zstd => { let mut decompressed = vec![0u8; uncompressed_size]; let actual = zstd::bulk::decompress_to_buffer(compressed, &mut decompressed) @@ -681,13 +725,63 @@ mod tests { fn constrained_decompression_rejects_header_size_before_decoding() { let source = vec![0u8; 4096]; let (compressed, compression) = - compress_block(&source, BlockCompressionType::Zstd, 1).unwrap(); + compress_codec_block(&source, BlockCompressionType::Zstd, 1).unwrap(); assert_eq!(compression, BlockCompressionType::Zstd); - let error = decompress_block_with_expected_size(&compressed, compression, 64) - .expect_err("the validated outer size must constrain allocation"); + let error = decompress_codec_block(&compressed, compression, 64) + .expect_err("the caller's metadata size must constrain allocation"); assert_eq!(error.kind(), io::ErrorKind::InvalidData); - assert!(error.to_string().contains("does not match its metadata")); + } + + #[test] + fn codec_block_is_the_varint_block_without_its_prefix() { + // Java's FM index stores `BlockCompressor` output verbatim, while the SST and + // bitmap formats prepend an uncompressed-size varint to the same bytes. + let source: Vec = (0..32768u32).map(|value| (value / 97) as u8).collect(); + for compression in [ + BlockCompressionType::Zstd, + BlockCompressionType::Lz4, + BlockCompressionType::Lzo, + ] { + let (with_prefix, kept) = compress_block(&source, compression, 1).unwrap(); + assert_eq!(kept, compression, "{compression:?} must stay compressed"); + let (codec_only, kept) = compress_codec_block(&source, compression, 1).unwrap(); + assert_eq!(kept, compression, "{compression:?} must stay compressed"); + + let mut prefix = Vec::new(); + encode_var_int(&mut prefix, source.len() as i32).unwrap(); + assert_eq!( + with_prefix.as_ref(), + [prefix.as_slice(), codec_only.as_ref()].concat(), + "{compression:?}" + ); + assert_eq!( + decompress_codec_block(codec_only.as_ref(), compression, source.len()).unwrap(), + source, + "{compression:?}" + ); + assert_eq!( + decompress_block(with_prefix.as_ref(), compression).unwrap(), + source, + "{compression:?}" + ); + } + } + + #[test] + fn codec_block_rejects_the_varint_prefixed_form() { + // Reading a Java FM block with the SST reader, or the reverse, must fail + // rather than silently return the wrong bytes. + let source = vec![7u8; 4096]; + let (with_prefix, _) = compress_block(&source, BlockCompressionType::Lz4, 1).unwrap(); + assert!(decompress_codec_block( + with_prefix.as_ref(), + BlockCompressionType::Lz4, + source.len() + ) + .is_err()); + let (codec_only, _) = compress_codec_block(&source, BlockCompressionType::Lz4, 1).unwrap(); + assert!(decompress_block(codec_only.as_ref(), BlockCompressionType::Lz4).is_err()); } #[test] diff --git a/crates/paimon/src/btree/mod.rs b/crates/paimon/src/btree/mod.rs index 37b7040a6..ad66ddfed 100644 --- a/crates/paimon/src/btree/mod.rs +++ b/crates/paimon/src/btree/mod.rs @@ -48,7 +48,7 @@ mod writer; pub use block::BlockCompressionType; pub(crate) use block::{ - compress_block, compute_crc32, decompress_block, decompress_block_with_expected_size, + compress_block, compress_codec_block, compute_crc32, decompress_block, decompress_codec_block, }; pub use footer::BTreeFileFooter; pub use key_serde::{make_key_comparator, serialize_datum}; diff --git a/crates/paimon/src/fm_index/format.rs b/crates/paimon/src/fm_index/format.rs index d4bba5684..62dbbe2f9 100644 --- a/crates/paimon/src/fm_index/format.rs +++ b/crates/paimon/src/fm_index/format.rs @@ -18,7 +18,7 @@ //! Java-compatible V1 FM-index container format. use crate::btree::{ - compress_block, compute_crc32, decompress_block_with_expected_size, BlockCompressionType, + compress_codec_block, compute_crc32, decompress_codec_block, BlockCompressionType, }; use crate::io::FileRead; use std::io; @@ -404,8 +404,11 @@ pub(crate) fn write_block( let offset = base_offset .checked_add(out.len() as u64) .ok_or_else(|| invalid_input("FM output offset overflow"))?; + // Java's `FMIndexFile#writeBlock` stores the compressor's output verbatim and + // records `storedLength = compressedLength`, so this format carries no outer + // uncompressed-size varint — unlike the SST and bitmap block formats. let (stored, actual_compression) = - compress_block(uncompressed, compression, compression_level)?; + compress_codec_block(uncompressed, compression, compression_level)?; let checksum = compute_crc32(stored.as_ref(), actual_compression); let stored_length = stored.len(); out.extend_from_slice(stored.as_ref()); @@ -1078,7 +1081,7 @@ fn decode_stored_block(stored: &[u8], block: BlockInfo) -> io::Result> { } return Ok(stored.to_vec()); } - decompress_block_with_expected_size(stored, block.compression, block.uncompressed_length) + decompress_codec_block(stored, block.compression, block.uncompressed_length) } fn validate_footer_common(bytes: &[u8], magic: u32, scope: &str) -> io::Result<()> { @@ -1638,4 +1641,65 @@ mod tests { 0x00c4_9e49 ); } + + /// A block the size of a packed rank block, compressible enough that both + /// writers keep the compressed form (Java and Rust both require a 12.5% saving). + fn compressible_block() -> Vec { + (0..BLOCK_WORDS as u32 * 8) + .map(|i| (i / 97) as u8) + .collect() + } + + /// What Java's `FMIndexFile#writeBlock` stores: `BlockCompressor#compress` + /// output verbatim. For LZ4 that is `[le32 compressedLength][le32 srcLen]` + /// followed by the payload; for ZSTD the bare frame. + fn java_stored_block(plain: &[u8], compression: BlockCompressionType) -> Vec { + match compression { + BlockCompressionType::Lz4 => { + let payload = lz4_flex::block::compress(plain); + let mut stored = Vec::with_capacity(8 + payload.len()); + stored.extend_from_slice(&(payload.len() as i32).to_le_bytes()); + stored.extend_from_slice(&(plain.len() as i32).to_le_bytes()); + stored.extend_from_slice(&payload); + stored + } + BlockCompressionType::Zstd => zstd::bulk::compress(plain, 1).unwrap(), + other => panic!("unsupported in this fixture: {other:?}"), + } + } + + #[test] + fn stores_blocks_in_javas_codec_envelope() { + let plain = compressible_block(); + for compression in [BlockCompressionType::Lz4, BlockCompressionType::Zstd] { + let expected = java_stored_block(&plain, compression); + + // Read side: a block written by Java must decode. + let block = BlockInfo { + offset: 0, + stored_length: expected.len(), + uncompressed_length: plain.len(), + compression, + checksum: compute_crc32(&expected, compression), + }; + assert_eq!( + decode_stored_block(&expected, block).unwrap(), + plain, + "{compression:?}" + ); + + // Write side: the bytes Rust stores must be the ones Java expects. + let mut out = Vec::new(); + let info = write_block(&mut out, 0, &plain, compression, 1).unwrap(); + assert_eq!(info.compression, compression, "{compression:?}"); + assert_eq!(info.stored_length, expected.len(), "{compression:?}"); + assert_eq!(info.uncompressed_length, plain.len(), "{compression:?}"); + assert_eq!(out, expected, "{compression:?}"); + assert_eq!( + decode_stored_block(&out, info).unwrap(), + plain, + "{compression:?}" + ); + } + } } From ab95c47c3dbf5a4a56faf6fb74a54367d3fe3ff9 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Mon, 14 Sep 2026 15:54:57 +0800 Subject: [PATCH 2/2] reserve the codec envelope capacity after compressing --- crates/paimon/src/btree/block.rs | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/crates/paimon/src/btree/block.rs b/crates/paimon/src/btree/block.rs index 9e41cbcef..950995f75 100644 --- a/crates/paimon/src/btree/block.rs +++ b/crates/paimon/src/btree/block.rs @@ -85,7 +85,7 @@ pub(crate) fn compress_block( return Ok((Cow::Borrowed(data), BlockCompressionType::None)); } - let mut encoded = Vec::with_capacity(13 + data.len()); + let mut encoded = Vec::with_capacity(5); encode_var_int(&mut encoded, block_length(data.len())?)?; append_codec_envelope(&mut encoded, data, compression_type, compression_level)?; Ok(retain_if_smaller(encoded, data, compression_type)) @@ -104,7 +104,7 @@ pub(crate) fn compress_codec_block( return Ok((Cow::Borrowed(data), BlockCompressionType::None)); } - let mut encoded = Vec::with_capacity(8 + data.len()); + let mut encoded = Vec::new(); append_codec_envelope(&mut encoded, data, compression_type, compression_level)?; Ok(retain_if_smaller(encoded, data, compression_type)) } @@ -135,6 +135,7 @@ fn append_codec_envelope( compression_type, BlockCompressionType::Lz4 | BlockCompressionType::Lzo ) { + encoded.reserve(8 + payload.len()); encoded.extend_from_slice( &i32::try_from(payload.len()) .map_err(|_| { @@ -146,6 +147,8 @@ fn append_codec_envelope( .to_le_bytes(), ); encoded.extend_from_slice(&block_length(data.len())?.to_le_bytes()); + } else { + encoded.reserve(payload.len()); } encoded.extend_from_slice(&payload); Ok(())