Skip to main content

operator/req_convert/common/
partitioner.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 api::v1::region::{DeleteRequest, InsertRequest};
16use api::v1::{PartitionExprVersion, Rows};
17use partition::manager::PartitionRuleManager;
18use snafu::ResultExt;
19use store_api::storage::RegionId;
20use table::metadata::TableInfo;
21
22use crate::error::{Result, SplitDeleteSnafu, SplitInsertSnafu};
23
24pub struct Partitioner<'a> {
25    partition_manager: &'a PartitionRuleManager,
26}
27
28impl<'a> Partitioner<'a> {
29    pub fn new(partition_manager: &'a PartitionRuleManager) -> Self {
30        Self { partition_manager }
31    }
32
33    pub async fn partition_insert_requests(
34        &self,
35        table_info: &TableInfo,
36        rows: Rows,
37        skip_wal: bool,
38    ) -> Result<Vec<InsertRequest>> {
39        let table_id = table_info.table_id();
40        let requests = self
41            .partition_manager
42            .split_rows(table_info, rows)
43            .await
44            .context(SplitInsertSnafu)?
45            .into_iter()
46            .map(
47                |(region_number, (rows, partition_expr_version))| InsertRequest {
48                    skip_wal,
49                    region_id: RegionId::new(table_id, region_number).into(),
50                    rows: Some(rows),
51                    partition_expr_version: partition_expr_version
52                        .map(|value| PartitionExprVersion { value }),
53                },
54            )
55            .collect();
56        Ok(requests)
57    }
58
59    pub async fn partition_delete_requests(
60        &self,
61        table_info: &TableInfo,
62        rows: Rows,
63    ) -> Result<Vec<DeleteRequest>> {
64        let table_id = table_info.table_id();
65
66        let requests = self
67            .partition_manager
68            .split_rows(table_info, rows)
69            .await
70            .context(SplitDeleteSnafu)?
71            .into_iter()
72            .map(
73                |(region_number, (rows, partition_expr_version))| DeleteRequest {
74                    region_id: RegionId::new(table_id, region_number).into(),
75                    rows: Some(rows),
76                    partition_expr_version: partition_expr_version
77                        .map(|value| PartitionExprVersion { value }),
78                },
79            )
80            .collect();
81        Ok(requests)
82    }
83}