From f599828d2a5585aae79043a340b426c7cd6e22ec Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Mon, 28 Sep 2026 20:14:47 +0800 Subject: [PATCH] feat(spec): support the Avro deflate codec The fast OCF reader (`spec/avro/ocf.rs`) only accepted the null, snappy and zstandard codecs, and the object-file writer (`spec/avro/objects_file.rs`) only produced those three. `deflate` is a standard Avro codec that Paimon accepts through `CodecFactory.fromString` (e.g. `manifest.compression=deflate` or `avro.output.codec=deflate`), so a manifest or object file written that way failed to read with `avro ocf: unsupported codec: deflate`, and could not be written at all. Decode deflate blocks with `flate2` (the Avro deflate codec is raw DEFLATE, matching the `miniz_oxide` output apache-avro produces on the write path) and map the `deflate` compression name to `Codec::Deflate`. Round-trip tests cover the reader (write with apache-avro, read back and decode the record) and a write-then-read of a manifest entry through the fast decoder. --- crates/paimon/src/spec/avro/ocf.rs | 49 ++++++++++++++++++++++++++ crates/paimon/src/spec/objects_file.rs | 17 ++++++++- 2 files changed, 65 insertions(+), 1 deletion(-) diff --git a/crates/paimon/src/spec/avro/ocf.rs b/crates/paimon/src/spec/avro/ocf.rs index 58b89ca7b..1004c23df 100644 --- a/crates/paimon/src/spec/avro/ocf.rs +++ b/crates/paimon/src/spec/avro/ocf.rs @@ -36,6 +36,7 @@ pub struct OcfHeader { #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum OcfCodec { Null, + Deflate, Snappy, Zstandard, } @@ -113,6 +114,20 @@ impl<'a> OcfBlockIter<'a> { fn decompress(&mut self, data: &'a [u8]) -> crate::Result> { match self.codec { OcfCodec::Null => Ok(Cow::Borrowed(data)), + OcfCodec::Deflate => { + // The Avro `deflate` codec is raw DEFLATE (RFC 1951), with no + // zlib or gzip wrapper. + use std::io::Read; + let mut decoder = flate2::read::DeflateDecoder::new(data); + let mut decompressed = Vec::new(); + decoder + .read_to_end(&mut decompressed) + .map_err(|e| Error::UnexpectedError { + message: format!("avro ocf: deflate decompression failed: {e}"), + source: None, + })?; + Ok(Cow::Owned(decompressed)) + } OcfCodec::Snappy => { if data.len() < 4 { return Err(Error::UnexpectedError { @@ -176,6 +191,7 @@ pub fn parse_ocf_streaming(bytes: &[u8]) -> crate::Result<(OcfHeader, OcfBlockIt let codec = match meta.get("avro.codec").map(|s| s.as_str()) { None | Some("null") => OcfCodec::Null, + Some("deflate") => OcfCodec::Deflate, Some("snappy") => OcfCodec::Snappy, Some("zstandard") => OcfCodec::Zstandard, Some(other) => { @@ -309,4 +325,37 @@ mod tests { assert_eq!(blocks.len(), 1); assert_eq!(blocks[0].object_count, 1); } + + #[test] + fn test_parse_ocf_deflate() { + use apache_avro::{from_avro_datum, types::Value, Codec, DeflateSettings, Schema, Writer}; + + let schema = Schema::parse_str( + r#"{"type": "record", "name": "test", "fields": [{"name": "x", "type": "long"}]}"#, + ) + .unwrap(); + let mut writer = Writer::with_codec( + &schema, + Vec::new(), + Codec::Deflate(DeflateSettings::default()), + ); + let mut record = apache_avro::types::Record::new(&schema).unwrap(); + record.put("x", 424242i64); + writer.append(record).unwrap(); + let bytes = writer.into_inner().unwrap(); + + let (header, blocks) = parse_ocf(&bytes).unwrap(); + assert_eq!(header.codec, OcfCodec::Deflate); + assert_eq!(blocks.len(), 1); + assert_eq!(blocks[0].object_count, 1); + + // The decompressed block must decode back to the original record, which + // proves the deflate output is correct rather than merely non-erroring. + let mut cursor = blocks[0].data.as_ref(); + let value = from_avro_datum(&schema, &mut cursor, None).unwrap(); + assert_eq!( + value, + Value::Record(vec![("x".to_string(), Value::Long(424242))]) + ); + } } diff --git a/crates/paimon/src/spec/objects_file.rs b/crates/paimon/src/spec/objects_file.rs index c30ecc8ff..a4e60c4f3 100644 --- a/crates/paimon/src/spec/objects_file.rs +++ b/crates/paimon/src/spec/objects_file.rs @@ -15,7 +15,9 @@ // specific language governing permissions and limitations // under the License. -use apache_avro::{from_value, to_value, Codec, Reader, Schema, Writer, ZstandardSettings}; +use apache_avro::{ + from_value, to_value, Codec, DeflateSettings, Reader, Schema, Writer, ZstandardSettings, +}; use serde::de::DeserializeOwned; use serde::Serialize; @@ -70,6 +72,7 @@ pub(crate) fn avro_codec(compression: &str) -> crate::Result { match compression.to_ascii_lowercase().as_str() { "zstd" | "zstandard" => Ok(Codec::Zstandard(ZstandardSettings::default())), "null" | "none" | "uncompressed" => Ok(Codec::Null), + "deflate" => Ok(Codec::Deflate(DeflateSettings::default())), "snappy" => Ok(Codec::Snappy), other => Err(crate::Error::Unsupported { message: format!("Unsupported Avro compression: {other}"), @@ -333,6 +336,18 @@ mod tests { assert_eq!(original, decoded); } + #[test] + fn test_roundtrip_manifest_entry_deflate() { + // Writing with the deflate codec and reading back through the fast + // decoder must round-trip, exercising both the write-side codec mapping + // and the read-side deflate decompression. + let original = vec![manifest_entry()]; + let bytes = + to_avro_bytes_with_compression(MANIFEST_ENTRY_SCHEMA, &original, "deflate").unwrap(); + let decoded = from_avro_bytes_fast::(&bytes).unwrap(); + assert_eq!(original, decoded); + } + #[test] fn test_roundtrip_index_manifest_field_order() { let entry: IndexManifestEntry = serde_json::from_value(serde_json::json!({