Skip to main content

meta_srv/gc/
handler.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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        // limit gc region scope to regions whose datanode have reported stats(by heartbeat)
48        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            // Get table route information to map regions to peers
173            let (phy_table_id, table_peer) = self.ctx.get_table_route(table_id).await?;
174
175            if phy_table_id != table_id {
176                // Skip logical tables
177                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    /// Process multiple datanodes concurrently with limited parallelism.
213    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        // Create a stream of datanode GC tasks with limited concurrency
224        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) // Reuse table concurrency limit for datanodes
257        .collect()
258        .await;
259
260        // Process all datanode GC results and collect regions that need retry from table reports
261        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                    // Note: We don't have a direct way to map peer to table_id here,
269                    // so we just log the error. The table_reports will contain individual region failures.
270                    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    /// Process GC for a single datanode with all its candidate regions.
282    /// Returns the table reports for this datanode.
283    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        // Step 2: Run GC for all regions on this datanode in a single batch
304        let (gc_report, fully_listed_regions) = {
305            // Partition regions into full listing and fast listing in a single pass
306
307            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                    |(&region_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                    |(&region_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            // First process regions that can fast list
331            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                        // Add to need_retry_regions since it failed
357                        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                        // Add to need_retry_regions since it failed
390                        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        // Step 3: Process the combined GC report and update table reports
405        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 &region_id in region_ids {
428            if force_full_listing.contains(&region_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(&region_id) {
444                    if let Some(last_full_listing) = gc_info.last_full_listing_time {
445                        // check if pass cooling down interval after last full listing
446                        let elapsed = now.saturating_duration_since(last_full_listing);
447                        elapsed >= self.config.full_file_listing_interval
448                    } else {
449                        // Never did full listing for this region, do it now
450                        true
451                    }
452                } else {
453                    // First time GC for this region, skip doing full listing, for this time
454                    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}