diff --git a/java/lance-jni/src/blocking_dataset.rs b/java/lance-jni/src/blocking_dataset.rs index 889c3e572eb..0e4d21dee98 100644 --- a/java/lance-jni/src/blocking_dataset.rs +++ b/java/lance-jni/src/blocking_dataset.rs @@ -59,7 +59,7 @@ use lance_namespace::LanceNamespace; use lance_table::io::commit::CommitHandler; use lance_table::io::commit::external_manifest::ExternalManifestCommitHandler; use lance_table::io::commit::{ManifestLocation, ManifestNamingScheme}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::future::IntoFuture; use std::iter::empty; use std::sync::Arc; @@ -3438,6 +3438,14 @@ fn extract_cleanup_policy(env: &mut JNIEnv<'_>, jpolicy: &JObject) -> Result, jpolicy: &JObject) -> Result beforeTimestampMillis; private final Optional beforeVersion; + private final Optional> versions; private final Optional deleteUnverified; private final Optional errorIfTaggedOldVersions; private final Optional cleanReferencedBranches; @@ -32,12 +36,14 @@ public class CleanupPolicy { private CleanupPolicy( Optional beforeTimestampMillis, Optional beforeVersion, + Optional> versions, Optional deleteUnverified, Optional errorIfTaggedOldVersions, Optional cleanReferencedBranches, Optional deleteRateLimit) { this.beforeTimestampMillis = beforeTimestampMillis; this.beforeVersion = beforeVersion; + this.versions = versions; this.deleteUnverified = deleteUnverified; this.errorIfTaggedOldVersions = errorIfTaggedOldVersions; this.cleanReferencedBranches = cleanReferencedBranches; @@ -56,6 +62,10 @@ public Optional getBeforeVersion() { return beforeVersion; } + public Optional> getVersions() { + return versions; + } + public Optional getDeleteUnverified() { return deleteUnverified; } @@ -76,6 +86,7 @@ public Optional getDeleteRateLimit() { public static class Builder { private Optional beforeTimestampMillis = Optional.empty(); private Optional beforeVersion = Optional.empty(); + private Optional> versions = Optional.empty(); private Optional deleteUnverified = Optional.empty(); private Optional errorIfTaggedOldVersions = Optional.empty(); private Optional cleanReferencedBranches = Optional.empty(); @@ -95,6 +106,12 @@ public Builder withBeforeVersion(long beforeVersion) { return this; } + /** Set the exact dataset versions to clean. */ + public Builder withVersions(List versions) { + this.versions = Optional.of(Collections.unmodifiableList(new ArrayList<>(versions))); + return this; + } + /** If true, delete unverified data files even if they are recent. */ public Builder withDeleteUnverified(boolean deleteUnverified) { this.deleteUnverified = Optional.of(deleteUnverified); @@ -123,6 +140,7 @@ public CleanupPolicy build() { return new CleanupPolicy( beforeTimestampMillis, beforeVersion, + versions, deleteUnverified, errorIfTaggedOldVersions, cleanReferencedBranches, diff --git a/java/src/test/java/org/lance/CleanupTest.java b/java/src/test/java/org/lance/CleanupTest.java index 5fc8ceeaa3f..f6f3f0ef043 100644 --- a/java/src/test/java/org/lance/CleanupTest.java +++ b/java/src/test/java/org/lance/CleanupTest.java @@ -54,6 +54,31 @@ public void testCleanupBeforeVersion(@TempDir Path tempDir) { } } + @Test + public void testCleanupSpecificVersions(@TempDir Path tempDir) { + String datasetPath = tempDir.resolve("test_dataset_for_cleanup").toString(); + try (RootAllocator allocator = new RootAllocator(Long.MAX_VALUE)) { + TestUtils.SimpleTestDataset testDataset = + new TestUtils.SimpleTestDataset(allocator, datasetPath); + + testDataset.createEmptyDataset().close(); + + testDataset.write(1, 10).close(); + testDataset.write(2, 10).close(); + + try (Dataset dataset = testDataset.write(3, 10)) { + assertEquals(4, dataset.listVersions().size()); + + RemovalStats stats = + dataset.cleanupWithPolicy(CleanupPolicy.builder().withVersions(List.of(2L)).build()); + + assertEquals(1L, stats.getOldVersions()); + assertEquals(3, dataset.listVersions().size()); + assertTrue(dataset.listVersions().stream().noneMatch(version -> version.getId() == 2L)); + } + } + } + @Test public void testExplainCleanupBeforeVersion(@TempDir Path tempDir) { String datasetPath = tempDir.resolve("test_dataset_for_cleanup").toString(); diff --git a/python/python/lance/dataset.py b/python/python/lance/dataset.py index e386e6cb854..0be922aa889 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -3143,6 +3143,7 @@ def cleanup_old_versions( delete_unverified: bool = False, error_if_tagged_old_versions: bool = True, delete_rate_limit: Optional[int] = None, + versions: Optional[List[int]] = None, ) -> CleanupStats: """ Cleans up old versions of the dataset. @@ -3188,8 +3189,13 @@ def cleanup_old_versions( deletions run at full speed. Set this to a positive integer to avoid hitting object store request rate limits (e.g. S3 HTTP 503 SlowDown). For example, ``delete_rate_limit=100`` limits to 100 operations/second. + + versions: list[int], optional + Clean up only the specified dataset versions. The current version is + never removed, and tagged versions are still protected by + ``error_if_tagged_old_versions``. """ - if older_than is None and retain_versions is None: + if older_than is None and retain_versions is None and versions is None: older_than = timedelta(days=14) return self._ds.cleanup_old_versions( @@ -3198,6 +3204,7 @@ def cleanup_old_versions( delete_unverified, error_if_tagged_old_versions, delete_rate_limit, + versions, ) def explain_cleanup_old_versions( @@ -3208,6 +3215,7 @@ def explain_cleanup_old_versions( delete_unverified: bool = False, error_if_tagged_old_versions: bool = True, delete_rate_limit: Optional[int] = None, + versions: Optional[List[int]] = None, include_files: bool = False, max_files: int = 1000, ) -> CleanupExplanation: @@ -3235,6 +3243,9 @@ def explain_cleanup_old_versions( Accepted for parity with :meth:`cleanup_old_versions`; no deletes are issued by explain. + versions: list[int], optional + Explain cleanup only for the specified dataset versions. + include_files: bool, default False If `True`, include candidate files in the explanation up to ``max_files`` entries. Aggregate stats always include all candidates. @@ -3243,7 +3254,7 @@ def explain_cleanup_old_versions( Maximum number of candidate files to include when ``include_files`` is `True`. """ - if older_than is None and retain_versions is None: + if older_than is None and retain_versions is None and versions is None: older_than = timedelta(days=14) if max_files <= 0: raise ValueError("max_files must be positive") @@ -3254,6 +3265,7 @@ def explain_cleanup_old_versions( delete_unverified, error_if_tagged_old_versions, delete_rate_limit, + versions, include_files, max_files, ) diff --git a/python/python/tests/test_dataset.py b/python/python/tests/test_dataset.py index ae49b8cab3d..6093ff2ce38 100644 --- a/python/python/tests/test_dataset.py +++ b/python/python/tests/test_dataset.py @@ -1734,6 +1734,24 @@ def test_cleanup_with_retain_versions(tmp_path: Path): assert ds.count_rows() == len(ds.to_table()) +def test_cleanup_specific_versions(tmp_path: Path): + base_dir = tmp_path / "cleanup_specific_versions" + table = pa.Table.from_pydict({"a": range(100), "b": range(100)}) + lance.write_dataset(table, base_dir, mode="create") + time.sleep(0.05) + lance.write_dataset(table, base_dir, mode="overwrite") + time.sleep(0.05) + lance.write_dataset(table, base_dir, mode="overwrite") + time.sleep(0.05) + ds = lance.write_dataset(table, base_dir, mode="append") + + assert [v["version"] for v in ds.versions()] == [1, 2, 3, 4] + + stats = ds.cleanup_old_versions(versions=[2]) + assert stats.old_versions == 1 + assert [v["version"] for v in ds.versions()] == [1, 3, 4] + + def test_cleanup_with_older_than_and_retain_versions(tmp_path: Path): base_dir = tmp_path / "cleanup_policy" table = pa.Table.from_pydict({"a": range(100), "b": range(100)}) diff --git a/python/src/dataset.rs b/python/src/dataset.rs index 1191309a1f7..bd243fade58 100644 --- a/python/src/dataset.rs +++ b/python/src/dataset.rs @@ -796,6 +796,7 @@ impl Dataset { delete_unverified: Option, error_if_tagged_old_versions: Option, delete_rate_limit: Option, + versions: Option>, ) -> lance_core::Result { let mut builder = CleanupPolicyBuilder::default(); if let Some(v) = older_than_micros { @@ -814,6 +815,9 @@ impl Dataset { if let Some(v) = delete_rate_limit { builder = builder.delete_rate_limit(v)?; } + if let Some(v) = versions { + builder = builder.versions(v)?; + } Ok(builder.build()) } } @@ -2147,7 +2151,7 @@ impl Dataset { } /// Cleanup old versions from the dataset - #[pyo3(signature = (older_than_micros = None, retain_versions = None, delete_unverified = None, error_if_tagged_old_versions = None, delete_rate_limit = None))] + #[pyo3(signature = (older_than_micros = None, retain_versions = None, delete_unverified = None, error_if_tagged_old_versions = None, delete_rate_limit = None, versions = None))] fn cleanup_old_versions( &self, older_than_micros: Option, @@ -2155,6 +2159,7 @@ impl Dataset { delete_unverified: Option, error_if_tagged_old_versions: Option, delete_rate_limit: Option, + versions: Option>, ) -> PyResult { let stats = rt() .block_on(None, async { @@ -2165,6 +2170,7 @@ impl Dataset { delete_unverified, error_if_tagged_old_versions, delete_rate_limit, + versions, ) .await?; self.ds.cleanup_with_policy(policy).await @@ -2175,7 +2181,7 @@ impl Dataset { /// Explain cleanup old versions from the dataset without deleting files #[allow(clippy::too_many_arguments)] - #[pyo3(signature = (older_than_micros = None, retain_versions = None, delete_unverified = None, error_if_tagged_old_versions = None, delete_rate_limit = None, include_files = false, max_files = 1000))] + #[pyo3(signature = (older_than_micros = None, retain_versions = None, delete_unverified = None, error_if_tagged_old_versions = None, delete_rate_limit = None, versions = None, include_files = false, max_files = 1000))] fn explain_cleanup_old_versions( &self, older_than_micros: Option, @@ -2183,6 +2189,7 @@ impl Dataset { delete_unverified: Option, error_if_tagged_old_versions: Option, delete_rate_limit: Option, + versions: Option>, include_files: bool, max_files: usize, ) -> PyResult { @@ -2195,6 +2202,7 @@ impl Dataset { delete_unverified, error_if_tagged_old_versions, delete_rate_limit, + versions, ) .await?; self.ds diff --git a/rust/lance/src/dataset/cleanup.rs b/rust/lance/src/dataset/cleanup.rs index 36ede66bc28..a8ff7ae477b 100644 --- a/rust/lance/src/dataset/cleanup.rs +++ b/rust/lance/src/dataset/cleanup.rs @@ -1289,6 +1289,8 @@ pub struct CleanupPolicy { pub before_timestamp: Option>, /// If not none, cleanup all versions before the specified version. pub before_version: Option, + /// If not none, cleanup only the specified versions. + pub versions: Option>, /// If true, delete unverified data files even if they are recent pub delete_unverified: bool, /// If true, return an Error if a tagged version is old @@ -1312,6 +1314,9 @@ impl CleanupPolicy { if let Some(before_version) = self.before_version { should_clean &= manifest.version < before_version; } + if let Some(versions) = self.versions.as_ref() { + should_clean &= versions.contains(&manifest.version); + } should_clean } } @@ -1321,6 +1326,7 @@ impl Default for CleanupPolicy { Self { before_timestamp: None, before_version: None, + versions: None, delete_unverified: false, error_if_tagged_old_versions: true, clean_referenced_branches: false, @@ -1347,6 +1353,24 @@ impl CleanupPolicyBuilder { self } + /// Cleanup only the specified dataset versions. + /// + /// This is an exact-version filter. If other policy filters are also + /// configured, a manifest is removed only when it satisfies all of them. + /// + /// # Errors + /// + /// Returns an error if `versions` is empty. + pub fn versions(mut self, versions: Vec) -> Result { + if versions.is_empty() { + return Err(Error::invalid_input( + "versions must not be empty when specified", + )); + } + self.policy.versions = Some(versions.into_iter().collect()); + Ok(self) + } + /// Cleanup all versions except the last `n` versions of the dataset. /// /// # Errors @@ -3472,6 +3496,37 @@ mod tests { ); } + #[tokio::test] + async fn cleanup_specific_versions_only() { + let fixture = MockDatasetFixture::try_new().unwrap(); + fixture.create_some_data().await.unwrap(); + fixture.overwrite_some_data().await.unwrap(); + fixture.overwrite_some_data().await.unwrap(); + + let before_count = fixture.count_files().await.unwrap(); + assert_eq!(before_count.num_manifest_files, 3); + + let policy = CleanupPolicyBuilder::default() + .versions(vec![2]) + .unwrap() + .build(); + let removed = fixture.run_cleanup_with_policy(policy).await.unwrap(); + + assert_eq!(removed.old_versions, 1); + + let versions = fixture + .open() + .await + .unwrap() + .version_refs() + .await + .unwrap() + .iter() + .map(|version| version.version) + .collect::>(); + assert_eq!(versions, vec![1, 3]); + } + #[tokio::test] async fn cleanup_before_ts_and_retain_n_recent_versions() { let fixture = MockDatasetFixture::try_new().unwrap();