Skip to main content

common_meta/ddl/
truncate_table.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::to_pb_time_ranges;
16use api::v1::region::{
17    RegionRequest, RegionRequestHeader, TruncateRequest as PbTruncateRegionRequest, region_request,
18    truncate_request,
19};
20use async_trait::async_trait;
21use common_procedure::error::{FromJsonSnafu, ToJsonSnafu};
22use common_procedure::{
23    Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure,
24    Result as ProcedureResult, Status,
25};
26use common_telemetry::debug;
27use common_telemetry::tracing_context::TracingContext;
28use futures::future::join_all;
29use serde::{Deserialize, Serialize};
30use snafu::{ResultExt, ensure};
31use store_api::storage::RegionId;
32use strum::AsRefStr;
33use table::metadata::{TableId, TableInfo};
34use table::table_name::TableName;
35use table::table_reference::TableReference;
36
37use crate::ddl::DdlContext;
38use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator};
39use crate::ddl::utils::{add_peer_context_if_needed, map_to_procedure_error};
40use crate::error::{ConvertTimeRangesSnafu, Result, TableNotFoundSnafu};
41use crate::key::DeserializedValueWithBytes;
42use crate::key::table_info::TableInfoValue;
43use crate::key::table_name::TableNameKey;
44use crate::lock_key::{CatalogLock, SchemaLock, TableLock};
45use crate::metrics;
46use crate::rpc::ddl::TruncateTableTask;
47use crate::rpc::router::{find_leader_regions, find_leaders};
48
49pub struct TruncateTableProcedure {
50    context: DdlContext,
51    data: TruncateTableData,
52}
53
54#[async_trait]
55impl Procedure for TruncateTableProcedure {
56    fn type_name(&self) -> &str {
57        Self::TYPE_NAME
58    }
59
60    async fn execute(&mut self, _ctx: &ProcedureContext) -> ProcedureResult<Status> {
61        let state = &self.data.state;
62
63        let _timer = metrics::METRIC_META_PROCEDURE_TRUNCATE_TABLE
64            .with_label_values(&[state.as_ref()])
65            .start_timer();
66
67        match self.data.state {
68            TruncateTableState::Prepare => self.on_prepare().await,
69            TruncateTableState::DatanodeTruncateRegions => {
70                self.on_datanode_truncate_regions().await
71            }
72        }
73        .map_err(map_to_procedure_error)
74    }
75
76    fn dump(&self) -> ProcedureResult<String> {
77        serde_json::to_string(&self.data).context(ToJsonSnafu)
78    }
79
80    fn lock_key(&self) -> LockKey {
81        let table_ref = &self.data.table_ref();
82        let table_id = self.data.table_id();
83        let lock_key = vec![
84            CatalogLock::Read(table_ref.catalog).into(),
85            SchemaLock::read(table_ref.catalog, table_ref.schema).into(),
86            TableLock::Write(table_id).into(),
87        ];
88
89        LockKey::new(lock_key)
90    }
91
92    fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
93        if !ctx
94            .event_type_filter
95            .allows(TableDdlEventType::TruncateTable.as_str())
96        {
97            return None;
98        }
99        let task = &self.data.task;
100        let locator = TableDdlLocator::new(&task.catalog, &task.schema, &task.table)
101            .with_table_id(task.table_id);
102        let event = match &ctx.trigger {
103            EventTrigger::Submitted => {
104                TableDdlEvent::truncate_table_submitted(locator, task.time_ranges.len())
105            }
106            _ => TableDdlEvent::lifecycle(TableDdlEventType::TruncateTable, [locator]),
107        };
108
109        Some(Box::new(event))
110    }
111}
112
113impl TruncateTableProcedure {
114    pub(crate) const TYPE_NAME: &'static str = "metasrv-procedure::TruncateTable";
115
116    pub(crate) fn new(
117        task: TruncateTableTask,
118        table_info_value: DeserializedValueWithBytes<TableInfoValue>,
119        context: DdlContext,
120    ) -> Self {
121        Self {
122            context,
123            data: TruncateTableData::new(task, table_info_value),
124        }
125    }
126
127    pub(crate) fn from_json(json: &str, context: DdlContext) -> ProcedureResult<Self> {
128        let data = serde_json::from_str(json).context(FromJsonSnafu)?;
129        Ok(Self { context, data })
130    }
131
132    // Checks whether the table exists.
133    async fn on_prepare(&mut self) -> Result<Status> {
134        let table_ref = &self.data.table_ref();
135
136        let manager = &self.context.table_metadata_manager;
137
138        let exist = manager
139            .table_name_manager()
140            .exists(TableNameKey::new(
141                table_ref.catalog,
142                table_ref.schema,
143                table_ref.table,
144            ))
145            .await?;
146
147        ensure!(
148            exist,
149            TableNotFoundSnafu {
150                table_name: table_ref.to_string()
151            }
152        );
153
154        self.data.state = TruncateTableState::DatanodeTruncateRegions;
155
156        Ok(Status::executing(true))
157    }
158
159    async fn on_datanode_truncate_regions(&mut self) -> Result<Status> {
160        let table_id = self.data.table_id();
161
162        let (_, physical_table_route) = self
163            .context
164            .table_metadata_manager
165            .table_route_manager()
166            .get_physical_table_route(table_id)
167            .await?;
168        let leaders = find_leaders(&physical_table_route.region_routes);
169        let mut truncate_region_tasks = Vec::with_capacity(leaders.len());
170
171        for datanode in leaders {
172            let requester = self.context.node_manager.datanode(&datanode).await;
173            let regions = find_leader_regions(&physical_table_route.region_routes, &datanode);
174
175            for region in regions {
176                let region_id = RegionId::new(table_id, region);
177                debug!(
178                    "Truncating table {} region {} on Datanode {:?}",
179                    self.data.table_ref(),
180                    region_id,
181                    datanode
182                );
183
184                let time_ranges = &self.data.task.time_ranges;
185                let kind = if time_ranges.is_empty() {
186                    truncate_request::Kind::All(api::v1::region::All {})
187                } else {
188                    let pb_time_ranges =
189                        to_pb_time_ranges(time_ranges).context(ConvertTimeRangesSnafu)?;
190                    truncate_request::Kind::TimeRanges(pb_time_ranges)
191                };
192
193                let request = RegionRequest {
194                    header: Some(RegionRequestHeader {
195                        tracing_context: TracingContext::from_current_span().to_w3c(),
196                        ..Default::default()
197                    }),
198                    body: Some(region_request::Body::Truncate(PbTruncateRegionRequest {
199                        region_id: region_id.as_u64(),
200                        kind: Some(kind),
201                    })),
202                };
203
204                let datanode = datanode.clone();
205                let requester = requester.clone();
206
207                truncate_region_tasks.push(async move {
208                    requester
209                        .handle(request)
210                        .await
211                        .map_err(add_peer_context_if_needed(datanode))
212                });
213            }
214        }
215
216        join_all(truncate_region_tasks)
217            .await
218            .into_iter()
219            .collect::<Result<Vec<_>>>()?;
220
221        Ok(Status::done())
222    }
223}
224
225#[derive(Debug, Serialize, Deserialize)]
226pub struct TruncateTableData {
227    state: TruncateTableState,
228    task: TruncateTableTask,
229    table_info_value: DeserializedValueWithBytes<TableInfoValue>,
230}
231
232impl TruncateTableData {
233    pub fn new(
234        task: TruncateTableTask,
235        table_info_value: DeserializedValueWithBytes<TableInfoValue>,
236    ) -> Self {
237        Self {
238            state: TruncateTableState::Prepare,
239            task,
240            table_info_value,
241        }
242    }
243
244    pub fn table_ref(&self) -> TableReference<'_> {
245        self.task.table_ref()
246    }
247
248    pub fn table_name(&self) -> TableName {
249        self.task.table_name()
250    }
251
252    fn table_info(&self) -> &TableInfo {
253        &self.table_info_value.table_info
254    }
255
256    fn table_id(&self) -> TableId {
257        self.table_info().ident.table_id
258    }
259}
260
261#[derive(Debug, Serialize, Deserialize, AsRefStr)]
262enum TruncateTableState {
263    /// Prepares to truncate the table
264    Prepare,
265    /// Truncates regions on Datanode
266    DatanodeTruncateRegions,
267}