operator/req_convert/insert/
table_to_region.rs1use 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 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 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 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(®ion_id).unwrap();
162 assert_eq!(
163 region_request,
164 build_region_request(vec![Some(101)], region_id, versions[®ion_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(®ion_id).unwrap();
169 assert_eq!(
170 region_request,
171 build_region_request(vec![Some(11)], region_id, versions[®ion_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(®ion_id).unwrap();
176 assert_eq!(
177 region_request,
178 build_region_request(
179 vec![Some(1), None],
180 region_id,
181 versions[®ion_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}