common_meta/ddl/
truncate_table.rs1use 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 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 Prepare,
265 DatanodeTruncateRegions,
267}