mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 12:31:36 +00:00
fix(table-catalog): preserve data sequences during file rewrites
Merge the reviewed fix from pull request #8104.
This commit is contained in:
@@ -4761,12 +4761,14 @@ where
|
||||
"manifest changed entries must belong to the committed snapshot"
|
||||
));
|
||||
}
|
||||
// Rewritten files may preserve an older data sequence number.
|
||||
if status == 1
|
||||
&& manifest.location.sequence_number.is_some_and(|sequence_number| {
|
||||
reference.sequence_number != Some(sequence_number) || reference.file_sequence_number != Some(sequence_number)
|
||||
})
|
||||
&& manifest
|
||||
.location
|
||||
.sequence_number
|
||||
.is_some_and(|sequence_number| reference.file_sequence_number != Some(sequence_number))
|
||||
{
|
||||
return Err(s3_error!(InvalidRequest, "added manifest entry sequence must match its manifest"));
|
||||
return Err(s3_error!(InvalidRequest, "added manifest entry file sequence must match its manifest"));
|
||||
}
|
||||
if status == 2
|
||||
&& context
|
||||
|
||||
@@ -8129,6 +8129,118 @@ async fn row_level_conflict_rejects_changed_inherited_manifest_identity() {
|
||||
assert_eq!(unchanged.generation, current.generation);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn row_level_conflict_rewrite_data_sequence_numbers() {
|
||||
use crate::table_catalog::test_support::{
|
||||
manifest_avro_bytes_with_entry_sequences, manifest_list_avro_entries_with_min_sequences,
|
||||
};
|
||||
|
||||
let invalid_data_sequence = "Iceberg v2 manifest entry sequence number exceeds its manifest sequence";
|
||||
let invalid_file_sequence = "added manifest entry file sequence must match its manifest";
|
||||
// A rewrite preserves the data age while adding a new physical file.
|
||||
for (data_sequence, file_sequence, expected_error) in [
|
||||
(Some(1), None, None),
|
||||
(Some(1), Some(2), None),
|
||||
(None, None, None),
|
||||
(Some(0), Some(2), None),
|
||||
(Some(2), Some(2), None),
|
||||
(Some(1), Some(1), Some(invalid_file_sequence)),
|
||||
(Some(-1), Some(2), Some(invalid_data_sequence)),
|
||||
(Some(3), Some(2), Some(invalid_data_sequence)),
|
||||
] {
|
||||
let store = TestTableCatalogStore::default();
|
||||
let backend = TestTableCatalogObjectBackend::content_addressed();
|
||||
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
|
||||
let created = create_standard_events_table(&store, &backend, &namespace).await;
|
||||
let location = created.metadata["location"].as_str().expect("table location should exist");
|
||||
let old_file = format!("{location}/data/old.parquet");
|
||||
let old_list = format!("{location}/metadata/snap-10.avro");
|
||||
seed_test_snapshot_manifest(&backend, "warehouse", &old_list, 10, 1, &[(&old_file, 0, 1, 10, 1)]).await;
|
||||
let append = serde_json::from_value(serde_json::json!({
|
||||
"updates": [
|
||||
{"action": "add-snapshot", "snapshot": {
|
||||
"snapshot-id": 10, "sequence-number": 1, "timestamp-ms": 1234,
|
||||
"manifest-list": old_list, "summary": {"operation": "append"}
|
||||
}},
|
||||
{"action": "set-snapshot-ref", "ref-name": "main", "snapshot-id": 10, "type": "branch"}
|
||||
]
|
||||
}))
|
||||
.expect("append request should parse");
|
||||
commit_table_response(&store, &trusted_table_commit_backend(&backend), "warehouse", &namespace, "events", append)
|
||||
.await
|
||||
.expect("initial append should succeed");
|
||||
let before = store.load_table("warehouse", "analytics", "events").await.unwrap().unwrap();
|
||||
|
||||
let new_file = format!("{location}/data/rewritten.parquet");
|
||||
let manifest = format!("{location}/metadata/rewrite.avro");
|
||||
let manifest_list = format!("{location}/metadata/snap-11.avro");
|
||||
let manifest_bytes = manifest_avro_bytes_with_entry_sequences(&[
|
||||
(&old_file, 0, 2, 11, Some(1), Some(1)),
|
||||
(&new_file, 0, 1, 11, data_sequence, file_sequence),
|
||||
]);
|
||||
let list_bytes = manifest_list_avro_entries_with_min_sequences(&[(
|
||||
&manifest,
|
||||
manifest_bytes.len(),
|
||||
0,
|
||||
0,
|
||||
2,
|
||||
data_sequence.unwrap_or(2).clamp(0, 2),
|
||||
11,
|
||||
)]);
|
||||
backend
|
||||
.put_bytes("warehouse", &test_snapshot_object_key("warehouse", &manifest), manifest_bytes)
|
||||
.await;
|
||||
backend
|
||||
.put_bytes("warehouse", &test_snapshot_object_key("warehouse", &manifest_list), list_bytes)
|
||||
.await;
|
||||
backend
|
||||
.put_bytes("warehouse", &test_snapshot_object_key("warehouse", &new_file), b"data".to_vec())
|
||||
.await;
|
||||
let rewrite = serde_json::from_value(serde_json::json!({
|
||||
"requirements": [{"type": "assert-ref-snapshot-id", "ref": "main", "snapshot-id": 10}],
|
||||
"updates": [
|
||||
{"action": "add-snapshot", "snapshot": {
|
||||
"snapshot-id": 11, "parent-snapshot-id": 10, "sequence-number": 2,
|
||||
"timestamp-ms": 2234, "manifest-list": manifest_list,
|
||||
"summary": {"operation": "replace"}
|
||||
}},
|
||||
{"action": "set-snapshot-ref", "ref-name": "main", "snapshot-id": 11, "type": "branch"}
|
||||
]
|
||||
}))
|
||||
.expect("rewrite request should parse");
|
||||
let result = commit_table_response(
|
||||
&store,
|
||||
&trusted_table_commit_backend(&backend),
|
||||
"warehouse",
|
||||
&namespace,
|
||||
"events",
|
||||
rewrite,
|
||||
)
|
||||
.await;
|
||||
let after = store.load_table("warehouse", "analytics", "events").await.unwrap().unwrap();
|
||||
if expected_error.is_none() {
|
||||
let commit = result.unwrap_or_else(|err| panic!("data={data_sequence:?}, file={file_sequence:?}: {err}"));
|
||||
assert_eq!(commit.metadata["current-snapshot-id"], 11);
|
||||
assert_eq!(commit.metadata["last-sequence-number"], 2);
|
||||
assert_eq!(after.generation, before.generation + 1);
|
||||
let live_files = load_snapshot_live_files(&backend, "warehouse", &after, &commit.metadata, Some(11))
|
||||
.await
|
||||
.expect("committed snapshot should load");
|
||||
assert!(!live_files.data_files.contains_key(&old_file));
|
||||
let added = live_files.data_files.get(&new_file).expect("rewritten file should be live");
|
||||
assert_eq!(added.sequence_number, Some(data_sequence.unwrap_or(2)));
|
||||
assert_eq!(added.file_sequence_number, Some(2));
|
||||
} else {
|
||||
let error = result.expect_err("invalid sequence numbers must be rejected");
|
||||
assert_eq!(error.code(), &S3ErrorCode::InvalidRequest);
|
||||
assert_eq!(error.message(), expected_error);
|
||||
assert_eq!(after.metadata_location, before.metadata_location);
|
||||
assert_eq!(after.generation, before.generation);
|
||||
assert_eq!(after.version_token, before.version_token);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Table-driven coverage for stale or historical manifest sequence failures.
|
||||
#[tokio::test]
|
||||
async fn row_level_conflict_rejects_stale_or_historical_manifest_sequences() {
|
||||
@@ -8142,7 +8254,7 @@ async fn row_level_conflict_rejects_stale_or_historical_manifest_sequences() {
|
||||
"new manifest sequence must match the committed snapshot",
|
||||
),
|
||||
(
|
||||
"stale-added-entry-sequence",
|
||||
"stale-added-file-sequence",
|
||||
2,
|
||||
"11",
|
||||
11,
|
||||
|
||||
@@ -84,6 +84,16 @@ pub(crate) fn manifest_list_avro_entries_with_partition_specs(manifests: &[(&str
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_list_avro_entries_with_content(manifests: &[(&str, usize, i32, i32, i64, i64)]) -> Vec<u8> {
|
||||
let manifests = manifests
|
||||
.iter()
|
||||
.map(|(path, length, spec_id, content, sequence_number, snapshot_id)| {
|
||||
(*path, *length, *spec_id, *content, *sequence_number, *sequence_number, *snapshot_id)
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
manifest_list_avro_entries_with_min_sequences(&manifests)
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_list_avro_entries_with_min_sequences(manifests: &[(&str, usize, i32, i32, i64, i64, i64)]) -> Vec<u8> {
|
||||
let schema = apache_avro::Schema::parse_str(
|
||||
r#"
|
||||
{
|
||||
@@ -109,7 +119,9 @@ pub(crate) fn manifest_list_avro_entries_with_content(manifests: &[(&str, usize,
|
||||
)
|
||||
.expect("manifest list avro schema should parse");
|
||||
let mut writer = apache_avro::Writer::new(&schema, Vec::new()).expect("manifest list writer should initialize");
|
||||
for (manifest_path, manifest_length, partition_spec_id, content, sequence_number, snapshot_id) in manifests {
|
||||
for (manifest_path, manifest_length, partition_spec_id, content, sequence_number, min_sequence_number, snapshot_id) in
|
||||
manifests
|
||||
{
|
||||
writer
|
||||
.append_value(apache_avro::types::Value::Record(vec![
|
||||
(
|
||||
@@ -123,7 +135,7 @@ pub(crate) fn manifest_list_avro_entries_with_content(manifests: &[(&str, usize,
|
||||
("partition_spec_id".to_string(), apache_avro::types::Value::Int(*partition_spec_id)),
|
||||
("content".to_string(), apache_avro::types::Value::Int(*content)),
|
||||
("sequence_number".to_string(), apache_avro::types::Value::Long(*sequence_number)),
|
||||
("min_sequence_number".to_string(), apache_avro::types::Value::Long(*sequence_number)),
|
||||
("min_sequence_number".to_string(), apache_avro::types::Value::Long(*min_sequence_number)),
|
||||
("added_snapshot_id".to_string(), apache_avro::types::Value::Long(*snapshot_id)),
|
||||
("added_files_count".to_string(), apache_avro::types::Value::Int(1)),
|
||||
("existing_files_count".to_string(), apache_avro::types::Value::Int(0)),
|
||||
@@ -283,6 +295,18 @@ pub(crate) fn nullable_long(value: Option<i64>) -> apache_avro::types::Value {
|
||||
}
|
||||
|
||||
pub(crate) fn manifest_avro_bytes_with_nullable_sequences(files: &[(&str, i32, i32, i64, Option<i64>)]) -> Vec<u8> {
|
||||
let files = files
|
||||
.iter()
|
||||
.map(|(path, content, status, snapshot_id, sequence_number)| {
|
||||
(*path, *content, *status, *snapshot_id, *sequence_number, *sequence_number)
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
manifest_avro_bytes_with_entry_sequences(&files)
|
||||
}
|
||||
|
||||
pub(crate) type ManifestEntryWithSequences<'a> = (&'a str, i32, i32, i64, Option<i64>, Option<i64>);
|
||||
|
||||
pub(crate) fn manifest_avro_bytes_with_entry_sequences(files: &[ManifestEntryWithSequences<'_>]) -> Vec<u8> {
|
||||
let schema = apache_avro::Schema::parse_str(
|
||||
r#"
|
||||
{
|
||||
@@ -313,13 +337,13 @@ pub(crate) fn manifest_avro_bytes_with_nullable_sequences(files: &[(&str, i32, i
|
||||
)
|
||||
.expect("manifest avro schema should parse");
|
||||
let mut writer = apache_avro::Writer::new(&schema, Vec::new()).expect("manifest writer should initialize");
|
||||
for (file_path, content, status, snapshot_id, sequence_number) in files {
|
||||
for (file_path, content, status, snapshot_id, sequence_number, file_sequence_number) in files {
|
||||
writer
|
||||
.append_value(apache_avro::types::Value::Record(vec![
|
||||
("status".to_string(), apache_avro::types::Value::Int(*status)),
|
||||
("snapshot_id".to_string(), apache_avro::types::Value::Long(*snapshot_id)),
|
||||
("sequence_number".to_string(), nullable_long(*sequence_number)),
|
||||
("file_sequence_number".to_string(), nullable_long(*sequence_number)),
|
||||
("file_sequence_number".to_string(), nullable_long(*file_sequence_number)),
|
||||
(
|
||||
"data_file".to_string(),
|
||||
apache_avro::types::Value::Record(vec![
|
||||
|
||||
Reference in New Issue
Block a user