diff --git a/rust/lance/src/dataset/cleanup.rs b/rust/lance/src/dataset/cleanup.rs index 4fcf432ac1a..ec19284e9f1 100644 --- a/rust/lance/src/dataset/cleanup.rs +++ b/rust/lance/src/dataset/cleanup.rs @@ -39,7 +39,6 @@ use crate::{Dataset, utils::temporal::utc_now}; use chrono::{DateTime, TimeDelta, Utc}; use dashmap::DashSet; use futures::future::try_join_all; -use futures::stream::BoxStream; use futures::{StreamExt, TryStreamExt, stream}; use humantime::parse_duration; use lance_core::{ @@ -69,8 +68,7 @@ use std::{ }, time::Duration, }; -use tokio::time::{MissedTickBehavior, interval}; -use tokio_stream::wrappers::IntervalStream; +use tokio::time; use tracing::{Span, debug, info, instrument}; #[derive(Clone, Debug, Default)] @@ -542,38 +540,66 @@ impl<'a> CleanupTask<'a> { let all_paths_to_remove = stream::iter(vec![unreferenced_paths, old_manifests_stream]).flatten(); - let paths_to_delete: BoxStream> = if let Some(rate) = - self.policy.delete_rate_limit - { - let duration = calculate_duration(self.dataset.object_store.scheme().to_string(), rate); - let mut ticker = interval(duration); - ticker.set_missed_tick_behavior(MissedTickBehavior::Delay); - IntervalStream::new(ticker) - .zip(all_paths_to_remove) - .map(|(_, path)| path) - .boxed() - } else { - all_paths_to_remove.boxed() - }; - - let delete_count = AtomicU64::new(0); - let paths_with_progress = paths_to_delete.map_ok(|path| { - let count = delete_count.fetch_add(1, Ordering::Relaxed) + 1; - if count % 1000 == 0 { - info!("Delete progress: {} files submitted", count); + let total_deleted = if let Some(rate) = self.policy.delete_rate_limit { + // Manual batch-level rate limiting: collect paths into batches + // and feed each batch to remove_stream individually, sleeping + // between batches. For S3 this is necessary because the + // delete_stream implementation uses .buffered(20) internally, + // which defeats any per-path rate limiting on the input stream. + let scheme = self.dataset.object_store.scheme().to_string(); + let max_batch = delete_batch_size(&scheme); + let batch_interval = calculate_batch_interval(&scheme, rate); + let mut all_paths = all_paths_to_remove.boxed(); + let mut batch: Vec> = Vec::with_capacity(max_batch); + let mut total: u64 = 0; + + loop { + match all_paths.try_next().await? { + Some(path) => { + total += 1; + batch.push(Ok(path)); + if batch.len() >= max_batch { + info!("Delete progress: {} files submitted", total); + let batch_stream = stream::iter(batch.drain(..).collect::>()); + self.dataset + .object_store + .remove_stream(batch_stream.boxed()) + .try_for_each(|_| future::ready(Ok(()))) + .await?; + time::sleep(batch_interval).await; + } + } + None => break, + } } - path - }); - - let delete_fut = self - .dataset - .object_store - .remove_stream(paths_with_progress.boxed()) - .try_for_each(|_| future::ready(Ok(()))); + if !batch.is_empty() { + info!("Delete progress: {} files submitted (final batch)", total); + let batch_stream = stream::iter(batch); + self.dataset + .object_store + .remove_stream(batch_stream.boxed()) + .try_for_each(|_| future::ready(Ok(()))) + .await?; + } + total + } else { + let delete_count = AtomicU64::new(0); + let paths_with_progress = all_paths_to_remove.map_ok(|path| { + let count = delete_count.fetch_add(1, Ordering::Relaxed) + 1; + if count % 1000 == 0 { + info!("Delete progress: {} files submitted", count); + } + path + }); - delete_fut.await?; + self.dataset + .object_store + .remove_stream(paths_with_progress.boxed()) + .try_for_each(|_| future::ready(Ok(()))) + .await?; + delete_count.load(Ordering::Relaxed) + }; - let total_deleted = delete_count.load(Ordering::Relaxed); info!( "Delete phase complete: {} files deleted ({} old manifests)", total_deleted, num_old_manifests, @@ -1005,23 +1031,36 @@ impl<'a> CleanupTask<'a> { } } -fn calculate_duration(scheme: String, rate: u64) -> Duration { - let batch_size = if scheme.to_lowercase().contains("s3") { - S3_DELETE_STREAM_BATCH_SIZE - } else if scheme.to_lowercase().contains("az") { - AZURE_DELETE_STREAM_BATCH_SIZE +/// Return the scheme-dependent batch size for delete operations. +/// +/// S3 DeleteObjects handles up to 1000 keys per call; Azure up to 256. +/// Other backends (local, GCS) are treated as single-file deletes. +fn delete_batch_size(scheme: &str) -> usize { + let lower = scheme.to_lowercase(); + if lower.contains("s3") { + S3_DELETE_STREAM_BATCH_SIZE as usize + } else if lower.contains("az") { + AZURE_DELETE_STREAM_BATCH_SIZE as usize } else { 1 - }; + } +} + +/// Compute the sleep duration between batch delete API calls. +/// +/// `rate` is the number of batch API requests per second (e.g., S3 +/// DeleteObjects calls). Each call deletes up to 1000 keys for S3, 256 +/// for Azure. +fn calculate_batch_interval(scheme: &str, rate: u64) -> Duration { let effective_rate = rate.max(1); - let path_rate = effective_rate * batch_size; + let batch_size = delete_batch_size(scheme); + let interval_ms = 1_000u64.div_ceil(effective_rate); info!( - "delete_rate_limit enabled: limit {} delete requests/sec", - effective_rate + "delete_rate_limit enabled: {} API requests/sec, batch_size={}, \ + interval={}ms between batches", + effective_rate, batch_size, interval_ms, ); - // convert user given op/s to the rate of issuing paths - let duration_ns = 1_000_000_000u64.div_ceil(path_rate).max(1); - Duration::from_nanos(duration_ns) + Duration::from_millis(interval_ms) } #[derive(Clone, Debug)] @@ -3692,63 +3731,74 @@ mod tests { } #[test] - fn test_calculate_duration_s3() { - // Normal case: duration is computed from S3 batch size and configured rate. - let normal_rate = 100; - let expected_duration_ns = - 1_000_000_000u64.div_ceil(normal_rate * S3_DELETE_STREAM_BATCH_SIZE); + fn test_delete_batch_size() { + assert_eq!( + delete_batch_size("s3"), + S3_DELETE_STREAM_BATCH_SIZE as usize + ); + assert_eq!( + delete_batch_size("s3+ddb"), + S3_DELETE_STREAM_BATCH_SIZE as usize + ); + assert_eq!( + delete_batch_size("az"), + AZURE_DELETE_STREAM_BATCH_SIZE as usize + ); + assert_eq!(delete_batch_size("file"), 1); + assert_eq!(delete_batch_size("memory"), 1); + } + + #[test] + fn test_calculate_batch_interval_s3() { + // rate=3 → one batch every 334ms (ceil(1000/3)) + assert_eq!( + calculate_batch_interval("s3", 3), + Duration::from_millis(334), + ); + + // rate=1 → one batch per second assert_eq!( - calculate_duration("s3".to_string(), normal_rate), - Duration::from_nanos(expected_duration_ns) + calculate_batch_interval("s3", 1), + Duration::from_millis(1000), ); - // Edge case: rate too small should be clamped to 1. - let min_rate_duration = calculate_duration("s3".to_string(), 1); - assert_eq!(calculate_duration("s3".to_string(), 0), min_rate_duration); + // rate=0 clamped to 1 → same as rate=1 + assert_eq!( + calculate_batch_interval("s3", 0), + Duration::from_millis(1000), + ); - // Edge case: computed duration_ns too small should be clamped to at least 1ns. - let very_large_rate = 2_000_000; + // Large rate → 1ms floor assert_eq!( - calculate_duration("s3".to_string(), very_large_rate), - Duration::from_nanos(1) + calculate_batch_interval("s3", 2_000_000), + Duration::from_millis(1), ); } #[tokio::test] async fn test_cleanup_with_rate_limit() { - // Create multiple versions with data files that will be deleted. + // Create enough versions so that the unreferenced files span multiple + // batches, exercising the sleep-between-batches logic. let fixture = MockDatasetFixture::try_new().unwrap(); fixture.create_some_data().await.unwrap(); - // Create several old versions for _ in 0..4 { fixture.overwrite_some_data().await.unwrap(); } MockClock::set_system_time(TimeDelta::try_days(10).unwrap().to_std().unwrap()); - // Set rate limit to 1 ops/second so cleanup of several files must take at least ~1s let policy = CleanupPolicyBuilder::default() .before_timestamp(utc_now() - TimeDelta::try_days(8).unwrap()) .delete_rate_limit(1) .unwrap() .build(); - let start = std::time::Instant::now(); let db = fixture.open().await.unwrap(); let stats = cleanup_old_versions(&db, policy).await.unwrap(); - let elapsed = start.elapsed(); - // We deleted old versions, so there should be removed files assert!( stats.old_versions > 0, "expected some old versions to be removed" ); - // With rate=1 and multiple files, it must take at least 2s - // (even just 2 deletions at 1/s means ≥2s) - assert!( - elapsed.as_millis() >= 2000, - "expected cleanup to be rate-limited (elapsed: {:?})", - elapsed - ); } }