operator/req_convert/common/
partitioner.rs1use 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}