Skip to main content

iceberg/cow_rewrite/
plan.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use futures::TryStreamExt;
19
20use crate::scan::FileScanTask;
21use crate::spec::DataFile;
22
23/// A data file selected for COW rewrite.
24///
25/// This is a read-only output type: the planning entry points that produce
26/// candidates are crate-internal and reached through
27/// [`crate::cow_rewrite::CowRewriteBuilder`], and the fields are exposed
28/// through accessors only. Follow-up commit-adapter work (overwrite and
29/// row-delta actions) consumes planned candidates through those accessors;
30/// construction from outside the crate is intentionally not exposed yet and
31/// can be added when an adapter actually needs it.
32#[derive(Debug, Clone)]
33pub struct CowRewriteFile {
34    /// Original data file from the manifest entry.
35    pub(crate) old_data_file: DataFile,
36    /// Full-read task for visible rows in the original file.
37    pub(crate) scan_task: FileScanTask,
38}
39
40impl CowRewriteFile {
41    /// Original data file from the manifest entry.
42    pub fn old_data_file(&self) -> &DataFile {
43        &self.old_data_file
44    }
45
46    /// Full-read task for visible rows in the original file.
47    pub fn scan_task(&self) -> &FileScanTask {
48        &self.scan_task
49    }
50}
51
52/// Plan data files that may need copy-on-write rewrite.
53pub(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}