use crate::Result;
use crate::error::CoreError;
use crate::metadata::commit::HoodieWriteStat;
use apache_avro_derive::AvroSchema as DeriveAvroSchema;
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use std::collections::HashMap;
#[derive(Debug, Clone, Serialize, Deserialize, DeriveAvroSchema)]
#[serde(rename_all = "camelCase")]
#[avro(namespace = "org.apache.hudi.avro.model")]
pub struct HoodieReplaceCommitMetadata {
pub version: Option<i32>,
#[avro(rename = "operationType")]
pub operation_type: Option<String>,
#[avro(rename = "partitionToWriteStats")]
pub partition_to_write_stats: Option<HashMap<String, Vec<HoodieWriteStat>>>,
pub compacted: Option<bool>,
#[avro(rename = "extraMetadata")]
pub extra_metadata: Option<HashMap<String, String>>,
#[avro(rename = "partitionToReplaceFileIds")]
pub partition_to_replace_file_ids: Option<HashMap<String, Vec<String>>>,
}
impl HoodieReplaceCommitMetadata {
pub fn from_json_map(map: &Map<String, Value>) -> Result<Self> {
serde_json::from_value(Value::Object(map.clone()))
.map_err(|e| CoreError::CommitMetadata(format!("Failed to parse commit metadata: {e}")))
}
pub fn from_json_bytes(bytes: &[u8]) -> Result<Self> {
serde_json::from_slice(bytes)
.map_err(|e| CoreError::CommitMetadata(format!("Failed to parse commit metadata: {e}")))
}
pub fn iter_replace_file_ids(&self) -> impl Iterator<Item = (&String, &String)> {
self.partition_to_replace_file_ids
.as_ref()
.into_iter()
.flat_map(|replace_ids| {
replace_ids.iter().flat_map(|(partition, file_ids)| {
file_ids.iter().map(move |file_id| (partition, file_id))
})
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn test_parse_replace_commit() {
let json = json!({
"partitionToReplaceFileIds": {
"30": ["a-0"],
"20": ["b-0", "b-1"],
"": ["c-0"]
},
"extraMetadata": {"k":"v"},
"version": 1,
"operationType": "REPLACE_COMMIT"
});
let metadata: HoodieReplaceCommitMetadata = serde_json::from_value(json).unwrap();
let ids: Vec<(&String, &String)> = metadata.iter_replace_file_ids().collect();
assert_eq!(ids.len(), 4);
}
#[test]
fn test_from_json_bytes() {
let json_str = r#"{
"partitionToReplaceFileIds": {
"30": ["a-0"],
"20": ["b-0"]
},
"version": 1,
"operationType": "REPLACE_COMMIT"
}"#;
let metadata = HoodieReplaceCommitMetadata::from_json_bytes(json_str.as_bytes()).unwrap();
assert_eq!(metadata.version, Some(1));
assert_eq!(metadata.operation_type, Some("REPLACE_COMMIT".to_string()));
}
#[test]
fn test_from_json_bytes_invalid() {
let invalid_json = b"invalid json";
let result = HoodieReplaceCommitMetadata::from_json_bytes(invalid_json);
assert!(result.is_err());
assert!(matches!(result, Err(CoreError::CommitMetadata(_))));
}
#[test]
fn test_iter_replace_file_ids_empty() {
let json = json!({});
let metadata: HoodieReplaceCommitMetadata = serde_json::from_value(json).unwrap();
let count = metadata.iter_replace_file_ids().count();
assert_eq!(count, 0);
}
}