leaves12138 commented on code in PR #915:
URL: https://github.com/apache/paimon-rust/pull/915#discussion_r4070384634
##########
crates/paimon/src/table/table_commit.rs:
##########
@@ -2395,40 +2476,103 @@ impl TableCommit {
Ok(())
}
- fn check_deletion_vector_conflicts(
+ /// Validate DV references against the state this attempt will publish.
Index
+ /// merging separately verifies replacement identities and one DV per file.
+ async fn check_deletion_vector_references(
&self,
- latest_snapshot: Option<&Snapshot>,
+ latest_snapshot: &Option<Snapshot>,
+ data_entries: &[ManifestEntry],
index_entries: &[IndexManifestEntry],
- check_from_snapshot: Option<i64>,
+ merged_indexes: &[IndexManifestEntry],
) -> Result<()> {
- if !self.data_evolution_enabled {
- return Ok(());
- }
- let Some(check_from_snapshot) = check_from_snapshot else {
- return Ok(());
+ let file_key = |partition: &[u8], bucket: i32, name: &str| {
+ (
+ if is_empty_partition(partition) {
+ vec![]
+ } else {
+ partition.to_vec()
+ },
+ bucket,
+ name.to_string(),
+ )
};
- let has_deletion_vector_index_change = index_entries
+ let deleted_files = data_entries
.iter()
- .any(|entry| entry.index_file.index_type ==
DELETION_VECTORS_INDEX_TYPE);
- if !has_deletion_vector_index_change {
- return Ok(());
+ .filter(|entry| *entry.kind() == FileKind::Delete)
+ .map(|entry| file_key(entry.partition(), entry.bucket(),
&entry.file().file_name))
+ .collect::<HashSet<_>>();
+ let added_dvs = index_entries
+ .iter()
+ .filter(|entry| {
+ entry.kind == FileKind::Add
+ && entry.index_file.index_type ==
DELETION_VECTORS_INDEX_TYPE
+ })
+ .map(|entry| file_key(&entry.partition, entry.bucket,
&entry.index_file.file_name))
+ .collect::<HashSet<_>>();
+ let mut referenced_files = HashSet::new();
+ for entry in merged_indexes
+ .iter()
+ .filter(|entry| entry.index_file.index_type ==
DELETION_VECTORS_INDEX_TYPE)
+ {
+ let is_added = added_dvs.contains(&file_key(
+ &entry.partition,
+ entry.bucket,
+ &entry.index_file.file_name,
+ ));
+ if let Some(ranges) = &entry.index_file.deletion_vectors_ranges {
+ for name in ranges.keys() {
+ let key = file_key(&entry.partition, entry.bucket, name);
+ // Also forbid retaining a DV after deleting its data file.
+ if is_added || deleted_files.contains(&key) {
+ referenced_files.insert(key);
Review Comment:
[P2] Adapt the existing stale-DV integration fixture to the new invariant
Rejecting retained DVs that reference a deleted data file also changes the
setup used by
`crates/integrations/datafusion/tests/partition_count_pushdown.rs::test_removed_file_deletion_vector_does_not_reduce_count`.
That test deliberately removes a data file while retaining its DV through
`TableCommit`, then verifies that count pushdown ignores the stale vector. It
now fails at the commit `unwrap()` on line 718 with `Deletion vector references
missing data file`, before reaching the reader/count assertions; this is the
failure in the current `integration (datafusion)` CI job.
Please update that fixture alongside this validation change. If retaining
such metadata is intentionally no longer allowed through the public committer,
construct the legacy snapshot/index state directly for this reader regression
instead of using the newly restricted API. Keep the stale-DV count/read
coverage rather than simply dropping the test.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]