From 155cb546ace4e6b27ebe89ac2639582a75e1bd9d Mon Sep 17 00:00:00 2001 From: junxiang Mu <1948535941@qq.com> Date: Wed, 28 May 2025 10:57:59 +0000 Subject: [PATCH] improve scanner(1) Signed-off-by: junxiang Mu <1948535941@qq.com> --- ecstore/src/disk/local.rs | 16 ++++++++-------- ecstore/src/set_disk.rs | 20 +++++++++++++------- ecstore/src/store.rs | 2 +- 3 files changed, 22 insertions(+), 16 deletions(-) diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 601b4dbaf..ecb7f042c 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -2384,16 +2384,16 @@ impl DiskAPI for LocalDisk { } let stop_fn = ScannerMetrics::log(ScannerMetric::ScanObject); let mut res = HashMap::new(); - let done_sz = ScannerMetrics::time_size(ScannerMetric::ReadMetadata).await; + let done_sz = ScannerMetrics::time_size(ScannerMetric::ReadMetadata); let buf = match disk.read_metadata(item.path.clone()).await { Ok(buf) => buf, Err(err) => { res.insert("err".to_string(), err.to_string()); - stop_fn(&res).await; + stop_fn(&res); return Err(Error::from_string(ERR_SKIP_FILE)); } }; - done_sz(buf.len() as u64).await; + done_sz(buf.len() as u64); res.insert("metasize".to_string(), buf.len().to_string()); item.transform_meda_dir(); let meta_cache = MetaCacheEntry { @@ -2405,7 +2405,7 @@ impl DiskAPI for LocalDisk { Ok(fivs) => fivs, Err(err) => { res.insert("err".to_string(), err.to_string()); - stop_fn(&res).await; + stop_fn(&res); return Err(Error::from_string(ERR_SKIP_FILE)); } }; @@ -2415,7 +2415,7 @@ impl DiskAPI for LocalDisk { Ok(obj_infos) => obj_infos, Err(err) => { res.insert("err".to_string(), err.to_string()); - stop_fn(&res).await; + stop_fn(&res); return Err(Error::from_string(ERR_SKIP_FILE)); } }; @@ -2431,7 +2431,7 @@ impl DiskAPI for LocalDisk { let done = ScannerMetrics::time(ScannerMetric::ApplyVersion); let sz: usize; (obj_deleted, sz) = item.apply_actions(info, &size_s).await; - done().await; + done(); if obj_deleted { break; @@ -2461,14 +2461,14 @@ impl DiskAPI for LocalDisk { let _obj_info = frer_version.to_object_info(&item.bucket, &item.object_path().to_string_lossy(), versioned); let done = ScannerMetrics::time(ScannerMetric::TierObjSweep); - done().await; + done(); } // todo: global trace if obj_deleted { return Err(Error::from_string(ERR_IGNORE_FILE_CONTRIB)); } - done().await; + done(); Ok(size_s) }) }), diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 83b00b403..0e70f992b 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -2893,7 +2893,7 @@ impl SetDisks { } pub async fn ns_scanner( - &self, + self: Arc, buckets: &[BucketInfo], want_cycle: u32, updates: Sender, @@ -2911,7 +2911,7 @@ impl SetDisks { return Ok(()); } - let old_cache = DataUsageCache::load(self, DATA_USAGE_CACHE_NAME).await?; + let old_cache = DataUsageCache::load(&self, DATA_USAGE_CACHE_NAME).await?; let mut cache = DataUsageCache { info: DataUsageCacheInfo { name: DATA_USAGE_ROOT.to_string(), @@ -2934,6 +2934,7 @@ impl SetDisks { permutes.shuffle(&mut rng); permutes }; + // Add new buckets first for idx in permutes.iter() { let b = buckets[*idx].clone(); @@ -2954,6 +2955,7 @@ impl SetDisks { Duration::from_secs(30) + Duration::from_secs_f64(10.0 * rng.gen_range(0.0..1.0)) }; let mut ticker = interval(update_time); + let task = tokio::spawn(async move { let last_save = Some(SystemTime::now()); let mut need_loop = true; @@ -2983,8 +2985,8 @@ impl SetDisks { } } }); + // Restrict parallelism for disk usage scanner - // upto GOMAXPROCS if GOMAXPROCS is < len(disks) let max_procs = num_cpus::get(); if max_procs < disks.len() { disks = disks[0..max_procs].to_vec(); @@ -2997,15 +2999,16 @@ impl SetDisks { Some(disk) => disk.clone(), None => continue, }; + let self_clone = Arc::clone(&self); let bucket_rx_clone = bucket_rx.clone(); let buckets_results_tx_clone = buckets_results_tx.clone(); - futures.push(async move { + futures.push(tokio::spawn(async move { loop { match bucket_rx_clone.write().await.try_recv() { Err(_) => return, Ok(bucket_info) => { let cache_name = Path::new(&bucket_info.name).join(DATA_USAGE_CACHE_NAME); - let mut cache = match DataUsageCache::load(self, &cache_name.to_string_lossy()).await { + let mut cache = match DataUsageCache::load(&self_clone, &cache_name.to_string_lossy()).await { Ok(cache) => cache, Err(_) => continue, }; @@ -3022,6 +3025,7 @@ impl SetDisks { ..Default::default() }; } + // Collect updates. let (tx, mut rx) = mpsc::channel(1); let buckets_results_tx_inner_clone = buckets_results_tx_clone.clone(); @@ -3042,9 +3046,10 @@ impl SetDisks { } } }); + // Calc usage let before = cache.info.last_update; - let mut cache = match disk.clone().ns_scanner(&cache, tx, heal_scan_mode, None).await { + let mut cache = match disk.ns_scanner(&cache, tx, heal_scan_mode, None).await { Ok(cache) => cache, Err(_) => { if cache.info.last_update > before { @@ -3078,8 +3083,9 @@ impl SetDisks { } info!("continue scanner"); } - }); + })); } + info!("ns_scanner start"); let _ = join_all(futures).await; drop(buckets_results_tx); diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index dabdae4d2..4256e9848 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -827,7 +827,7 @@ impl ECStore { } } }); - if let Err(err) = set + if let Err(err) = set.clone() .ns_scanner(&all_buckets_clone, want_cycle as u32, tx, heal_scan_mode) .await {