use std::collections::HashMap;
use crate::table_changes::scan_file::{TableChangesFileAction, TableChangesScanFile};
use crate::{DeltaResult, Error};
struct Side {
is_add: bool,
file: TableChangesScanFile,
}
pub(super) fn collapse_net_changes(
actions: Vec<TableChangesFileAction>,
) -> DeltaResult<Vec<TableChangesFileAction>> {
let mut by_path: HashMap<String, Vec<Side>> = HashMap::new();
for action in actions {
if let Some(add) = action.add {
by_path.entry(add.path.clone()).or_default().push(Side {
is_add: true,
file: add,
});
}
if let Some(remove) = action.remove {
by_path.entry(remove.path.clone()).or_default().push(Side {
is_add: false,
file: remove,
});
}
}
let mut out: Vec<(String, TableChangesFileAction)> = Vec::with_capacity(by_path.len());
for (path, sides) in by_path {
let key = |s: &&Side| (s.file.commit_version, s.is_add);
let (Some(earliest), Some(latest)) =
(sides.iter().min_by_key(key), sides.iter().max_by_key(key))
else {
return Err(Error::internal_error(format!(
"net-changes collapse produced an empty side slot for path {path}"
)));
};
let first_remove = (!earliest.is_add).then(|| earliest.file.clone());
let last_add = latest.is_add.then(|| latest.file.clone());
let action = match (first_remove, last_add) {
(Some(remove), None) => TableChangesFileAction {
add: None,
remove: Some(remove),
},
(None, Some(add)) => TableChangesFileAction {
add: Some(add),
remove: None,
},
(Some(remove), Some(add)) if remove.deletion_vector != add.deletion_vector => {
TableChangesFileAction {
add: Some(add),
remove: Some(remove),
}
}
(Some(_), Some(_)) | (None, None) => continue,
};
out.push((path, action));
}
out.sort_unstable_by(|(a, _), (b, _)| a.cmp(b));
Ok(out.into_iter().map(|(_, action)| action).collect())
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use rstest::rstest;
use crate::actions::deletion_vector::{DeletionVectorDescriptor, DeletionVectorStorageType};
use crate::table_changes::net_changes::collapse_net_changes;
use crate::table_changes::scan_file::{TableChangesFileAction, TableChangesScanFile};
fn test_dv(id: &str) -> DeletionVectorDescriptor {
DeletionVectorDescriptor {
storage_type: DeletionVectorStorageType::PersistedRelative,
path_or_inline_dv: id.to_string(),
offset: Some(1),
size_in_bytes: 8,
cardinality: 1,
}
}
fn side(
path: &str,
commit_version: i64,
deletion_vector: Option<DeletionVectorDescriptor>,
) -> TableChangesScanFile {
TableChangesScanFile {
path: path.to_string(),
deletion_vector,
partition_values: HashMap::new(),
size: Some(10),
commit_version,
commit_timestamp: 0,
base_row_id: Some(0),
default_row_commit_version: Some(1),
}
}
fn insert(
path: &str,
commit_version: i64,
dv: Option<DeletionVectorDescriptor>,
) -> TableChangesFileAction {
TableChangesFileAction {
add: Some(side(path, commit_version, dv)),
remove: None,
}
}
fn delete(
path: &str,
commit_version: i64,
dv: Option<DeletionVectorDescriptor>,
) -> TableChangesFileAction {
TableChangesFileAction {
add: None,
remove: Some(side(path, commit_version, dv)),
}
}
fn paired_dv_update(
path: &str,
commit_version: i64,
add_dv: Option<DeletionVectorDescriptor>,
removed_dv: Option<DeletionVectorDescriptor>,
) -> TableChangesFileAction {
TableChangesFileAction {
add: Some(side(path, commit_version, add_dv)),
remove: Some(side(path, commit_version, removed_dv)),
}
}
type SideSummary = Option<(i64, Option<DeletionVectorDescriptor>)>;
fn summarize(a: &TableChangesFileAction) -> (String, SideSummary, SideSummary) {
let s = |side: &Option<TableChangesScanFile>| {
side.as_ref()
.map(|f| (f.commit_version, f.deletion_vector.clone()))
};
let path = a
.add
.as_ref()
.or(a.remove.as_ref())
.map(|f| f.path.clone())
.expect("an action always has at least one side");
(path, s(&a.add), s(&a.remove))
}
#[test]
fn net_changes_collapse_reduces_each_path_to_its_net_effect() {
let dv_a = test_dv("aaaaaaaaaaaaaaaaaaaa");
let dv_b = test_dv("bbbbbbbbbbbbbbbbbbbb");
let actions = vec![
insert("insert.parquet", 1, None),
delete("delete.parquet", 1, Some(dv_a.clone())),
delete("update.parquet", 1, Some(dv_a.clone())),
insert("update.parquet", 2, Some(dv_b.clone())),
delete("carryover.parquet", 1, Some(dv_a.clone())),
insert("carryover.parquet", 2, Some(dv_a.clone())),
insert("churn.parquet", 1, None),
delete("churn.parquet", 2, None),
insert("reinsert.parquet", 1, None),
delete("reinsert.parquet", 2, None),
insert("reinsert.parquet", 3, None),
];
let out = collapse_net_changes(actions).unwrap();
assert_eq!(
out.iter().map(summarize).collect::<Vec<_>>(),
vec![
(
"delete.parquet".to_string(),
None,
Some((1, Some(dv_a.clone())))
),
("insert.parquet".to_string(), Some((1, None)), None),
("reinsert.parquet".to_string(), Some((3, None)), None),
(
"update.parquet".to_string(),
Some((2, Some(dv_b))),
Some((1, Some(dv_a)))
),
]
);
}
#[rstest]
#[case::same_commit_dv_update(
vec![paired_dv_update("update.parquet", 3, Some(test_dv("b")), Some(test_dv("a")))],
vec![(
"update.parquet".to_string(),
Some((3, Some(test_dv("b")))),
Some((3, Some(test_dv("a")))),
)]
)]
#[case::multi_update_chain(
vec![
paired_dv_update("f.parquet", 1, Some(test_dv("b")), Some(test_dv("a"))),
paired_dv_update("f.parquet", 2, Some(test_dv("c")), Some(test_dv("b"))),
],
vec![(
"f.parquet".to_string(),
Some((2, Some(test_dv("c")))),
Some((1, Some(test_dv("a")))),
)]
)]
#[case::multi_commit_revert(
vec![
paired_dv_update("f.parquet", 1, Some(test_dv("b")), Some(test_dv("a"))),
paired_dv_update("f.parquet", 2, Some(test_dv("a")), Some(test_dv("b"))),
],
vec![]
)]
#[case::copy_on_write(
vec![delete("old.parquet", 1, None), insert("new.parquet", 1, None)],
vec![
("new.parquet".to_string(), Some((1, None)), None),
("old.parquet".to_string(), None, Some((1, None))),
]
)]
fn net_changes_collapse_reduces_paths_to_their_boundaries(
#[case] actions: Vec<TableChangesFileAction>,
#[case] expected: Vec<(String, SideSummary, SideSummary)>,
) {
let out = collapse_net_changes(actions).unwrap();
assert_eq!(out.iter().map(summarize).collect::<Vec<_>>(), expected);
}
#[test]
fn net_changes_collapse_preserves_row_tracking_fields() {
let dv_a = test_dv("aaaaaaaaaaaaaaaaaaaa");
let dv_b = test_dv("bbbbbbbbbbbbbbbbbbbb");
let actions = vec![
TableChangesFileAction {
add: None,
remove: Some(TableChangesScanFile {
base_row_id: Some(100),
default_row_commit_version: Some(1),
..side("f.parquet", 1, Some(dv_a))
}),
},
TableChangesFileAction {
add: Some(TableChangesScanFile {
base_row_id: Some(200),
default_row_commit_version: Some(2),
..side("f.parquet", 2, Some(dv_b))
}),
remove: None,
},
];
let out = collapse_net_changes(actions).unwrap();
assert_eq!(out.len(), 1, "expected one grouped net update: {out:?}");
let remove = out[0].remove.as_ref().expect("remove side");
let add = out[0].add.as_ref().expect("add side");
assert_eq!(remove.base_row_id, Some(100));
assert_eq!(remove.default_row_commit_version, Some(1));
assert_eq!(add.base_row_id, Some(200));
assert_eq!(add.default_row_commit_version, Some(2));
}
}