Skip to main content

operator/
request.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::sync::Arc;
16
17use api::helper::to_pb_time_unit;
18use api::v1::region::region_request::Body as RegionRequestBody;
19use api::v1::region::{
20    BuildIndexRequest, CompactRequest, CompactionTimeRange, FlushRequest, RegionRequestHeader,
21    TruncateRequest, Unflushed, build_index_request, truncate_request,
22};
23use catalog::CatalogManagerRef;
24use common_catalog::build_db_string;
25use common_catalog::consts::METRIC_ENGINE;
26use common_meta::node_manager::{AffectedRows, NodeManagerRef};
27use common_meta::peer::Peer;
28use common_telemetry::tracing_context::TracingContext;
29use common_telemetry::{debug, error, info};
30use common_time::range::TimestampRange;
31use common_time::timestamp::TimeUnit as TimestampUnit;
32use futures_util::future;
33use partition::cache::PhysicalPartitionInfo;
34use partition::manager::PartitionRuleManagerRef;
35use session::context::QueryContextRef;
36use snafu::prelude::*;
37use store_api::storage::RegionId;
38use table::requests::{BuildIndexTableRequest, CompactTableRequest, FlushTableRequest};
39use table::table_name::TableName;
40
41use crate::error::{
42    CatalogSnafu, FindRegionLeaderSnafu, FindTablePartitionRuleSnafu, JoinTaskSnafu,
43    NotSupportedSnafu, RequestRegionSnafu, Result, TableNotFoundSnafu,
44    UnsupportedRegionRequestSnafu,
45};
46use crate::region_req_factory::RegionRequestFactory;
47
48/// Region requester which processes flush, compact requests etc.
49pub struct Requester {
50    catalog_manager: CatalogManagerRef,
51    partition_manager: PartitionRuleManagerRef,
52    node_manager: NodeManagerRef,
53}
54
55pub type RequesterRef = Arc<Requester>;
56
57impl Requester {
58    pub fn new(
59        catalog_manager: CatalogManagerRef,
60        partition_manager: PartitionRuleManagerRef,
61        node_manager: NodeManagerRef,
62    ) -> Self {
63        Self {
64            catalog_manager,
65            partition_manager,
66            node_manager,
67        }
68    }
69
70    /// Handle the request to flush table.
71    pub async fn handle_table_flush(
72        &self,
73        request: FlushTableRequest,
74        ctx: QueryContextRef,
75    ) -> Result<AffectedRows> {
76        let partitions = &self
77            .get_table_partition_info(
78                &request.catalog_name,
79                &request.schema_name,
80                &request.table_name,
81            )
82            .await?
83            .partitions;
84
85        let requests = partitions
86            .iter()
87            .map(|partition| {
88                RegionRequestBody::Flush(FlushRequest {
89                    region_id: partition.id.into(),
90                })
91            })
92            .collect();
93
94        info!("Handle table manual flush request: {:?}", request);
95
96        self.do_request(
97            requests,
98            Some(build_db_string(&request.catalog_name, &request.schema_name)),
99            &ctx,
100        )
101        .await
102    }
103
104    /// Handle the request to build index for table.
105    pub async fn handle_table_build_index(
106        &self,
107        request: BuildIndexTableRequest,
108        ctx: QueryContextRef,
109    ) -> Result<AffectedRows> {
110        if matches!(
111            request.options,
112            Some(build_index_request::Options::SeriesIndex(_))
113        ) {
114            let table = self
115                .catalog_manager
116                .table(
117                    &request.catalog_name,
118                    &request.schema_name,
119                    &request.table_name,
120                    None,
121                )
122                .await
123                .context(CatalogSnafu)?;
124            let table = table.with_context(|| TableNotFoundSnafu {
125                table_name: common_catalog::format_full_table_name(
126                    &request.catalog_name,
127                    &request.schema_name,
128                    &request.table_name,
129                ),
130            })?;
131            let info = table.table_info();
132            ensure_build_series_index_supported(&info.meta.engine, info.is_physical_table())?;
133        }
134        let partitions = &self
135            .get_table_partition_info(
136                &request.catalog_name,
137                &request.schema_name,
138                &request.table_name,
139            )
140            .await?
141            .partitions;
142
143        let requests = partitions
144            .iter()
145            .map(|partition| {
146                RegionRequestBody::BuildIndex(BuildIndexRequest {
147                    region_id: partition.id.into(),
148                    options: request.options,
149                })
150            })
151            .collect();
152
153        info!(
154            "Handle table manual build index for table {}",
155            request.table_name
156        );
157        debug!("Request details: {:?}", request);
158
159        self.do_request(
160            requests,
161            Some(build_db_string(&request.catalog_name, &request.schema_name)),
162            &ctx,
163        )
164        .await
165    }
166
167    /// Handle the request to compact table.
168    pub async fn handle_table_compaction(
169        &self,
170        request: CompactTableRequest,
171        ctx: QueryContextRef,
172    ) -> Result<AffectedRows> {
173        let partitions = &self
174            .get_table_partition_info(
175                &request.catalog_name,
176                &request.schema_name,
177                &request.table_name,
178            )
179            .await?
180            .partitions;
181
182        let time_range = request
183            .time_range
184            .map(to_pb_compaction_time_range)
185            .transpose()?;
186        let requests = partitions
187            .iter()
188            .map(|partition| {
189                RegionRequestBody::Compact(CompactRequest {
190                    region_id: partition.id.into(),
191                    parallelism: request.parallelism,
192                    options: Some(request.compact_options),
193                    time_range,
194                })
195            })
196            .collect();
197
198        info!("Handle table manual compaction request: {:?}", request);
199
200        self.do_request(
201            requests,
202            Some(build_db_string(&request.catalog_name, &request.schema_name)),
203            &ctx,
204        )
205        .await
206    }
207
208    /// Handle the request to flush the region.
209    pub async fn handle_region_flush(
210        &self,
211        region_id: RegionId,
212        ctx: QueryContextRef,
213    ) -> Result<AffectedRows> {
214        let request = RegionRequestBody::Flush(FlushRequest {
215            region_id: region_id.into(),
216        });
217
218        info!("Handle region manual flush request: {region_id}");
219        self.do_request(vec![request], None, &ctx).await
220    }
221
222    /// Handle the request to compact the region.
223    pub async fn handle_region_compaction(
224        &self,
225        region_id: RegionId,
226        ctx: QueryContextRef,
227    ) -> Result<AffectedRows> {
228        let request = RegionRequestBody::Compact(CompactRequest {
229            region_id: region_id.into(),
230            parallelism: 1,
231            options: None, // todo(hl): maybe also support parameters in region compaction.
232            time_range: None,
233        });
234
235        info!("Handle region manual compaction request: {region_id}");
236        self.do_request(vec![request], None, &ctx).await
237    }
238
239    /// Discard all unflushed data from the region.
240    pub async fn handle_discard_unflushed_data(
241        &self,
242        region_id: RegionId,
243        ctx: QueryContextRef,
244    ) -> Result<AffectedRows> {
245        let request = RegionRequestBody::Truncate(TruncateRequest {
246            region_id: region_id.into(),
247            kind: Some(truncate_request::Kind::Unflushed(Unflushed {})),
248        });
249
250        info!("Handle region discard unflushed data request: {region_id}");
251        self.do_request(vec![request], None, &ctx).await
252    }
253
254    /// Discard all unflushed data from all regions of the table.
255    pub async fn handle_discard_unflushed_data_by_table(
256        &self,
257        table_name: TableName,
258        ctx: QueryContextRef,
259    ) -> Result<AffectedRows> {
260        let table = self
261            .catalog_manager
262            .table(
263                &table_name.catalog_name,
264                &table_name.schema_name,
265                &table_name.table_name,
266                None,
267            )
268            .await
269            .context(CatalogSnafu)?;
270        let table = table.with_context(|| TableNotFoundSnafu {
271            table_name: table_name.to_string(),
272        })?;
273        let table_info = table.table_info();
274        ensure_discard_unflushed_supported(
275            &table_info.meta.engine,
276            table_info.is_physical_table(),
277        )?;
278
279        let partitions = &self
280            .partition_manager
281            .find_physical_partition_info(table_info.ident.table_id)
282            .await
283            .with_context(|_| FindTablePartitionRuleSnafu {
284                table_name: table_name.to_string(),
285            })?
286            .partitions;
287        let requests = partitions
288            .iter()
289            .map(|partition| {
290                RegionRequestBody::Truncate(TruncateRequest {
291                    region_id: partition.id.into(),
292                    kind: Some(truncate_request::Kind::Unflushed(Unflushed {})),
293                })
294            })
295            .collect();
296
297        info!("Handle table discard unflushed data request: {table_name}");
298        self.do_request(
299            requests,
300            Some(build_db_string(
301                &table_name.catalog_name,
302                &table_name.schema_name,
303            )),
304            &ctx,
305        )
306        .await
307    }
308}
309
310fn to_pb_compaction_time_range(range: TimestampRange) -> Result<CompactionTimeRange> {
311    let (Some(start), Some(end)) = (*range.start(), *range.end()) else {
312        return crate::error::InvalidTimestampRangeSnafu {
313            start: format!("{:?}", range.start()),
314            end: format!("{:?}", range.end()),
315        }
316        .fail();
317    };
318    let original_start = start;
319    let original_end = end;
320    let start = start.convert_to(TimestampUnit::Second).with_context(|| {
321        crate::error::InvalidTimestampRangeSnafu {
322            start: format!("{original_start:?}"),
323            end: format!("{original_end:?}"),
324        }
325    })?;
326    let end = end
327        .convert_to_ceil(TimestampUnit::Second)
328        .with_context(|| crate::error::InvalidTimestampRangeSnafu {
329            start: format!("{original_start:?}"),
330            end: format!("{original_end:?}"),
331        })?;
332    ensure!(
333        start < end,
334        crate::error::InvalidTimestampRangeSnafu {
335            start: format!("{original_start:?}"),
336            end: format!("{original_end:?}"),
337        }
338    );
339
340    Ok(CompactionTimeRange {
341        start: start.value(),
342        end: end.value(),
343        time_unit: to_pb_time_unit(TimestampUnit::Second) as i32,
344    })
345}
346
347impl Requester {
348    async fn do_request(
349        &self,
350        requests: Vec<RegionRequestBody>,
351        db_string: Option<String>,
352        ctx: &QueryContextRef,
353    ) -> Result<AffectedRows> {
354        let request_factory = RegionRequestFactory::new(RegionRequestHeader {
355            tracing_context: TracingContext::from_current_span().to_w3c(),
356            dbname: db_string.unwrap_or_else(|| ctx.get_db_string()),
357            ..Default::default()
358        });
359
360        let tasks = requests.into_iter().map(|req_body| {
361            let request = request_factory.build_request(req_body.clone());
362            let partition_manager = self.partition_manager.clone();
363            let node_manager = self.node_manager.clone();
364            common_runtime::spawn_global(async move {
365                let peer =
366                    Self::find_region_leader_by_request(partition_manager, &req_body).await?;
367                node_manager
368                    .datanode(&peer)
369                    .await
370                    .handle(request)
371                    .await
372                    .context(RequestRegionSnafu)
373            })
374        });
375        let results = future::try_join_all(tasks).await.context(JoinTaskSnafu)?;
376
377        let affected_rows = results
378            .into_iter()
379            .map(|resp| resp.map(|r| r.affected_rows))
380            .sum::<Result<AffectedRows>>()?;
381
382        Ok(affected_rows)
383    }
384
385    async fn find_region_leader_by_request(
386        partition_manager: PartitionRuleManagerRef,
387        req: &RegionRequestBody,
388    ) -> Result<Peer> {
389        let region_id = match req {
390            RegionRequestBody::Flush(req) => req.region_id,
391            RegionRequestBody::Compact(req) => req.region_id,
392            RegionRequestBody::BuildIndex(req) => req.region_id,
393            RegionRequestBody::Truncate(req) => req.region_id,
394            _ => {
395                error!("Unsupported region request: {:?}", req);
396                return UnsupportedRegionRequestSnafu {}.fail();
397            }
398        };
399
400        partition_manager
401            .find_region_leader(region_id.into())
402            .await
403            .context(FindRegionLeaderSnafu)
404    }
405
406    async fn get_table_partition_info(
407        &self,
408        catalog: &str,
409        schema: &str,
410        table_name: &str,
411    ) -> Result<Arc<PhysicalPartitionInfo>> {
412        let table = self
413            .catalog_manager
414            .table(catalog, schema, table_name, None)
415            .await
416            .context(CatalogSnafu)?;
417
418        let table = table.with_context(|| TableNotFoundSnafu {
419            table_name: common_catalog::format_full_table_name(catalog, schema, table_name),
420        })?;
421        let table_info = table.table_info();
422
423        self.partition_manager
424            .find_physical_partition_info(table_info.ident.table_id)
425            .await
426            .with_context(|_| FindTablePartitionRuleSnafu {
427                table_name: common_catalog::format_full_table_name(catalog, schema, table_name),
428            })
429    }
430}
431
432fn ensure_build_series_index_supported(engine: &str, is_physical_table: bool) -> Result<()> {
433    ensure!(
434        engine == METRIC_ENGINE && is_physical_table,
435        NotSupportedSnafu {
436            feat: "building series indexes requires a physical metric table",
437        }
438    );
439    Ok(())
440}
441
442fn ensure_discard_unflushed_supported(engine: &str, is_physical_table: bool) -> Result<()> {
443    ensure!(
444        engine != METRIC_ENGINE || is_physical_table,
445        NotSupportedSnafu {
446            feat: "discarding unflushed data from a Metric Engine logical table"
447        }
448    );
449    Ok(())
450}
451
452#[cfg(test)]
453mod tests {
454    use api::v1::TimeUnit;
455    use common_time::Timestamp;
456    use common_time::range::TimestampRange;
457
458    use super::*;
459
460    #[test]
461    fn test_build_series_index_requires_physical_metric_table() {
462        ensure_build_series_index_supported(METRIC_ENGINE, true).unwrap();
463        assert!(ensure_build_series_index_supported(METRIC_ENGINE, false).is_err());
464        assert!(ensure_build_series_index_supported("mito", false).is_err());
465        assert!(ensure_build_series_index_supported("mito", true).is_err());
466    }
467
468    #[test]
469    fn test_to_pb_compaction_time_range_normalizes_mixed_units() {
470        let range = TimestampRange::new(
471            Timestamp::new_millisecond(1_500),
472            Timestamp::new_microsecond(2_500_000),
473        )
474        .unwrap();
475
476        let pb_range = to_pb_compaction_time_range(range).unwrap();
477        assert_eq!(1, pb_range.start);
478        assert_eq!(3, pb_range.end);
479        assert_eq!(TimeUnit::Second as i32, pb_range.time_unit);
480    }
481
482    #[test]
483    fn test_discard_unflushed_rejects_metric_logical_table() {
484        let error = ensure_discard_unflushed_supported(METRIC_ENGINE, false).unwrap_err();
485        assert!(matches!(error, crate::error::Error::NotSupported { .. }));
486
487        ensure_discard_unflushed_supported(METRIC_ENGINE, true).unwrap();
488        ensure_discard_unflushed_supported("mito", false).unwrap();
489    }
490
491    #[test]
492    fn test_to_pb_compaction_time_range() {
493        let range = TimestampRange::new(
494            Timestamp::new_microsecond(1_000),
495            Timestamp::new_microsecond(2_000),
496        )
497        .unwrap();
498
499        let pb_range = to_pb_compaction_time_range(range).unwrap();
500        assert_eq!(0, pb_range.start);
501        assert_eq!(1, pb_range.end);
502        assert_eq!(TimeUnit::Second as i32, pb_range.time_unit);
503    }
504}