Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
195 changes: 146 additions & 49 deletions crates/paimon/src/btree/block.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -80,25 +85,57 @@ pub(crate) fn compress_block(
return Ok((Cow::Borrowed(data), BlockCompressionType::None));
}

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))
}

/// 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::new();
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> {
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<u8>,
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
) {
encoded.reserve(8 + payload.len());
encoded.extend_from_slice(
&i32::try_from(payload.len())
.map_err(|_| {
Expand All @@ -109,71 +146,81 @@ 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());
} else {
encoded.reserve(payload.len());
}
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<u8>,
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<Vec<u8>> {
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<Vec<u8>> {
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<usize>,
uncompressed_size: usize,
) -> io::Result<Vec<u8>> {
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",
));
}
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<Vec<u8>> {
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)
Expand Down Expand Up @@ -681,13 +728,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<u8> = (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]
Expand Down
2 changes: 1 addition & 1 deletion crates/paimon/src/btree/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down
70 changes: 67 additions & 3 deletions crates/paimon/src/fm_index/format.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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());
Expand Down Expand Up @@ -1078,7 +1081,7 @@ fn decode_stored_block(stored: &[u8], block: BlockInfo) -> io::Result<Vec<u8>> {
}
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<()> {
Expand Down Expand Up @@ -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<u8> {
(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<u8> {
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:?}"
);
}
}
}
Loading