Skip to content
Merged
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
49 changes: 49 additions & 0 deletions crates/paimon/src/spec/avro/ocf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ pub struct OcfHeader {
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OcfCodec {
Null,
Deflate,
Snappy,
Zstandard,
}
Expand Down Expand Up @@ -113,6 +114,20 @@ impl<'a> OcfBlockIter<'a> {
fn decompress(&mut self, data: &'a [u8]) -> crate::Result<Cow<'a, [u8]>> {
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 {
Expand Down Expand Up @@ -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) => {
Expand Down Expand Up @@ -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))])
);
}
}
17 changes: 16 additions & 1 deletion crates/paimon/src/spec/objects_file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -70,6 +72,7 @@ pub(crate) fn avro_codec(compression: &str) -> crate::Result<Codec> {
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}"),
Expand Down Expand Up @@ -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::<ManifestEntry>(&bytes).unwrap();
assert_eq!(original, decoded);
}

#[test]
fn test_roundtrip_index_manifest_field_order() {
let entry: IndexManifestEntry = serde_json::from_value(serde_json::json!({
Expand Down
Loading