iceberg/cow_rewrite/
plan.rs1use futures::TryStreamExt;
19
20use crate::scan::FileScanTask;
21use crate::spec::DataFile;
22
23#[derive(Debug, Clone)]
33pub struct CowRewriteFile {
34 pub(crate) old_data_file: DataFile,
36 pub(crate) scan_task: FileScanTask,
38}
39
40impl CowRewriteFile {
41 pub fn old_data_file(&self) -> &DataFile {
43 &self.old_data_file
44 }
45
46 pub fn scan_task(&self) -> &FileScanTask {
48 &self.scan_task
49 }
50}
51
52pub(crate) async fn plan_cow_rewrite_files(
54 table: &crate::table::Table,
55 predicate: Option<crate::expr::Predicate>,
56 snapshot_id: Option<i64>,
57 case_sensitive: bool,
58) -> crate::Result<Vec<CowRewriteFile>> {
59 let mut scan_builder = table
60 .scan()
61 .select_all()
62 .with_case_sensitive(case_sensitive);
63
64 if let Some(predicate) = predicate {
65 scan_builder = scan_builder.with_filter(predicate);
66 }
67
68 if let Some(snapshot_id) = snapshot_id {
69 scan_builder = scan_builder.snapshot_id(snapshot_id);
70 }
71
72 scan_builder
73 .build()?
74 .plan_cow_rewrite_files()
75 .await?
76 .try_collect()
77 .await
78}
79
80#[cfg(test)]
81mod tests {
82 use futures::TryStreamExt;
83
84 use crate::Result;
85 use crate::expr::Predicate;
86 use crate::test_utils::scan::TableTestFixture;
87
88 #[tokio::test]
89 async fn cow_planner_returns_old_file_and_full_read_task() -> Result<()> {
90 let mut fixture = TableTestFixture::new();
91 fixture.setup_manifest_files().await;
92
93 let mut files =
94 super::plan_cow_rewrite_files(&fixture.table, Some(Predicate::AlwaysTrue), None, true)
95 .await?;
96
97 assert_eq!(files.len(), 2);
98
99 files.sort_by_key(|file| file.old_data_file.file_path().to_string());
100 assert_eq!(
101 files[0].old_data_file.file_path(),
102 format!("{}/1.parquet", fixture.table_location)
103 );
104 assert_eq!(
105 files[1].old_data_file.file_path(),
106 format!("{}/3.parquet", fixture.table_location)
107 );
108
109 for file in files {
110 assert_eq!(
111 file.old_data_file.file_path(),
112 file.scan_task.data_file_path()
113 );
114 assert!(file.scan_task.predicate().is_none());
115 assert_eq!(
116 Some(file.old_data_file.record_count()),
117 file.scan_task.record_count()
118 );
119 assert_eq!(0, file.scan_task.start());
120 assert_eq!(
121 file.old_data_file.file_size_in_bytes(),
122 file.scan_task.length()
123 );
124 }
125
126 Ok(())
127 }
128
129 #[tokio::test]
130 async fn cow_planner_preserves_delete_files() -> Result<()> {
131 let mut fixture = TableTestFixture::new();
132 fixture.setup_deadlock_manifests().await;
133
134 let scan = fixture
135 .table
136 .scan()
137 .select_all()
138 .with_concurrency_limit(1)
139 .build()?;
140
141 let files = tokio::time::timeout(std::time::Duration::from_secs(5), async {
142 scan.plan_cow_rewrite_files()
143 .await?
144 .try_collect::<Vec<_>>()
145 .await
146 })
147 .await
148 .expect("COW planning should not deadlock")?;
149
150 assert_eq!(files.len(), 10);
151 for file in files {
152 assert_eq!(file.scan_task.deletes().len(), 1);
153 assert_eq!(
154 file.scan_task.deletes()[0].file_path(),
155 format!("{}/del.parquet", fixture.table_location)
156 );
157 }
158
159 Ok(())
160 }
161}