1use std::collections::{HashMap, HashSet};
16use std::time::Instant;
17
18use common_catalog::consts::MITO_ENGINE;
19use common_meta::datanode::{RegionManifestInfo, RegionStat};
20use common_meta::peer::Peer;
21use common_procedure::ProcedureContext;
22use common_telemetry::tracing::Instrument as _;
23use common_telemetry::{debug, error, info, warn};
24use futures::StreamExt;
25use itertools::Itertools;
26use ordered_float::OrderedFloat;
27use store_api::region_engine::RegionRole;
28use store_api::storage::{GcReport, RegionId};
29use table::metadata::TableId;
30
31use crate::error::Result;
32use crate::gc::Region2Peers;
33use crate::gc::candidate::GcCandidate;
34use crate::gc::dropped::DroppedRegionCollector;
35use crate::gc::scheduler::{GcJobReport, GcScheduler};
36use crate::gc::tracker::RegionGcInfo;
37use crate::metrics::METRIC_META_GC_CANDIDATE_REGIONS;
38
39impl GcScheduler {
40 pub(crate) async fn trigger_gc(
41 &self,
42 procedure_context: ProcedureContext,
43 ) -> Result<GcJobReport> {
44 let start_time = Instant::now();
45 info!("Starting GC cycle");
46
47 let table_to_region_stats = self
49 .ctx
50 .get_table_to_region_stats()
51 .instrument(common_telemetry::tracing::info_span!(
52 "meta_gc_get_region_stats"
53 ))
54 .await?;
55 info!(
56 "Fetched region stats for {} tables",
57 table_to_region_stats.len()
58 );
59
60 let per_table_candidates = self
61 .select_gc_candidates(&table_to_region_stats)
62 .instrument(common_telemetry::tracing::info_span!(
63 "meta_gc_select_candidates"
64 ))
65 .await?;
66
67 let table_reparts = self
68 .ctx
69 .get_table_reparts()
70 .instrument(common_telemetry::tracing::info_span!(
71 "meta_gc_get_table_reparts"
72 ))
73 .await?;
74 let dropped_collector =
75 DroppedRegionCollector::new(self.ctx.as_ref(), &self.config, &self.region_gc_tracker);
76 let dropped_assignment = dropped_collector
77 .collect_and_assign(&table_reparts)
78 .instrument(common_telemetry::tracing::info_span!(
79 "meta_gc_collect_dropped_regions"
80 ))
81 .await?;
82 let candidate_count: usize = per_table_candidates.values().map(|c| c.len()).sum();
83 let dropped_count: usize = dropped_assignment
84 .regions_by_peer
85 .values()
86 .map(|regions| regions.len())
87 .sum();
88 METRIC_META_GC_CANDIDATE_REGIONS.set((candidate_count + dropped_count) as i64);
89
90 if per_table_candidates.is_empty() && dropped_assignment.regions_by_peer.is_empty() {
91 info!("No GC candidates found, skipping GC cycle");
92 return Ok(Default::default());
93 }
94
95 let mut datanode_to_candidates = self
96 .aggregate_candidates_by_datanode(per_table_candidates)
97 .instrument(common_telemetry::tracing::info_span!(
98 "meta_gc_aggregate_by_datanode"
99 ))
100 .await?;
101
102 self.merge_dropped_regions(&mut datanode_to_candidates, &dropped_assignment);
103
104 if datanode_to_candidates.is_empty() {
105 info!("No valid datanode candidates found, skipping GC cycle");
106 return Ok(Default::default());
107 }
108
109 let report = self
110 .parallel_process_datanodes(
111 datanode_to_candidates,
112 dropped_assignment.force_full_listing,
113 dropped_assignment.region_routes_override,
114 procedure_context,
115 )
116 .instrument(common_telemetry::tracing::info_span!(
117 "meta_gc_dispatch_to_datanodes"
118 ))
119 .await;
120
121 let duration = start_time.elapsed();
122 match &report {
123 GcJobReport::PerDatanode {
124 per_datanode_reports,
125 failed_datanodes,
126 } => {
127 info!(
128 "Finished GC cycle. Processed {} datanodes ({} failed). Duration: {:?}",
129 per_datanode_reports.len(),
130 failed_datanodes.len(),
131 duration
132 );
133 }
134 GcJobReport::Combined { report } => {
135 info!(
136 "Finished GC cycle with combined report. Deleted files: {}, deleted indexes: {}. Duration: {:?}",
137 report.deleted_files.len(),
138 report.deleted_indexes.len(),
139 duration
140 );
141 }
142 }
143 debug!("Detailed GC Job Report: {report:#?}");
144
145 Ok(report)
146 }
147
148 fn merge_dropped_regions(
149 &self,
150 datanode_to_candidates: &mut HashMap<Peer, Vec<(TableId, GcCandidate)>>,
151 assignment: &crate::gc::dropped::DroppedRegionAssignment,
152 ) {
153 for (peer, dropped_infos) in &assignment.regions_by_peer {
154 let entry = datanode_to_candidates.entry(peer.clone()).or_default();
155 for info in dropped_infos {
156 entry.push((info.table_id, dropped_candidate(info.region_id)));
157 }
158 }
159 }
160
161 pub(crate) async fn aggregate_candidates_by_datanode(
162 &self,
163 per_table_candidates: HashMap<TableId, Vec<GcCandidate>>,
164 ) -> Result<HashMap<Peer, Vec<(TableId, GcCandidate)>>> {
165 let mut datanode_to_candidates: HashMap<Peer, Vec<(TableId, GcCandidate)>> = HashMap::new();
166
167 for (table_id, candidates) in per_table_candidates {
168 if candidates.is_empty() {
169 continue;
170 }
171
172 let (phy_table_id, table_peer) = self.ctx.get_table_route(table_id).await?;
174
175 if phy_table_id != table_id {
176 continue;
178 }
179
180 let region_to_peer = table_peer
181 .region_routes
182 .iter()
183 .filter_map(|r| {
184 r.leader_peer
185 .as_ref()
186 .map(|peer| (r.region.id, peer.clone()))
187 })
188 .collect::<HashMap<RegionId, Peer>>();
189
190 for candidate in candidates {
191 if let Some(peer) = region_to_peer.get(&candidate.region_id) {
192 datanode_to_candidates
193 .entry(peer.clone())
194 .or_default()
195 .push((table_id, candidate));
196 } else {
197 warn!(
198 "Skipping region {} for table {}: no leader peer found",
199 candidate.region_id, table_id
200 );
201 }
202 }
203 }
204
205 info!(
206 "Aggregated GC candidates for {} datanodes",
207 datanode_to_candidates.len()
208 );
209 Ok(datanode_to_candidates)
210 }
211
212 pub(crate) async fn parallel_process_datanodes(
214 &self,
215 datanode_to_candidates: HashMap<Peer, Vec<(TableId, GcCandidate)>>,
216 force_full_listing_by_peer: HashMap<Peer, HashSet<RegionId>>,
217 region_routes_override_by_peer: HashMap<Peer, Region2Peers>,
218 procedure_context: ProcedureContext,
219 ) -> GcJobReport {
220 let mut per_datanode_reports = HashMap::new();
221 let mut failed_datanodes: HashMap<_, Vec<_>> = HashMap::new();
222
223 let results: Vec<_> = futures::stream::iter(
225 datanode_to_candidates
226 .into_iter()
227 .filter(|(_, candidates)| !candidates.is_empty()),
228 )
229 .map(|(peer, candidates)| {
230 let scheduler = self;
231 let peer_clone = peer.clone();
232 let force_full_listing = force_full_listing_by_peer
233 .get(&peer)
234 .cloned()
235 .unwrap_or_default();
236 let region_routes_override = region_routes_override_by_peer
237 .get(&peer)
238 .cloned()
239 .unwrap_or_default();
240 let procedure_context = procedure_context.clone();
241 async move {
242 (
243 peer,
244 scheduler
245 .process_datanode_gc(
246 peer_clone,
247 candidates,
248 force_full_listing,
249 region_routes_override,
250 procedure_context,
251 )
252 .await,
253 )
254 }
255 })
256 .buffer_unordered(self.config.max_concurrent_tables) .collect()
258 .await;
259
260 for (peer, result) in results {
262 match result {
263 Ok(dn_report) => {
264 per_datanode_reports.insert(peer.id, dn_report);
265 }
266 Err(e) => {
267 error!(e; "Failed to process datanode GC for peer {}", peer);
268 failed_datanodes.entry(peer.id).or_default().push(e);
271 }
272 }
273 }
274
275 GcJobReport::PerDatanode {
276 per_datanode_reports,
277 failed_datanodes,
278 }
279 }
280
281 pub(crate) async fn process_datanode_gc(
284 &self,
285 peer: Peer,
286 candidates: Vec<(TableId, GcCandidate)>,
287 force_full_listing: HashSet<RegionId>,
288 region_routes_override: Region2Peers,
289 procedure_context: ProcedureContext,
290 ) -> Result<GcReport> {
291 info!(
292 "Starting GC for datanode {} with {} candidate regions",
293 peer,
294 candidates.len()
295 );
296
297 if candidates.is_empty() {
298 return Ok(Default::default());
299 }
300
301 let all_region_ids: Vec<RegionId> = candidates.iter().map(|(_, c)| c.region_id).collect();
302
303 let (gc_report, fully_listed_regions) = {
305 let batch_full_listing_decisions = self
308 .batch_should_use_full_listing(&all_region_ids, &force_full_listing)
309 .await;
310
311 let need_full_list_regions = batch_full_listing_decisions
312 .iter()
313 .filter_map(
314 |(®ion_id, &need_full)| {
315 if need_full { Some(region_id) } else { None }
316 },
317 )
318 .collect_vec();
319 let fast_list_regions = batch_full_listing_decisions
320 .iter()
321 .filter_map(
322 |(®ion_id, &need_full)| {
323 if !need_full { Some(region_id) } else { None }
324 },
325 )
326 .collect_vec();
327
328 let mut combined_report = GcReport::default();
329
330 if !fast_list_regions.is_empty() {
332 match self
333 .ctx
334 .gc_regions(
335 &fast_list_regions,
336 false,
337 self.config.mailbox_timeout,
338 region_routes_override.clone(),
339 procedure_context.clone(),
340 )
341 .instrument(common_telemetry::tracing::info_span!(
342 "meta_gc_call_datanode",
343 peer = %peer,
344 mode = "fast",
345 region_count = fast_list_regions.len()
346 ))
347 .await
348 {
349 Ok(report) => combined_report.merge(report),
350 Err(e) => {
351 error!(
352 e; "Failed to GC regions {:?} on datanode {}",
353 fast_list_regions, peer,
354 );
355
356 combined_report
358 .need_retry_regions
359 .extend(fast_list_regions.clone());
360 }
361 }
362 }
363
364 if !need_full_list_regions.is_empty() {
365 match self
366 .ctx
367 .gc_regions(
368 &need_full_list_regions,
369 true,
370 self.config.mailbox_timeout,
371 region_routes_override,
372 procedure_context,
373 )
374 .instrument(common_telemetry::tracing::info_span!(
375 "meta_gc_call_datanode",
376 peer = %peer,
377 mode = "full",
378 region_count = need_full_list_regions.len()
379 ))
380 .await
381 {
382 Ok(report) => combined_report.merge(report),
383 Err(e) => {
384 error!(
385 e; "Failed to GC regions {:?} on datanode {}",
386 need_full_list_regions, peer,
387 );
388
389 combined_report
391 .need_retry_regions
392 .extend(need_full_list_regions.clone());
393 }
394 }
395 }
396 let fully_listed_regions = need_full_list_regions
397 .into_iter()
398 .filter(|r| !combined_report.need_retry_regions.contains(r))
399 .collect::<HashSet<_>>();
400
401 (combined_report, fully_listed_regions)
402 };
403
404 for region_id in &all_region_ids {
406 self.update_full_listing_time(*region_id, fully_listed_regions.contains(region_id))
407 .await;
408 }
409
410 info!(
411 "Completed GC for datanode {}: {} regions processed",
412 peer,
413 all_region_ids.len()
414 );
415
416 Ok(gc_report)
417 }
418
419 async fn batch_should_use_full_listing(
420 &self,
421 region_ids: &[RegionId],
422 force_full_listing: &HashSet<RegionId>,
423 ) -> HashMap<RegionId, bool> {
424 let mut result = HashMap::new();
425 let mut gc_tracker = self.region_gc_tracker.lock().await;
426 let now = Instant::now();
427 for ®ion_id in region_ids {
428 if force_full_listing.contains(®ion_id) {
429 gc_tracker
430 .entry(region_id)
431 .and_modify(|info| {
432 info.last_full_listing_time = Some(now);
433 info.last_gc_time = now;
434 })
435 .or_insert_with(|| RegionGcInfo {
436 last_gc_time: now,
437 last_full_listing_time: Some(now),
438 });
439 result.insert(region_id, true);
440 continue;
441 }
442 let use_full_listing = {
443 if let Some(gc_info) = gc_tracker.get(®ion_id) {
444 if let Some(last_full_listing) = gc_info.last_full_listing_time {
445 let elapsed = now.saturating_duration_since(last_full_listing);
447 elapsed >= self.config.full_file_listing_interval
448 } else {
449 true
451 }
452 } else {
453 gc_tracker.insert(
455 region_id,
456 RegionGcInfo {
457 last_gc_time: now,
458 last_full_listing_time: Some(now),
459 },
460 );
461 false
462 }
463 };
464 result.insert(region_id, use_full_listing);
465 }
466 result
467 }
468}
469
470fn dropped_candidate(region_id: RegionId) -> GcCandidate {
471 GcCandidate {
472 region_id,
473 score: OrderedFloat(0.0),
474 region_stat: dropped_region_stat(region_id),
475 }
476}
477
478fn dropped_region_stat(region_id: RegionId) -> RegionStat {
479 RegionStat {
480 id: region_id,
481 rcus: 0,
482 wcus: 0,
483 approximate_bytes: 0,
484 engine: MITO_ENGINE.to_string(),
485 role: RegionRole::Leader,
486 num_rows: 0,
487 memtable_size: 0,
488 manifest_size: 0,
489 sst_size: 0,
490 sst_num: 0,
491 index_size: 0,
492 region_manifest: RegionManifestInfo::Mito {
493 manifest_version: 0,
494 flushed_entry_id: 0,
495 file_removed_cnt: 0,
496 },
497 written_bytes: 0,
498 query_cpu_time: 0,
499 query_scanned_bytes: 0,
500 data_topic_latest_entry_id: 0,
501 metadata_topic_latest_entry_id: 0,
502 }
503}