mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-07 05:49:12 +00:00
test(rust): cover concurrent appends during compaction (#3878)
## Summary - add a LanceDB core regression for compaction overlapping appends through separate table handles - verify concurrent commits preserve fragment ID order on an indexed table - run the follow-up compaction that exposed the original row-ID ordering failure and verify all rows remain ## Root cause Older Lance versions could reserve fragment IDs for compaction, allow concurrent appends to commit later IDs, and then commit the reserved compaction fragments at the end of the manifest. A later compaction could consequently receive row IDs out of order. Current Lance sorts fragments at the transaction boundary; this adds the missing LanceDB-level regression coverage for the Node-visible concurrency contract. ## Validation - `cargo fmt --all` - focused regression passed once with output and 20 repeated runs - `cargo test --quiet --features remote -p lancedb table::optimize::tests` (14 passed) - `cargo check --quiet --features remote --tests --examples` - `cargo clippy --quiet --features remote --tests --examples` - `cargo test --quiet --features remote --tests` (867 passed, 1 ignored) Fixes #1498 <!-- lance-gatekeeper-fix:v1 agent=93aaefb15507dca52d064e15388773d7 generation=1 --> Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
This commit is contained in:
@@ -304,6 +304,96 @@ mod tests {
|
||||
assert_eq!(all_values, expected);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_compact_with_concurrent_add() {
|
||||
const NUM_FRAGMENTS: usize = 5;
|
||||
const ROWS_PER_FRAGMENT: i32 = 300;
|
||||
|
||||
let tmpdir = tempfile::tempdir().unwrap();
|
||||
let conn = connect(tmpdir.path().to_str().unwrap())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
|
||||
let batch = RecordBatch::try_new(
|
||||
schema,
|
||||
vec![Arc::new(Int32Array::from_iter_values(0..ROWS_PER_FRAGMENT))],
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let table = conn
|
||||
.create_table("test_concurrent_compact", batch.clone())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
table
|
||||
.create_index(&["id"], Index::BTree(BTreeIndexBuilder::default()))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
for _ in 0..NUM_FRAGMENTS {
|
||||
table.add(batch.clone()).execute().await.unwrap();
|
||||
}
|
||||
|
||||
// Use separate handles so the two writes actually overlap, as they can
|
||||
// when different Node connections operate on the same S3 table.
|
||||
let compact_table = conn
|
||||
.open_table("test_concurrent_compact")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let append_table = conn
|
||||
.open_table("test_concurrent_compact")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let compact_task = tokio::spawn(async move {
|
||||
compact_table
|
||||
.optimize(OptimizeAction::Compact {
|
||||
options: CompactionOptions {
|
||||
target_rows_per_fragment: 1_000,
|
||||
..Default::default()
|
||||
},
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
});
|
||||
tokio::task::yield_now().await;
|
||||
for _ in 0..NUM_FRAGMENTS {
|
||||
append_table.add(batch.clone()).execute().await.unwrap();
|
||||
}
|
||||
compact_task.await.unwrap().unwrap();
|
||||
|
||||
let table = conn
|
||||
.open_table("test_concurrent_compact")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let dataset = table.dataset().unwrap().get().await.unwrap();
|
||||
let fragment_ids = dataset
|
||||
.get_fragments()
|
||||
.iter()
|
||||
.map(|fragment| fragment.id())
|
||||
.collect::<Vec<_>>();
|
||||
assert!(fragment_ids.windows(2).all(|ids| ids[0] < ids[1]));
|
||||
|
||||
// A second compaction exposed the original out-of-order row-id bug.
|
||||
table
|
||||
.optimize(OptimizeAction::Compact {
|
||||
options: CompactionOptions {
|
||||
target_rows_per_fragment: 1_000,
|
||||
..Default::default()
|
||||
},
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
table.count_rows(None).await.unwrap(),
|
||||
ROWS_PER_FRAGMENT as usize * (NUM_FRAGMENTS * 2 + 1)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_optimize_prune_versions() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
|
||||
Reference in New Issue
Block a user