use std::collections::{BTreeMap, HashMap};
use std::sync::Arc;
use arrow_array::cast::AsArray;
use arrow_array::types::Int32Type;
use arrow_array::{Array, ArrayRef, BooleanArray, Int64Array, MapArray, RecordBatch, StructArray};
use arrow_buffer::OffsetBuffer;
use arrow_schema::{DataType, Field, Fields};
use crate::Result;
use crate::error::CoreError;
use crate::file_group::reader_v2::buffered_record::{BufferedRecord, DeleteRecord, RecordPayload};
use crate::file_group::reader_v2::record_merger::BufferedRecordMerger;
use crate::file_group::reader_v2::resolver::{PAYLOAD_CLASS_KEYS, RECORD_MERGE_STRATEGY_ID_KEYS};
const METADATA_PAYLOAD_CLASS: &str = "org.apache.hudi.metadata.HoodieMetadataPayload";
const PAYLOAD_BASED_STRATEGY_ID: &str = "00000000-0000-0000-0000-000000000000";
const TYPE_ALL_PARTITIONS: i32 = 1;
const TYPE_FILES: i32 = 2;
const TYPE_COLUMN_STATS: i32 = 3;
const TYPE_PARTITION_STATS: i32 = 6;
const KEY_COLUMN: &str = "key";
const TYPE_COLUMN: &str = "type";
const FILESYSTEM_METADATA_COLUMN: &str = "filesystemMetadata";
const COLUMN_STATS_COLUMN: &str = "ColumnStatsMetadata";
const MIN_VALUE_FIELD: &str = "minValue";
const MAX_VALUE_FIELD: &str = "maxValue";
const IS_TIGHT_BOUND_FIELD: &str = "isTightBound";
const WRAPPED_VALUE_FIELD: &str = "value";
const COUNTER_FIELDS: [&str; 4] = [
"valueCount",
"nullCount",
"totalSize",
"totalUncompressedSize",
];
const FILE_SIZE_FIELD: &str = "size";
const FILE_IS_DELETED_FIELD: &str = "isDeleted";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CustomMerger {
MetadataPayload,
}
impl CustomMerger {
pub(crate) fn required_field_names(&self) -> &'static [&'static str] {
match self {
Self::MetadataPayload => &[KEY_COLUMN, TYPE_COLUMN],
}
}
}
pub fn resolve_custom_merger(table_config: &HashMap<String, String>) -> Option<CustomMerger> {
let lookup = |keys: &[&str]| -> String {
keys.iter()
.find_map(|k| table_config.get(*k))
.map(|s| s.trim().to_string())
.unwrap_or_default()
};
let payload_class = lookup(&PAYLOAD_CLASS_KEYS);
let strategy_id = lookup(&RECORD_MERGE_STRATEGY_ID_KEYS);
let defers_to_payload =
strategy_id.is_empty() || strategy_id.eq_ignore_ascii_case(PAYLOAD_BASED_STRATEGY_ID);
if !defers_to_payload {
return None;
}
match payload_class.as_str() {
METADATA_PAYLOAD_CLASS => Some(CustomMerger::MetadataPayload),
_ => None,
}
}
#[derive(Debug)]
pub struct MetadataPayloadMerger;
impl BufferedRecordMerger for MetadataPayloadMerger {
fn delta_merge(
&self,
new_record: &BufferedRecord,
existing_record: Option<&BufferedRecord>,
) -> Result<Option<BufferedRecord>> {
match existing_record {
Some(existing) => Ok(Some(self.final_merge(existing, new_record)?)),
None => Ok(Some(new_record.clone())),
}
}
fn delta_merge_delete(
&self,
delete_record: &DeleteRecord,
_existing_record: Option<&BufferedRecord>,
) -> Result<Option<DeleteRecord>> {
Ok(Some(delete_record.clone()))
}
fn final_merge(
&self,
older_record: &BufferedRecord,
newer_record: &BufferedRecord,
) -> Result<BufferedRecord> {
if newer_record.is_delete() || older_record.is_delete() {
return Ok(newer_record.clone());
}
let (Some(older), Some(newer)) = (
older_record.payload.get_record(),
newer_record.payload.get_record(),
) else {
return Ok(newer_record.clone());
};
match record_type(&newer)? {
TYPE_ALL_PARTITIONS | TYPE_FILES => {
let folded = fold_filesystem_metadata(&older, &newer)?;
Ok(BufferedRecord {
record_key: newer_record.record_key.clone(),
payload: RecordPayload::Owned(folded),
ordering_value: newer_record.ordering_value.clone(),
})
}
TYPE_COLUMN_STATS | TYPE_PARTITION_STATS => {
let folded = fold_column_stats(&older, &newer)?;
Ok(BufferedRecord {
record_key: newer_record.record_key.clone(),
payload: RecordPayload::Owned(folded),
ordering_value: newer_record.ordering_value.clone(),
})
}
_ => Ok(newer_record.clone()),
}
}
}
fn record_type(batch: &RecordBatch) -> Result<i32> {
let column = batch.column_by_name(TYPE_COLUMN).ok_or_else(|| {
CoreError::Unsupported(format!(
"A metadata record must carry a '{TYPE_COLUMN}' column to be merged."
))
})?;
let typed = column.as_primitive_opt::<Int32Type>().ok_or_else(|| {
CoreError::Unsupported(format!(
"The metadata '{TYPE_COLUMN}' column must be Int32, found {:?}.",
column.data_type()
))
})?;
if typed.is_empty() || typed.is_null(0) {
return Err(CoreError::Unsupported(
"A metadata record's partition type is null.".to_string(),
));
}
Ok(typed.value(0))
}
#[derive(Clone, Copy)]
struct FileInfo {
size: i64,
is_deleted: bool,
}
fn fold_filesystem_metadata(older: &RecordBatch, newer: &RecordBatch) -> Result<RecordBatch> {
let Some((column_idx, _)) = newer
.schema()
.column_with_name(FILESYSTEM_METADATA_COLUMN)
.map(|(i, f)| (i, f.clone()))
else {
return Ok(newer.clone());
};
let newer_column = newer.column(column_idx);
let older_column = older.column_by_name(FILESYSTEM_METADATA_COLUMN);
let mut merged: BTreeMap<String, FileInfo> = match older_column {
Some(column) => read_file_map(column)?,
None => BTreeMap::new(),
};
for (name, new_info) in read_file_map(newer_column)? {
match merged.get(&name) {
Some(old_info) => {
if new_info.is_deleted {
if old_info.is_deleted {
merged.insert(name, new_info);
} else {
merged.remove(&name);
}
} else {
merged.insert(
name,
FileInfo {
size: old_info.size.max(new_info.size),
is_deleted: false,
},
);
}
}
None => {
merged.insert(name, new_info);
}
}
}
let folded = build_file_map(newer_column.data_type(), &merged)?;
let mut columns = newer.columns().to_vec();
columns[column_idx] = folded;
RecordBatch::try_new(newer.schema(), columns).map_err(CoreError::ArrowError)
}
fn read_file_map(column: &ArrayRef) -> Result<BTreeMap<String, FileInfo>> {
let map = column.as_map_opt().ok_or_else(|| {
CoreError::Unsupported(format!(
"'{FILESYSTEM_METADATA_COLUMN}' must be a Map column, found {:?}.",
column.data_type()
))
})?;
let mut entries = BTreeMap::new();
if map.is_empty() || map.is_null(0) {
return Ok(entries);
}
let row = map.value(0);
let keys = row.column(0).as_string_opt::<i32>().ok_or_else(|| {
CoreError::Unsupported(format!(
"'{FILESYSTEM_METADATA_COLUMN}' keys must be Utf8 file names."
))
})?;
let values = row.column(1).as_struct_opt().ok_or_else(|| {
CoreError::Unsupported("'{FILESYSTEM_METADATA_COLUMN}' values must be structs.".to_string())
})?;
let sizes = struct_field(values, FILE_SIZE_FIELD)?
.as_primitive_opt::<arrow_array::types::Int64Type>()
.ok_or_else(|| {
CoreError::Unsupported(format!("'{FILE_SIZE_FIELD}' must be Int64.").to_string())
})?
.clone();
let deleted = struct_field(values, FILE_IS_DELETED_FIELD)?
.as_boolean_opt()
.ok_or_else(|| {
CoreError::Unsupported(
format!("'{FILE_IS_DELETED_FIELD}' must be Boolean.").to_string(),
)
})?
.clone();
for i in 0..keys.len() {
if keys.is_null(i) {
continue;
}
entries.insert(
keys.value(i).to_string(),
FileInfo {
size: if sizes.is_null(i) { 0 } else { sizes.value(i) },
is_deleted: !deleted.is_null(i) && deleted.value(i),
},
);
}
Ok(entries)
}
fn struct_field<'a>(values: &'a StructArray, name: &str) -> Result<&'a ArrayRef> {
values.column_by_name(name).ok_or_else(|| {
CoreError::Unsupported(format!(
"A '{FILESYSTEM_METADATA_COLUMN}' value must carry a '{name}' field."
))
})
}
fn build_file_map(data_type: &DataType, entries: &BTreeMap<String, FileInfo>) -> Result<ArrayRef> {
let DataType::Map(entries_field, sorted) = data_type else {
return Err(CoreError::Unsupported(format!(
"'{FILESYSTEM_METADATA_COLUMN}' must be a Map column, found {data_type:?}."
)));
};
let DataType::Struct(entry_fields) = entries_field.data_type() else {
return Err(CoreError::Unsupported(
"A Map column's entries must be a struct.".to_string(),
));
};
let key_field = entry_fields[0].clone();
let value_field = entry_fields[1].clone();
let DataType::Struct(value_fields) = value_field.data_type() else {
return Err(CoreError::Unsupported(format!(
"'{FILESYSTEM_METADATA_COLUMN}' values must be structs."
)));
};
let names: Vec<&str> = entries.keys().map(String::as_str).collect();
let sizes: Vec<i64> = entries.values().map(|f| f.size).collect();
let deleted: Vec<bool> = entries.values().map(|f| f.is_deleted).collect();
let mut value_columns: Vec<ArrayRef> = Vec::with_capacity(value_fields.len());
for field in value_fields.iter() {
match field.name().as_str() {
FILE_SIZE_FIELD => value_columns.push(Arc::new(Int64Array::from(sizes.clone()))),
FILE_IS_DELETED_FIELD => {
value_columns.push(Arc::new(BooleanArray::from(deleted.clone())))
}
other => {
return Err(CoreError::Unsupported(format!(
"Unexpected '{FILESYSTEM_METADATA_COLUMN}' value field '{other}'."
)));
}
}
}
let values = StructArray::try_new(value_fields.clone(), value_columns, None)
.map_err(CoreError::ArrowError)?;
let keys = arrow_array::StringArray::from(names);
let entry_struct = StructArray::try_new(
Fields::from(vec![key_field, value_field]),
vec![Arc::new(keys) as ArrayRef, Arc::new(values) as ArrayRef],
None,
)
.map_err(CoreError::ArrowError)?;
let offsets = OffsetBuffer::new(vec![0i32, entries.len() as i32].into());
let map = MapArray::try_new(
Arc::new(Field::clone(entries_field)),
offsets,
entry_struct,
None,
*sorted,
)
.map_err(CoreError::ArrowError)?;
Ok(Arc::new(map))
}
fn fold_column_stats(older: &RecordBatch, newer: &RecordBatch) -> Result<RecordBatch> {
let Some((column_idx, _)) = newer
.schema()
.column_with_name(COLUMN_STATS_COLUMN)
.map(|(i, f)| (i, f.clone()))
else {
return Ok(newer.clone());
};
let newer_stats = stats_struct(newer.column(column_idx))?;
let Some(older_column) = older.column_by_name(COLUMN_STATS_COLUMN) else {
return Ok(newer.clone());
};
let older_stats = stats_struct(older_column)?;
let (Some(newer_stats), Some(older_stats)) = (newer_stats, older_stats) else {
return Ok(newer.clone());
};
if flag(&newer_stats, FILE_IS_DELETED_FIELD) || flag(&older_stats, FILE_IS_DELETED_FIELD) {
return Ok(newer.clone());
}
if flag(&newer_stats, IS_TIGHT_BOUND_FIELD) {
return Ok(newer.clone());
}
let fields = match newer_stats.data_type() {
DataType::Struct(fields) => fields.clone(),
other => {
return Err(CoreError::Unsupported(format!(
"'{COLUMN_STATS_COLUMN}' must be a struct, found {other:?}."
)));
}
};
let mut columns: Vec<ArrayRef> = Vec::with_capacity(fields.len());
for (idx, field) in fields.iter().enumerate() {
let name = field.name().as_str();
let column = if name == MIN_VALUE_FIELD {
pick_bound(&older_stats, &newer_stats, name, std::cmp::Ordering::Less)?
} else if name == MAX_VALUE_FIELD {
pick_bound(
&older_stats,
&newer_stats,
name,
std::cmp::Ordering::Greater,
)?
} else if COUNTER_FIELDS.contains(&name) {
sum_counter(&older_stats, &newer_stats, name)?
} else {
newer_stats.column(idx).clone()
};
columns.push(column);
}
let folded = StructArray::try_new(fields, columns, newer_stats.nulls().cloned())
.map_err(CoreError::ArrowError)?;
let mut batch_columns = newer.columns().to_vec();
batch_columns[column_idx] = Arc::new(folded);
RecordBatch::try_new(newer.schema(), batch_columns).map_err(CoreError::ArrowError)
}
fn stats_struct(column: &ArrayRef) -> Result<Option<StructArray>> {
let stats = column.as_struct_opt().ok_or_else(|| {
CoreError::Unsupported(format!(
"'{COLUMN_STATS_COLUMN}' must be a struct column, found {:?}.",
column.data_type()
))
})?;
if stats.is_empty() || stats.is_null(0) {
return Ok(None);
}
Ok(Some(stats.clone()))
}
fn flag(stats: &StructArray, name: &str) -> bool {
stats
.column_by_name(name)
.and_then(|c| c.as_boolean_opt().cloned())
.is_some_and(|values| !values.is_null(0) && values.value(0))
}
fn sum_counter(older: &StructArray, newer: &StructArray, name: &str) -> Result<ArrayRef> {
let read = |stats: &StructArray| -> Option<i64> {
stats
.column_by_name(name)?
.as_primitive_opt::<arrow_array::types::Int64Type>()
.filter(|values| !values.is_null(0))
.map(|values| values.value(0))
};
let summed = match (read(older), read(newer)) {
(None, None) => None,
(a, b) => Some(a.unwrap_or(0).saturating_add(b.unwrap_or(0))),
};
Ok(Arc::new(Int64Array::from(vec![summed])))
}
fn pick_bound(
older: &StructArray,
newer: &StructArray,
name: &str,
want: std::cmp::Ordering,
) -> Result<ArrayRef> {
let (Some(older_bound), Some(newer_bound)) =
(older.column_by_name(name), newer.column_by_name(name))
else {
return Err(CoreError::Unsupported(format!(
"'{COLUMN_STATS_COLUMN}' must carry a '{name}' field to be merged."
)));
};
let older_union = bound_union(older_bound, name)?;
let newer_union = bound_union(newer_bound, name)?;
let older_value = wrapped_value(older_union, name)?;
let newer_value = wrapped_value(newer_union, name)?;
match (older_value, newer_value) {
(None, _) => Ok(newer_bound.clone()),
(_, None) => Ok(older_bound.clone()),
(Some((older_values, older_idx)), Some((newer_values, newer_idx))) => {
if older_values.data_type() != newer_values.data_type() {
return Err(CoreError::Unsupported(format!(
"'{name}' changes type between the two records being merged ({:?} and {:?}); this crate does not promote column statistics across types.",
older_values.data_type(),
newer_values.data_type()
)));
}
let compare = arrow_ord::ord::make_comparator(
older_values.as_ref(),
newer_values.as_ref(),
arrow_schema::SortOptions::default(),
)
.map_err(CoreError::ArrowError)?;
if compare(older_idx, newer_idx) == want {
Ok(older_bound.clone())
} else {
Ok(newer_bound.clone())
}
}
}
}
fn bound_union<'a>(bound: &'a ArrayRef, name: &str) -> Result<&'a arrow_array::UnionArray> {
bound
.as_any()
.downcast_ref::<arrow_array::UnionArray>()
.ok_or_else(|| {
CoreError::Unsupported(format!(
"'{name}' must be a union of typed wrappers, found {:?}.",
bound.data_type()
))
})
}
fn wrapped_value(bound: &arrow_array::UnionArray, name: &str) -> Result<Option<(ArrayRef, usize)>> {
if bound.is_empty() {
return Ok(None);
}
let type_id = bound.type_id(0);
let offset = bound.value_offset(0);
let child = bound.child(type_id);
if matches!(child.data_type(), DataType::Null) || child.is_null(offset) {
return Ok(None);
}
let wrapper = child.as_struct_opt().ok_or_else(|| {
CoreError::Unsupported(format!(
"'{name}' union branches must wrap their value in a struct, found {:?}.",
child.data_type()
))
})?;
let values = wrapper.column_by_name(WRAPPED_VALUE_FIELD).ok_or_else(|| {
CoreError::Unsupported(format!(
"'{name}' union branches must carry a '{WRAPPED_VALUE_FIELD}' field."
))
})?;
if values.is_null(offset) {
return Ok(None);
}
Ok(Some((values.clone(), offset)))
}
#[cfg(test)]
mod tests {
use super::*;
use arrow_array::{Int32Array, StringArray};
use arrow_schema::Schema;
fn map_type() -> DataType {
let value = DataType::Struct(Fields::from(vec![
Field::new(FILE_SIZE_FIELD, DataType::Int64, false),
Field::new(FILE_IS_DELETED_FIELD, DataType::Boolean, false),
]));
DataType::Map(
Arc::new(Field::new(
"entries",
DataType::Struct(Fields::from(vec![
Field::new("key", DataType::Utf8, false),
Field::new("value", value, false),
])),
false,
)),
false,
)
}
fn record(record_type: i32, entries: &[(&str, i64, bool)]) -> RecordBatch {
let map_type = map_type();
let map = build_file_map(
&map_type,
&entries
.iter()
.map(|(name, size, is_deleted)| {
(
(*name).to_string(),
FileInfo {
size: *size,
is_deleted: *is_deleted,
},
)
})
.collect(),
)
.unwrap();
let schema = Arc::new(Schema::new(vec![
Field::new("key", DataType::Utf8, false),
Field::new(TYPE_COLUMN, DataType::Int32, false),
Field::new(FILESYSTEM_METADATA_COLUMN, map_type, false),
]));
RecordBatch::try_new(
schema,
vec![
Arc::new(StringArray::from(vec!["city=chennai"])),
Arc::new(Int32Array::from(vec![record_type])),
map,
],
)
.unwrap()
}
fn folded(
older: &[(&str, i64, bool)],
newer: &[(&str, i64, bool)],
) -> Vec<(String, i64, bool)> {
let merged =
fold_filesystem_metadata(&record(TYPE_FILES, older), &record(TYPE_FILES, newer))
.unwrap();
read_file_map(merged.column_by_name(FILESYSTEM_METADATA_COLUMN).unwrap())
.unwrap()
.into_iter()
.map(|(name, info)| (name, info.size, info.is_deleted))
.collect()
}
#[test]
fn a_file_on_both_sides_keeps_the_larger_size() {
assert_eq!(
folded(&[("a.parquet", 900, false)], &[("a.parquet", 100, false)]),
vec![("a.parquet".to_string(), 900, false)]
);
assert_eq!(
folded(&[("a.parquet", 100, false)], &[("a.parquet", 900, false)]),
vec![("a.parquet".to_string(), 900, false)]
);
}
#[test]
fn a_tombstone_cancels_a_live_entry() {
assert_eq!(
folded(&[("a.parquet", 900, false)], &[("a.parquet", 0, true)]),
vec![]
);
}
#[test]
fn a_tombstone_over_a_tombstone_keeps_the_newer() {
assert_eq!(
folded(&[("a.parquet", 1, true)], &[("a.parquet", 2, true)]),
vec![("a.parquet".to_string(), 2, true)]
);
}
#[test]
fn a_live_entry_over_a_tombstone_revives_the_file() {
assert_eq!(
folded(&[("a.parquet", 5, true)], &[("a.parquet", 900, false)]),
vec![("a.parquet".to_string(), 900, false)]
);
}
#[test]
fn entries_on_one_side_only_pass_through() {
assert_eq!(
folded(
&[("old.parquet", 10, false)],
&[("new.parquet", 20, false), ("gone.parquet", 0, true)]
),
vec![
("gone.parquet".to_string(), 0, true),
("new.parquet".to_string(), 20, false),
("old.parquet".to_string(), 10, false),
]
);
}
#[test]
fn selection_types_take_the_newer_record() {
for record_type in [
4,
5,
7,
] {
let older = BufferedRecord::new_data(
"k".to_string(),
record(record_type, &[("old.parquet", 1, false)]),
None,
);
let newer = BufferedRecord::new_data(
"k".to_string(),
record(record_type, &[("new.parquet", 2, false)]),
None,
);
let merged = MetadataPayloadMerger.final_merge(&older, &newer).unwrap();
let batch = merged.payload.get_record().unwrap();
assert_eq!(
read_file_map(batch.column_by_name(FILESYSTEM_METADATA_COLUMN).unwrap())
.unwrap()
.keys()
.cloned()
.collect::<Vec<_>>(),
vec!["new.parquet".to_string()],
"partition type {record_type} must take the newer record whole"
);
}
}
fn stats_union_fields() -> arrow_schema::UnionFields {
arrow_schema::UnionFields::try_new(
vec![0i8, 3, 7],
vec![
Field::new("null", DataType::Null, true),
Field::new(
"LongWrapper",
DataType::Struct(Fields::from(vec![Field::new(
WRAPPED_VALUE_FIELD,
DataType::Int64,
true,
)])),
false,
),
Field::new(
"StringWrapper",
DataType::Struct(Fields::from(vec![Field::new(
WRAPPED_VALUE_FIELD,
DataType::Utf8,
true,
)])),
false,
),
],
)
.unwrap()
}
fn bound(value: Option<Bound>) -> ArrayRef {
let (type_id, long, string) = match value {
None => (0i8, None, None),
Some(Bound::Long(v)) => (3, Some(v), None),
Some(Bound::Text(v)) => (7, None, Some(v)),
};
let long_child: ArrayRef = Arc::new(
StructArray::try_new(
Fields::from(vec![Field::new(WRAPPED_VALUE_FIELD, DataType::Int64, true)]),
vec![Arc::new(Int64Array::from(long.map_or(vec![], |v| vec![v])))],
None,
)
.unwrap(),
);
let string_child: ArrayRef = Arc::new(
StructArray::try_new(
Fields::from(vec![Field::new(WRAPPED_VALUE_FIELD, DataType::Utf8, true)]),
vec![Arc::new(StringArray::from(
string.map_or(vec![], |v| vec![v]),
))],
None,
)
.unwrap(),
);
let null_child: ArrayRef = Arc::new(arrow_array::NullArray::new(if type_id == 0 {
1
} else {
0
}));
Arc::new(
arrow_array::UnionArray::try_new(
stats_union_fields(),
vec![type_id].into(),
Some(vec![0i32].into()),
vec![null_child, long_child, string_child],
)
.unwrap(),
)
}
#[derive(Clone, Copy)]
enum Bound {
Long(i64),
Text(&'static str),
}
fn stats_record(
min: Option<Bound>,
max: Option<Bound>,
counters: [i64; 4],
is_deleted: bool,
is_tight_bound: bool,
) -> RecordBatch {
let union_type = DataType::Union(stats_union_fields(), arrow_schema::UnionMode::Dense);
let mut fields = vec![
Field::new(MIN_VALUE_FIELD, union_type.clone(), true),
Field::new(MAX_VALUE_FIELD, union_type, true),
];
let mut columns: Vec<ArrayRef> = vec![bound(min), bound(max)];
for (name, value) in COUNTER_FIELDS.iter().zip(counters) {
fields.push(Field::new(*name, DataType::Int64, true));
columns.push(Arc::new(Int64Array::from(vec![value])));
}
fields.push(Field::new(FILE_IS_DELETED_FIELD, DataType::Boolean, false));
columns.push(Arc::new(BooleanArray::from(vec![is_deleted])));
fields.push(Field::new(IS_TIGHT_BOUND_FIELD, DataType::Boolean, false));
columns.push(Arc::new(BooleanArray::from(vec![is_tight_bound])));
let stats = StructArray::try_new(Fields::from(fields), columns, None).unwrap();
let schema = Arc::new(Schema::new(vec![
Field::new(TYPE_COLUMN, DataType::Int32, false),
Field::new(COLUMN_STATS_COLUMN, stats.data_type().clone(), true),
]));
RecordBatch::try_new(
schema,
vec![
Arc::new(Int32Array::from(vec![TYPE_COLUMN_STATS])),
Arc::new(stats),
],
)
.unwrap()
}
fn folded_stats(
older: &RecordBatch,
newer: &RecordBatch,
) -> (Option<String>, Option<String>, Vec<i64>) {
let merged = fold_column_stats(older, newer).unwrap();
let stats = merged
.column_by_name(COLUMN_STATS_COLUMN)
.unwrap()
.as_struct()
.clone();
let render = |name: &str| -> Option<String> {
let union = bound_union(stats.column_by_name(name).unwrap(), name).unwrap();
wrapped_value(union, name).unwrap().map(|(values, idx)| {
arrow_cast::display::array_value_to_string(&values, idx).unwrap()
})
};
let counters = COUNTER_FIELDS
.iter()
.map(|name| {
stats
.column_by_name(name)
.unwrap()
.as_primitive::<arrow_array::types::Int64Type>()
.value(0)
})
.collect();
(render(MIN_VALUE_FIELD), render(MAX_VALUE_FIELD), counters)
}
#[test]
fn stats_bounds_widen_and_counters_sum() {
let older = stats_record(
Some(Bound::Long(5)),
Some(Bound::Long(50)),
[1, 2, 3, 4],
false,
false,
);
let newer = stats_record(
Some(Bound::Long(9)),
Some(Bound::Long(90)),
[10, 20, 30, 40],
false,
false,
);
assert_eq!(
folded_stats(&older, &newer),
(
Some("5".to_string()),
Some("90".to_string()),
vec![11, 22, 33, 44]
)
);
assert_eq!(
folded_stats(&newer, &older),
(
Some("5".to_string()),
Some("90".to_string()),
vec![11, 22, 33, 44]
)
);
}
#[test]
fn stats_bounds_compare_within_the_union_branch() {
let older = stats_record(
Some(Bound::Text("m")),
Some(Bound::Text("m")),
[0; 4],
false,
false,
);
let newer = stats_record(
Some(Bound::Text("b")),
Some(Bound::Text("z")),
[0; 4],
false,
false,
);
let (min, max, _) = folded_stats(&older, &newer);
assert_eq!((min, max), (Some("b".to_string()), Some("z".to_string())));
}
#[test]
fn a_null_bound_loses_to_a_value() {
let older = stats_record(None, None, [0; 4], false, false);
let newer = stats_record(
Some(Bound::Long(7)),
Some(Bound::Long(7)),
[0; 4],
false,
false,
);
assert_eq!(
folded_stats(&older, &newer).0,
Some("7".to_string()),
"a null older bound must not erase the newer value"
);
assert_eq!(
folded_stats(&newer, &older).0,
Some("7".to_string()),
"a null newer bound must not erase the older value"
);
}
#[test]
fn a_tight_bound_newer_record_supersedes() {
let older = stats_record(
Some(Bound::Long(1)),
Some(Bound::Long(100)),
[9, 9, 9, 9],
false,
false,
);
let newer = stats_record(
Some(Bound::Long(5)),
Some(Bound::Long(50)),
[1, 1, 1, 1],
false,
true,
);
assert_eq!(
folded_stats(&older, &newer),
(
Some("5".to_string()),
Some("50".to_string()),
vec![1, 1, 1, 1]
)
);
}
#[test]
fn a_deleted_stats_record_on_either_side_overwrites() {
let live = stats_record(
Some(Bound::Long(1)),
Some(Bound::Long(100)),
[9, 9, 9, 9],
false,
false,
);
let deleted = stats_record(
Some(Bound::Long(5)),
Some(Bound::Long(50)),
[1, 1, 1, 1],
true,
false,
);
for (older, newer) in [(&live, &deleted), (&deleted, &live)] {
let (_, _, counters) = folded_stats(older, newer);
assert_eq!(
counters,
COUNTER_FIELDS
.iter()
.enumerate()
.map(|(i, _)| {
newer
.column_by_name(COLUMN_STATS_COLUMN)
.unwrap()
.as_struct()
.column(2 + i)
.as_primitive::<arrow_array::types::Int64Type>()
.value(0)
})
.collect::<Vec<_>>(),
"a tombstone must overwrite rather than sum"
);
}
}
#[test]
fn a_bound_that_changes_type_is_refused() {
let older = stats_record(
Some(Bound::Long(5)),
Some(Bound::Long(5)),
[0; 4],
false,
false,
);
let newer = stats_record(
Some(Bound::Text("b")),
Some(Bound::Text("b")),
[0; 4],
false,
false,
);
let err = fold_column_stats(&older, &newer).expect_err("a type change must be refused");
assert!(
format!("{err:?}").contains("changes type between the two records"),
"the error must name the type change, got: {err:?}"
);
}
#[test]
fn only_the_metadata_payload_selects_a_merger() {
let with = |pairs: &[(&str, &str)]| -> Option<CustomMerger> {
resolve_custom_merger(
&pairs
.iter()
.map(|(k, v)| ((*k).to_string(), (*v).to_string()))
.collect(),
)
};
assert_eq!(
with(&[
("hoodie.compaction.payload.class", METADATA_PAYLOAD_CLASS),
("hoodie.record.merge.strategy.id", PAYLOAD_BASED_STRATEGY_ID),
]),
Some(CustomMerger::MetadataPayload)
);
assert_eq!(
with(&[
("hoodie.compaction.payload.class", "com.example.Payload"),
("hoodie.record.merge.strategy.id", PAYLOAD_BASED_STRATEGY_ID),
]),
None
);
assert_eq!(
with(&[
("hoodie.compaction.payload.class", METADATA_PAYLOAD_CLASS),
(
"hoodie.record.merge.strategy.id",
"eeb8d96f-b1e4-49fd-bbf8-28ac514178e5"
),
]),
None
);
assert_eq!(
with(&[("hoodie.compaction.payload.class", METADATA_PAYLOAD_CLASS)]),
Some(CustomMerger::MetadataPayload)
);
assert_eq!(with(&[]), None);
assert_eq!(
with(&[
("hoodie.compaction.payload.class", METADATA_PAYLOAD_CLASS),
(
"hoodie.compaction.record.merger.strategy",
"eeb8d96f-b1e4-49fd-bbf8-28ac514178e5"
),
]),
None
);
for key in [
"hoodie.datasource.write.payload.class",
"hoodie.table.legacy.payload.class",
] {
assert_eq!(
with(&[(key, METADATA_PAYLOAD_CLASS)]),
Some(CustomMerger::MetadataPayload),
"the payload class must be read from '{key}' too"
);
}
}
}