@@ -38,14 +38,15 @@ use crate::table::Table;
3838/// Represents an operation to perform on a snapshot.
3939#[ derive( Debug , Clone ) ]
4040pub enum Operation {
41- /// Add rows with the given `n` values and `data` values. Example: `Add(vec![1, 2, 3],
42- /// vec!["a", "b", "c"])` adds three rows with n=1,2,3 and data="a","b","c"
43- Add ( Vec < i32 > , Vec < String > ) ,
44-
45- /// Delete rows by their global positions (uses positional deletes).
46- /// Positions are global indices across all data files in order of their creation.
47- /// Example: `Delete(vec![0, 1])` deletes the first and second rows from the table
48- Delete ( Vec < i64 > ) ,
41+ /// Add rows with the given (n, data) tuples, and write to the specified parquet file name.
42+ /// Example: `Add(vec![(1, "a".to_string()), (2, "b".to_string())], "data-1.parquet".to_string())`
43+ /// adds two rows with n=1,2 and data="a","b" to a file named "data-1.parquet"
44+ Add ( Vec < ( i32 , String ) > , String ) ,
45+
46+ /// Delete rows by their positions within specific parquet files (uses positional deletes).
47+ /// Takes a vector of (position, file_name) tuples specifying which position in which file to delete.
48+ /// Example: `Delete(vec![(0, "data-1.parquet"), (1, "data-1.parquet")])` deletes positions 0 and 1 from data-1.parquet
49+ Delete ( Vec < ( i64 , String ) > ) ,
4950}
5051
5152/// Tracks the state of data files across snapshots
@@ -62,9 +63,16 @@ struct DataFileInfo {
6263/// # Example
6364/// ```
6465/// let fixture = IncrementalTestFixture::new(vec![
65- /// Operation::Add(vec![], vec![]), // Empty snapshot
66- /// Operation::Add(vec![1, 2, 3], vec!["1", "2", "3"]), // Add 3 rows
67- /// Operation::Delete(vec![2]), // Delete row with n=2
66+ /// Operation::Add(vec![], "empty.parquet".to_string()), // Empty snapshot
67+ /// Operation::Add(
68+ /// vec![
69+ /// (1, "1".to_string()),
70+ /// (2, "2".to_string()),
71+ /// (3, "3".to_string()),
72+ /// ],
73+ /// "data-1.parquet".to_string(),
74+ /// ), // Add 3 rows
75+ /// Operation::Delete(vec![(1, "data-1.parquet".to_string())]), // Delete position 1 from data-1.parquet
6876/// ])
6977/// .await;
7078/// ```
@@ -259,7 +267,11 @@ impl IncrementalTestFixture {
259267 } ;
260268
261269 match operation {
262- Operation :: Add ( n_values, data_values) => {
270+ Operation :: Add ( rows, file_name) => {
271+ // Extract n_values and data_values from tuples
272+ let n_values: Vec < i32 > = rows. iter ( ) . map ( |( n, _) | * n) . collect ( ) ;
273+ let data_values: Vec < String > = rows. iter ( ) . map ( |( _, d) | d. clone ( ) ) . collect ( ) ;
274+
263275 // Create data manifest
264276 let mut data_writer = ManifestWriterBuilder :: new (
265277 self . next_manifest_file ( ) ,
@@ -299,9 +311,8 @@ impl IncrementalTestFixture {
299311
300312 // Add new data if not empty
301313 if !n_values. is_empty ( ) {
302- let data_file_path =
303- format ! ( "{}/data/data-{}.parquet" , & self . table_location, snapshot_id) ;
304- self . write_parquet_file ( & data_file_path, n_values, data_values)
314+ let data_file_path = format ! ( "{}/data/{}" , & self . table_location, file_name) ;
315+ self . write_parquet_file ( & data_file_path, & n_values, & data_values)
305316 . await ;
306317
307318 data_writer
@@ -330,7 +341,7 @@ impl IncrementalTestFixture {
330341 path : data_file_path,
331342 snapshot_id,
332343 sequence_number,
333- n_values : n_values . clone ( ) ,
344+ n_values,
334345 } ) ;
335346 }
336347
@@ -404,26 +415,15 @@ impl IncrementalTestFixture {
404415 }
405416
406417 Operation :: Delete ( positions_to_delete) => {
407- // Map global positions to file-specific positions
418+ // Group deletes by file
408419 let mut deletes_by_file: HashMap < String , Vec < i64 > > = HashMap :: new ( ) ;
409420
410- for global_pos in positions_to_delete {
411- // Find which data file contains this global position
412- let mut file_offset = 0i64 ;
413- for data_file in & data_files {
414- let file_size = data_file. n_values . len ( ) as i64 ;
415- if global_pos >= & file_offset && global_pos < & ( file_offset + file_size)
416- {
417- // This position belongs to this file
418- let local_pos = global_pos - file_offset;
419- deletes_by_file
420- . entry ( data_file. path . clone ( ) )
421- . or_default ( )
422- . push ( local_pos) ;
423- break ;
424- }
425- file_offset += file_size;
426- }
421+ for ( position, file_name) in positions_to_delete {
422+ let data_file_path = format ! ( "{}/data/{}" , & self . table_location, file_name) ;
423+ deletes_by_file
424+ . entry ( data_file_path)
425+ . or_default ( )
426+ . push ( * position) ;
427427 }
428428
429429 // Create data manifest with existing data files
@@ -655,7 +655,7 @@ impl IncrementalTestFixture {
655655 from_snapshot_id : i64 ,
656656 to_snapshot_id : i64 ,
657657 expected_appends : Vec < ( i32 , & str ) > ,
658- expected_deletes : Vec < u64 > ,
658+ expected_deletes : Vec < ( u64 , & str ) > ,
659659 ) {
660660 use arrow_array:: cast:: AsArray ;
661661 use arrow_select:: concat:: concat_batches;
@@ -717,9 +717,38 @@ impl IncrementalTestFixture {
717717 . column ( 0 )
718718 . as_primitive :: < arrow_array:: types:: UInt64Type > ( ) ;
719719
720- let mut deleted_pairs: Vec < u64 > = pos_array. iter ( ) . filter_map ( |v| v) . collect ( ) ;
720+ // The file path column is a RunArray (Run-End Encoded), so we need to decode it
721+ // RunArray stores repeated values efficiently. To access individual values, we use
722+ // the Array trait which handles the run-length decoding automatically.
723+ use arrow_array:: Array ;
724+ let file_path_column = delete_batch. column ( 1 ) ;
725+
726+ let mut deleted_pairs: Vec < ( u64 , String ) > = ( 0 ..delete_batch. num_rows ( ) )
727+ . map ( |i| {
728+ let pos = pos_array. value ( i) ;
729+ // Use Array::to_data() to get the decoded data, then cast back
730+ let file_path = {
731+ // Get a slice of the run array for this single row
732+ let slice = file_path_column. slice ( i, 1 ) ;
733+ // Cast the slice to a run array and get its values
734+ let run_arr = slice
735+ . as_any ( )
736+ . downcast_ref :: < arrow_array:: RunArray < arrow_array:: types:: Int32Type > > ( )
737+ . unwrap ( ) ;
738+ let values = run_arr. values ( ) ;
739+ let str_arr = values. as_string :: < i32 > ( ) ;
740+ str_arr. value ( 0 ) . to_string ( )
741+ } ;
742+ ( pos, file_path)
743+ } )
744+ . collect ( ) ;
721745 deleted_pairs. sort ( ) ;
722746
747+ let expected_deletes: Vec < ( u64 , String ) > = expected_deletes
748+ . into_iter ( )
749+ . map ( |( pos, file) | ( pos, file. to_string ( ) ) )
750+ . collect ( ) ;
751+
723752 assert_eq ! ( deleted_pairs, expected_deletes) ;
724753 } else {
725754 assert ! ( expected_deletes. is_empty( ) , "Expected deletes but got none" ) ;
@@ -730,13 +759,16 @@ impl IncrementalTestFixture {
730759#[ tokio:: test]
731760async fn test_incremental_fixture_simple ( ) {
732761 let fixture = IncrementalTestFixture :: new ( vec ! [
733- Operation :: Add ( vec![ ] , vec![ ] ) ,
734- Operation :: Add ( vec![ 1 , 2 , 3 ] , vec![
735- "1" . to_string( ) ,
736- "2" . to_string( ) ,
737- "3" . to_string( ) ,
738- ] ) ,
739- Operation :: Delete ( vec![ 1 ] ) , // Delete position 1 (n=2, data="2")
762+ Operation :: Add ( vec![ ] , "empty.parquet" . to_string( ) ) ,
763+ Operation :: Add (
764+ vec![
765+ ( 1 , "1" . to_string( ) ) ,
766+ ( 2 , "2" . to_string( ) ) ,
767+ ( 3 , "3" . to_string( ) ) ,
768+ ] ,
769+ "data-2.parquet" . to_string( ) ,
770+ ) ,
771+ Operation :: Delete ( vec![ ( 1 , "data-2.parquet" . to_string( ) ) ] ) , // Delete position 1 (n=2, data="2")
740772 ] )
741773 . await ;
742774
@@ -764,28 +796,44 @@ async fn test_incremental_fixture_simple() {
764796 . await ;
765797
766798 // Verify incremental scan from snapshot 2 to snapshot 3.
767- fixture. verify_incremental_scan ( 2 , 3 , vec ! [ ] , vec ! [ 1 ] ) . await ;
768-
769- // Verify incremental scan from snapshot 1 to snapshot 1.
799+ let data_file_path = format ! ( "{}/data/data-2.parquet" , fixture. table_location) ;
770800 fixture
771- . verify_incremental_scan ( 1 , 1 , vec ! [ ] , vec ! [ ] )
801+ . verify_incremental_scan ( 2 , 3 , vec ! [ ] , vec ! [ ( 1 , & data_file_path ) ] )
772802 . await ;
803+
804+ // Verify incremental scan from snapshot 1 to snapshot 1.
805+ fixture. verify_incremental_scan ( 1 , 1 , vec ! [ ] , vec ! [ ] ) . await ;
773806}
774807
775808#[ tokio:: test]
776809async fn test_incremental_fixture_complex ( ) {
777810 let fixture = IncrementalTestFixture :: new ( vec ! [
778- Operation :: Add ( vec![ ] , vec![ ] ) , // Snapshot 1: Empty
779- Operation :: Add ( vec![ 1 , 2 , 3 , 4 , 5 ] , vec![
780- "a" . to_string( ) ,
781- "b" . to_string( ) ,
782- "c" . to_string( ) ,
783- "d" . to_string( ) ,
784- "e" . to_string( ) ,
785- ] ) , // Snapshot 2: Add 5 rows (positions 0-4)
786- Operation :: Delete ( vec![ 1 , 3 ] ) , // Snapshot 3: Delete positions 1,3 (n=2,4; data=b,d)
787- Operation :: Add ( vec![ 6 , 7 ] , vec![ "f" . to_string( ) , "g" . to_string( ) ] ) , // Snapshot 4: Add 2 more rows (positions 5-6)
788- Operation :: Delete ( vec![ 0 , 2 , 4 , 5 , 6 ] ) , // Snapshot 5: Delete positions 0,2,4,5,6 (all remaining rows: n=1,3,5,6,7)
811+ Operation :: Add ( vec![ ] , "empty.parquet" . to_string( ) ) , // Snapshot 1: Empty
812+ Operation :: Add (
813+ vec![
814+ ( 1 , "a" . to_string( ) ) ,
815+ ( 2 , "b" . to_string( ) ) ,
816+ ( 3 , "c" . to_string( ) ) ,
817+ ( 4 , "d" . to_string( ) ) ,
818+ ( 5 , "e" . to_string( ) ) ,
819+ ] ,
820+ "data-2.parquet" . to_string( ) ,
821+ ) , // Snapshot 2: Add 5 rows (positions 0-4)
822+ Operation :: Delete ( vec![
823+ ( 1 , "data-2.parquet" . to_string( ) ) ,
824+ ( 3 , "data-2.parquet" . to_string( ) ) ,
825+ ] ) , // Snapshot 3: Delete positions 1,3 (n=2,4; data=b,d)
826+ Operation :: Add (
827+ vec![ ( 6 , "f" . to_string( ) ) , ( 7 , "g" . to_string( ) ) ] ,
828+ "data-4.parquet" . to_string( ) ,
829+ ) , // Snapshot 4: Add 2 more rows (positions 5-6)
830+ Operation :: Delete ( vec![
831+ ( 0 , "data-2.parquet" . to_string( ) ) ,
832+ ( 2 , "data-2.parquet" . to_string( ) ) ,
833+ ( 4 , "data-2.parquet" . to_string( ) ) ,
834+ ( 0 , "data-4.parquet" . to_string( ) ) ,
835+ ( 1 , "data-4.parquet" . to_string( ) ) ,
836+ ] ) , // Snapshot 5: Delete positions 0,2,4,5,6 (all remaining rows: n=1,3,5,6,7)
789837 ] )
790838 . await ;
791839
0 commit comments