Skip to main content

operator/req_convert/insert/
table_to_region.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::helper::vectors_to_rows;
16use api::v1::Rows;
17use api::v1::region::InsertRequests as RegionInsertRequests;
18use partition::manager::PartitionRuleManager;
19use table::metadata::TableInfo;
20use table::requests::InsertRequest as TableInsertRequest;
21
22use crate::error::Result;
23use crate::insert::InstantAndNormalInsertRequests;
24use crate::req_convert::common::partitioner::Partitioner;
25use crate::req_convert::common::{column_schema, row_count};
26
27pub struct TableToRegion<'a> {
28    table_info: &'a TableInfo,
29    partition_manager: &'a PartitionRuleManager,
30}
31
32impl<'a> TableToRegion<'a> {
33    pub fn new(table_info: &'a TableInfo, partition_manager: &'a PartitionRuleManager) -> Self {
34        Self {
35            table_info,
36            partition_manager,
37        }
38    }
39
40    /// Converts column vectors into rows without partition routing.
41    pub fn prepare(&self, request: TableInsertRequest) -> Result<Rows> {
42        let row_count = row_count(&request.columns_values)?;
43        let schema = column_schema(self.table_info, &request.columns_values)?;
44        let rows = vectors_to_rows(request.columns_values.values(), row_count);
45        Ok(Rows { schema, rows })
46    }
47
48    pub async fn convert(
49        &self,
50        request: TableInsertRequest,
51    ) -> Result<InstantAndNormalInsertRequests> {
52        let skip_wal = request.skip_wal;
53        let rows = self.prepare(request)?;
54        self.partition(rows, skip_wal).await
55    }
56
57    /// Routes already prepared rows while retaining TTL and WAL behavior.
58    pub async fn partition(
59        &self,
60        rows: Rows,
61        skip_wal: bool,
62    ) -> Result<InstantAndNormalInsertRequests> {
63        let requests = Partitioner::new(self.partition_manager)
64            .partition_insert_requests(self.table_info, rows, skip_wal)
65            .await?;
66
67        let requests = RegionInsertRequests { requests };
68        if self.table_info.is_ttl_instant_table() {
69            Ok(InstantAndNormalInsertRequests {
70                normal_requests: Default::default(),
71                instant_requests: requests,
72            })
73        } else {
74            Ok(InstantAndNormalInsertRequests {
75                normal_requests: requests,
76                instant_requests: Default::default(),
77            })
78        }
79    }
80}
81
82#[cfg(test)]
83mod tests {
84    use std::collections::HashMap;
85    use std::sync::Arc;
86
87    use api::v1::helper::tag_column_schema;
88    use api::v1::region::InsertRequest as RegionInsertRequest;
89    use api::v1::value::ValueData;
90    use api::v1::{ColumnDataType, PartitionExprVersion, Row, Rows, Value};
91    use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
92    use datatypes::vectors::{Int32Vector, VectorRef};
93    use store_api::storage::RegionId;
94    use table::requests::InsertRequest as TableInsertRequest;
95
96    use crate::req_convert::insert::table_to_region::TableToRegion;
97    use crate::test_util::{
98        create_partition_rule_manager, new_test_table_info, prepare_mocked_backend,
99    };
100
101    #[tokio::test]
102    async fn test_prepare_preserves_rows_before_routing() {
103        let backend = prepare_mocked_backend().await;
104        let partition_manager = create_partition_rule_manager(backend).await;
105        let table_info = new_test_table_info(1, "table_1", vec![0u32, 1, 2].into_iter());
106        let converter = TableToRegion::new(&table_info, &partition_manager);
107        let values = vec![Some(1), None, Some(11), Some(101)];
108        let request = build_table_request(Arc::new(Int32Vector::from(values.clone())));
109        let expected = build_region_request(values, 0, None, false).rows.unwrap();
110        assert_eq!(converter.prepare(request).unwrap(), expected);
111    }
112
113    #[tokio::test]
114    async fn test_insert_request_table_to_region() {
115        check_insert_request_table_to_region(false).await;
116        check_insert_request_table_to_region(true).await;
117    }
118
119    async fn check_insert_request_table_to_region(skip_wal: bool) {
120        // region to datanode placement:
121        // 1 -> 1
122        // 2 -> 2
123        // 3 -> 3
124        //
125        // region value ranges:
126        // 1 -> [50, max)
127        // 2 -> [10, 50)
128        // 3 -> (min, 10)
129
130        let backend = prepare_mocked_backend().await;
131        let partition_manager = create_partition_rule_manager(backend.clone()).await;
132        let table_info = new_test_table_info(1, "table_1", vec![0u32, 1, 2].into_iter());
133
134        let converter = TableToRegion::new(&table_info, &partition_manager);
135
136        let mut table_request = build_table_request(Arc::new(Int32Vector::from(vec![
137            Some(1),
138            None,
139            Some(11),
140            Some(101),
141        ])));
142        table_request.skip_wal = skip_wal;
143        let versions = partition_manager
144            .find_physical_partition_info(1)
145            .await
146            .unwrap()
147            .partitions
148            .iter()
149            .map(|p| (p.id.as_u64(), p.partition_expr_version))
150            .collect::<HashMap<_, _>>();
151
152        let region_requests = converter.convert(table_request).await.unwrap();
153        let mut region_id_to_region_requests = region_requests
154            .normal_requests
155            .requests
156            .into_iter()
157            .map(|r| (r.region_id, r))
158            .collect::<HashMap<_, _>>();
159
160        let region_id = RegionId::new(1, 1).as_u64();
161        let region_request = region_id_to_region_requests.remove(&region_id).unwrap();
162        assert_eq!(
163            region_request,
164            build_region_request(vec![Some(101)], region_id, versions[&region_id], skip_wal)
165        );
166
167        let region_id = RegionId::new(1, 2).as_u64();
168        let region_request = region_id_to_region_requests.remove(&region_id).unwrap();
169        assert_eq!(
170            region_request,
171            build_region_request(vec![Some(11)], region_id, versions[&region_id], skip_wal)
172        );
173
174        let region_id = RegionId::new(1, 3).as_u64();
175        let region_request = region_id_to_region_requests.remove(&region_id).unwrap();
176        assert_eq!(
177            region_request,
178            build_region_request(
179                vec![Some(1), None],
180                region_id,
181                versions[&region_id],
182                skip_wal
183            )
184        );
185    }
186
187    fn build_table_request(vector: VectorRef) -> TableInsertRequest {
188        TableInsertRequest {
189            catalog_name: DEFAULT_CATALOG_NAME.to_string(),
190            schema_name: DEFAULT_SCHEMA_NAME.to_string(),
191            table_name: "table_1".to_string(),
192            columns_values: HashMap::from([("a".to_string(), vector)]),
193            skip_wal: false,
194        }
195    }
196
197    fn build_region_request(
198        rows: Vec<Option<i32>>,
199        region_id: u64,
200        version: Option<u64>,
201        skip_wal: bool,
202    ) -> RegionInsertRequest {
203        RegionInsertRequest {
204            skip_wal,
205            region_id,
206            rows: Some(Rows {
207                schema: vec![tag_column_schema("a", ColumnDataType::Int32)],
208                rows: rows
209                    .into_iter()
210                    .map(|v| Row {
211                        values: vec![Value {
212                            value_data: v.map(ValueData::I32Value),
213                        }],
214                    })
215                    .collect(),
216            }),
217            partition_expr_version: version.map(|value| PartitionExprVersion { value }),
218        }
219    }
220}