Skip to main content

common_meta/ddl/
alter_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
15mod executor;
16mod metadata;
17mod region_request;
18
19use std::collections::HashSet;
20use std::vec;
21
22use api::region::RegionResponse;
23use api::v1::alter_table_expr::Kind;
24use api::v1::{RenameTable, SetTableOptions};
25use async_trait::async_trait;
26use common_catalog::consts::{METRIC_ENGINE, MITO_ENGINE};
27use common_error::ext::BoxedError;
28use common_procedure::error::{FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu};
29use common_procedure::{
30    Context as ProcedureContext, ContextProvider, Error as ProcedureError, EventContext,
31    EventTrigger, LockKey, PoisonKey, PoisonKeys, Procedure, ProcedureId, Status, StringKey,
32};
33use common_telemetry::{error, info, warn};
34use serde::{Deserialize, Serialize};
35use snafu::{ResultExt, ensure};
36use store_api::metadata::ColumnMetadata;
37use store_api::metric_engine_consts::TABLE_COLUMN_METADATA_EXTENSION_KEY;
38use store_api::storage::RegionId;
39use strum::AsRefStr;
40use table::metadata::{TableId, TableInfo};
41use table::requests::SKIP_WAL_KEY;
42use table::table_reference::TableReference;
43
44use crate::ddl::DdlContext;
45use crate::ddl::alter_table::executor::AlterTableExecutor;
46use crate::ddl::event::table::{
47    TableDdlEvent, TableDdlEventType, TableDdlLocator, alter_table_kind_name,
48};
49use crate::ddl::utils::{
50    MultipleResults, extract_column_metadatas, handle_multiple_results, map_to_procedure_error,
51    sync_follower_regions,
52};
53use crate::error::{
54    AbortProcedureSnafu, ConvertAlterTableRequestSnafu, NoLeaderSnafu, PutPoisonSnafu, Result,
55    RetryLaterSnafu, UnsupportedSnafu,
56};
57use crate::key::table_info::TableInfoValue;
58use crate::key::{DeserializedValueWithBytes, RegionDistribution};
59use crate::lock_key::{CatalogLock, RegionLock, SchemaLock, TableLock, TableNameLock};
60use crate::metrics;
61use crate::poison_key::table_poison_key;
62use crate::rpc::ddl::AlterTableTask;
63use crate::rpc::router::{RegionRoute, find_leaders, region_distribution};
64
65/// The alter table procedure
66pub struct AlterTableProcedure {
67    /// The runtime context.
68    context: DdlContext,
69    /// The serialized data.
70    data: AlterTableData,
71    /// Cached new table metadata in the prepare step.
72    /// If we recover the procedure from json, then the table info value is not cached.
73    /// But we already validated it in the prepare step.
74    new_table_info: Option<TableInfo>,
75    /// The alter table executor.
76    executor: AlterTableExecutor,
77}
78
79#[derive(Debug)]
80pub(crate) struct RegionRouteChanged;
81
82/// Builds the executor from the [`AlterTableData`].
83///
84/// # Panics
85/// - If the alter kind is not set.
86fn build_executor_from_alter_expr(alter_data: &AlterTableData) -> AlterTableExecutor {
87    let table_name = alter_data.table_ref().into();
88    let table_id = alter_data.table_id;
89    let alter_kind = alter_data.task.alter_table.kind.as_ref().unwrap();
90    let new_table_name = if let Kind::RenameTable(RenameTable { new_table_name }) = alter_kind {
91        Some(new_table_name.clone())
92    } else {
93        None
94    };
95    AlterTableExecutor::new(table_name, table_id, new_table_name)
96}
97
98impl AlterTableProcedure {
99    pub const TYPE_NAME: &'static str = "metasrv-procedure::AlterTable";
100
101    pub fn new(table_id: TableId, task: AlterTableTask, context: DdlContext) -> Result<Self> {
102        Self::new_with_region_locks(table_id, task, vec![], context)
103    }
104
105    pub(crate) fn new_with_region_locks(
106        table_id: TableId,
107        task: AlterTableTask,
108        region_locks: Vec<RegionId>,
109        context: DdlContext,
110    ) -> Result<Self> {
111        task.validate()?;
112        let data = AlterTableData::new(task, table_id, region_locks);
113        let executor = build_executor_from_alter_expr(&data);
114        Ok(Self {
115            context,
116            data,
117            new_table_info: None,
118            executor,
119        })
120    }
121
122    pub fn from_json(json: &str, context: DdlContext) -> ProcedureResult<Self> {
123        let data: AlterTableData = serde_json::from_str(json).context(FromJsonSnafu)?;
124        let executor = build_executor_from_alter_expr(&data);
125
126        Ok(AlterTableProcedure {
127            context,
128            data,
129            new_table_info: None,
130            executor,
131        })
132    }
133
134    // Checks whether the table exists.
135    pub(crate) async fn on_prepare(&mut self) -> Result<Status> {
136        self.executor
137            .on_prepare(&self.context.table_metadata_manager)
138            .await?;
139        self.fill_table_info().await?;
140
141        // Safety: filled in `fill_table_info`.
142        let table_info_value = self.data.table_info_value.as_ref().unwrap();
143        let new_table_info = AlterTableExecutor::validate_alter_table_expr(
144            &table_info_value.table_info,
145            self.data.task.alter_table.clone(),
146        )?;
147        self.new_table_info = Some(new_table_info);
148
149        // Safety: Checked in `AlterTableProcedure::new`.
150        let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap();
151        if enables_skip_wal(alter_kind) {
152            ensure!(
153                only_enables_skip_wal(alter_kind),
154                UnsupportedSnafu {
155                    operation: "combining skip_wal = 'true' with other table options".to_string()
156                }
157            );
158            let engine = table_info_value.table_info.meta.engine.as_str();
159            ensure!(
160                engine == MITO_ENGINE || engine == METRIC_ENGINE,
161                UnsupportedSnafu {
162                    operation: format!("setting skip_wal on {engine} engine tables")
163                }
164            );
165            // Persist the irreversible intent before any region stops writing WAL.
166            let (physical_table_id, physical_table_route) = self
167                .context
168                .table_metadata_manager
169                .table_route_manager()
170                .get_physical_table_route(self.data.table_id())
171                .await?;
172            ensure!(
173                physical_table_id == self.data.table_id(),
174                UnsupportedSnafu {
175                    operation: "setting skip_wal on logical tables".to_string()
176                }
177            );
178            let current_region_ids = physical_table_route
179                .region_routes
180                .iter()
181                .map(|route| route.region.id)
182                .collect::<HashSet<_>>();
183            let locked_region_ids = self
184                .data
185                .region_locks
186                .iter()
187                .copied()
188                .collect::<HashSet<_>>();
189            if physical_table_route.region_routes.len() != self.data.region_locks.len()
190                || current_region_ids != locked_region_ids
191            {
192                // Procedure lock keys are fixed at submission. Finish this attempt so the
193                // DDL manager can submit another procedure with locks for the current route.
194                return Ok(Status::done_with_output(RegionRouteChanged));
195            }
196            ensure!(
197                !find_leaders(&physical_table_route.region_routes).is_empty(),
198                NoLeaderSnafu {
199                    table_id: physical_table_id
200                }
201            );
202            self.data.region_distribution =
203                Some(region_distribution(&physical_table_route.region_routes));
204        }
205        self.data.state = self.data.flow()?.after_prepare();
206        Ok(Status::executing(true))
207    }
208
209    fn table_poison_key(&self) -> PoisonKey {
210        table_poison_key(self.data.table_id())
211    }
212
213    async fn put_poison(
214        &self,
215        ctx_provider: &dyn ContextProvider,
216        procedure_id: ProcedureId,
217    ) -> Result<()> {
218        let poison_key = self.table_poison_key();
219        ctx_provider
220            .try_put_poison(&poison_key, procedure_id)
221            .await
222            .context(PutPoisonSnafu)
223    }
224
225    pub async fn submit_alter_region_requests(
226        &mut self,
227        procedure_id: ProcedureId,
228        ctx_provider: &dyn ContextProvider,
229    ) -> Result<Status> {
230        let table_id = self.data.table_id();
231        let (_, physical_table_route) = self
232            .context
233            .table_metadata_manager
234            .table_route_manager()
235            .get_physical_table_route(table_id)
236            .await?;
237
238        self.data.region_distribution =
239            Some(region_distribution(&physical_table_route.region_routes));
240        let leaders = find_leaders(&physical_table_route.region_routes);
241        let alter_kind = self.make_region_alter_kind()?;
242
243        info!(
244            "Submitting alter region requests for table {}, table_id: {}, alter_kind: {:?}",
245            self.data.table_ref(),
246            table_id,
247            alter_kind,
248        );
249
250        ensure!(!leaders.is_empty(), NoLeaderSnafu { table_id });
251        // Puts the poison before submitting alter region requests to datanodes.
252        self.put_poison(ctx_provider, procedure_id).await?;
253        let flow = self.data.flow()?;
254        if flow == AlterTableFlow::MetadataFirst {
255            let results = self
256                .executor
257                .on_alter_skip_wal_regions(
258                    &self.context.node_manager,
259                    &physical_table_route.region_routes,
260                    alter_kind,
261                )
262                .await;
263
264            return match handle_multiple_results(results) {
265                MultipleResults::PartialRetryable(error) => Err(error),
266                MultipleResults::PartialNonRetryable(error)
267                | MultipleResults::AllNonRetryable(error) => {
268                    // The metadata already enables skip-WAL. Retry the idempotent request
269                    // so later attempts can update the remaining replicas.
270                    Err(BoxedError::new(error)).context(RetryLaterSnafu {
271                        clean_poisons: true,
272                    })
273                }
274                MultipleResults::AllRetryable(error) => {
275                    Err(BoxedError::new(error)).context(RetryLaterSnafu {
276                        clean_poisons: true,
277                    })
278                }
279                MultipleResults::Ok(results) => {
280                    self.handle_alter_region_response(results)?;
281                    Ok(Status::executing_with_clean_poisons(true))
282                }
283            };
284        }
285
286        let results = self
287            .executor
288            .on_alter_regions(
289                &self.context.node_manager,
290                &physical_table_route.region_routes,
291                alter_kind,
292            )
293            .await;
294
295        match handle_multiple_results(results) {
296            MultipleResults::PartialRetryable(error) => {
297                // Just returns the error, and wait for the next try.
298                Err(error)
299            }
300            MultipleResults::PartialNonRetryable(error) => {
301                error!(error; "Partial non-retryable errors occurred during alter table, table {}, table_id: {}", self.data.table_ref(), self.data.table_id());
302                // No retry will be done.
303                Ok(Status::poisoned(
304                    Some(self.table_poison_key()),
305                    ProcedureError::external(error),
306                ))
307            }
308            MultipleResults::AllRetryable(error) => {
309                // Just returns the error, and wait for the next try.
310                let err = BoxedError::new(error);
311                Err(err).context(RetryLaterSnafu {
312                    clean_poisons: true,
313                })
314            }
315            MultipleResults::Ok(results) => {
316                self.submit_sync_region_requests(&results, &physical_table_route.region_routes)
317                    .await;
318                self.handle_alter_region_response(results)?;
319                Ok(Status::executing_with_clean_poisons(true))
320            }
321            MultipleResults::AllNonRetryable(error) => {
322                error!(error; "All alter requests returned non-retryable errors for table {}, table_id: {}", self.data.table_ref(), self.data.table_id());
323                // It assumes the metadata on datanode is not changed.
324                // Case: The alter region request is sent but not applied. (e.g., InvalidArgument)
325
326                let err = BoxedError::new(error);
327                Err(err).context(AbortProcedureSnafu {
328                    clean_poisons: true,
329                })
330            }
331        }
332    }
333
334    fn handle_alter_region_response(&mut self, mut results: Vec<RegionResponse>) -> Result<()> {
335        if let Some(column_metadatas) =
336            extract_column_metadatas(&mut results, TABLE_COLUMN_METADATA_EXTENSION_KEY)?
337        {
338            self.data.column_metadatas = column_metadatas;
339        } else {
340            warn!(
341                "altering table result doesn't contains extension key `{TABLE_COLUMN_METADATA_EXTENSION_KEY}`,leaving the table's column metadata unchanged"
342            );
343        }
344        self.data.state = self.data.flow()?.after_regions();
345        Ok(())
346    }
347
348    async fn submit_sync_region_requests(
349        &mut self,
350        results: &[RegionResponse],
351        region_routes: &[RegionRoute],
352    ) {
353        // Safety: filled in `prepare` step.
354        let table_info = self.data.table_info().unwrap();
355        if let Err(err) = sync_follower_regions(
356            &self.context,
357            self.data.table_id(),
358            results,
359            region_routes,
360            table_info.meta.engine.as_str(),
361        )
362        .await
363        {
364            error!(err; "Failed to sync regions for table {}, table_id: {}", self.data.table_ref(), self.data.table_id());
365        }
366    }
367
368    /// Update table metadata.
369    pub(crate) async fn on_update_metadata(&mut self) -> Result<Status> {
370        let table_id = self.data.table_id();
371        let table_ref = self.data.table_ref();
372        // Safety: filled in `fill_table_info`.
373        let table_info_value = self.data.table_info_value.as_ref().unwrap();
374        // Safety: Checked in `AlterTableProcedure::new`.
375        let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap();
376        let flow = self.data.flow()?;
377        let metadata_only_alter = flow == AlterTableFlow::MetadataOnly;
378
379        // Gets the table info from the cache or builds it.
380        let  new_info = match &self.new_table_info {
381            Some(cached) => cached.clone(),
382            None => AlterTableExecutor::validate_alter_table_expr(
383                &table_info_value.table_info,
384                self.data.task.alter_table.clone(),
385               )
386                .inspect_err(|e| {
387                    // We already check the table info in the prepare step so this should not happen.
388                    error!(e; "Unable to build info for table {} in update metadata step, table_id: {}", table_ref, table_id);
389                })?,
390        };
391
392        // Safety: region distribution is set in `submit_alter_region_requests`.
393        self.executor
394            .on_alter_metadata(
395                &self.context.table_metadata_manager,
396                table_info_value,
397                self.data.region_distribution.as_ref(),
398                new_info,
399                &self.data.column_metadatas,
400                metadata_only_alter,
401            )
402            .await?;
403
404        info!(
405            "Updated table metadata for table {table_ref}, table_id: {table_id}, kind: {alter_kind:?}"
406        );
407        self.data.state = flow.after_metadata();
408        Ok(Status::executing(true))
409    }
410
411    /// Broadcasts the invalidating table cache instructions.
412    async fn on_broadcast(&mut self) -> Result<Status> {
413        self.executor
414            .invalidate_table_cache(&self.context.cache_invalidator)
415            .await?;
416        Ok(Status::done())
417    }
418
419    fn lock_key_inner(&self) -> Vec<StringKey> {
420        let mut lock_key = vec![];
421        let table_ref = self.data.table_ref();
422        let table_id = self.data.table_id();
423        lock_key.push(CatalogLock::Read(table_ref.catalog).into());
424        lock_key.push(SchemaLock::read(table_ref.catalog, table_ref.schema).into());
425        lock_key.push(TableLock::Write(table_id).into());
426
427        for region_id in &self.data.region_locks {
428            lock_key.push(RegionLock::Write(*region_id).into());
429        }
430
431        // Safety: Checked in `AlterTableProcedure::new`.
432        let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap();
433        if let Kind::RenameTable(RenameTable { new_table_name }) = alter_kind {
434            lock_key.push(
435                TableNameLock::new(table_ref.catalog, table_ref.schema, new_table_name).into(),
436            )
437        }
438
439        lock_key
440    }
441
442    #[cfg(test)]
443    pub(crate) fn data(&self) -> &AlterTableData {
444        &self.data
445    }
446
447    #[cfg(test)]
448    pub(crate) fn mut_data(&mut self) -> &mut AlterTableData {
449        &mut self.data
450    }
451}
452
453fn enables_skip_wal(alter_kind: &Kind) -> bool {
454    let Kind::SetTableOptions(SetTableOptions { table_options }) = alter_kind else {
455        return false;
456    };
457
458    table_options
459        .iter()
460        .any(|option| option.key == SKIP_WAL_KEY && option.value == "true")
461}
462
463pub(crate) fn only_enables_skip_wal(alter_kind: &Kind) -> bool {
464    let Kind::SetTableOptions(SetTableOptions { table_options }) = alter_kind else {
465        return false;
466    };
467
468    table_options.len() == 1
469        && table_options[0].key == SKIP_WAL_KEY
470        && table_options[0].value == "true"
471}
472
473fn is_metadata_only_alter(alter_kind: &Kind) -> Result<bool> {
474    // A mixed annotation batch is an error, never "not metadata-only": falling
475    // through to the region-first flow would dispatch it to regions.
476    let family = common_grpc_expr::annotation_alter_family(alter_kind)
477        .context(ConvertAlterTableRequestSnafu)?;
478    Ok(family.is_some() || matches!(alter_kind, Kind::RenameTable { .. }))
479}
480
481#[derive(Debug, Clone, Copy, PartialEq, Eq)]
482enum AlterTableFlow {
483    RegionFirst,
484    MetadataFirst,
485    MetadataOnly,
486}
487
488impl AlterTableFlow {
489    fn from_kind(kind: &Kind) -> Result<Self> {
490        Ok(if only_enables_skip_wal(kind) {
491            Self::MetadataFirst
492        } else if is_metadata_only_alter(kind)? {
493            Self::MetadataOnly
494        } else {
495            Self::RegionFirst
496        })
497    }
498
499    fn after_prepare(self) -> AlterTableState {
500        match self {
501            Self::RegionFirst => AlterTableState::SubmitAlterRegionRequests,
502            Self::MetadataFirst | Self::MetadataOnly => AlterTableState::UpdateMetadata,
503        }
504    }
505
506    fn after_regions(self) -> AlterTableState {
507        match self {
508            Self::RegionFirst => AlterTableState::UpdateMetadata,
509            Self::MetadataFirst | Self::MetadataOnly => AlterTableState::InvalidateTableCache,
510        }
511    }
512
513    fn after_metadata(self) -> AlterTableState {
514        match self {
515            Self::MetadataFirst => AlterTableState::SubmitAlterRegionRequests,
516            Self::RegionFirst | Self::MetadataOnly => AlterTableState::InvalidateTableCache,
517        }
518    }
519}
520
521#[async_trait]
522impl Procedure for AlterTableProcedure {
523    fn type_name(&self) -> &str {
524        Self::TYPE_NAME
525    }
526
527    async fn execute(&mut self, ctx: &ProcedureContext) -> ProcedureResult<Status> {
528        let state = &self.data.state;
529
530        let step = state.as_ref();
531
532        let _timer = metrics::METRIC_META_PROCEDURE_ALTER_TABLE
533            .with_label_values(&[step])
534            .start_timer();
535
536        match state {
537            AlterTableState::Prepare => self.on_prepare().await,
538            AlterTableState::SubmitAlterRegionRequests => {
539                self.submit_alter_region_requests(ctx.procedure_id, ctx.provider.as_ref())
540                    .await
541            }
542            AlterTableState::UpdateMetadata => self.on_update_metadata().await,
543            AlterTableState::InvalidateTableCache => self.on_broadcast().await,
544        }
545        .map_err(map_to_procedure_error)
546    }
547
548    fn dump(&self) -> ProcedureResult<String> {
549        serde_json::to_string(&self.data).context(ToJsonSnafu)
550    }
551
552    fn lock_key(&self) -> LockKey {
553        let key = self.lock_key_inner();
554
555        LockKey::new(key)
556    }
557
558    fn poison_keys(&self) -> PoisonKeys {
559        PoisonKeys::new(vec![self.table_poison_key()])
560    }
561
562    fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
563        if !ctx
564            .event_type_filter
565            .allows(TableDdlEventType::AlterTable.as_str())
566        {
567            return None;
568        }
569        let table_ref = self.data.table_ref();
570        let locator = TableDdlLocator::new(table_ref.catalog, table_ref.schema, table_ref.table)
571            .with_table_id(self.data.table_id());
572        let event = match &ctx.trigger {
573            EventTrigger::Submitted => {
574                let kind = self
575                    .data
576                    .task
577                    .alter_table
578                    .kind
579                    .as_ref()
580                    .and_then(alter_table_kind_name);
581                TableDdlEvent::alter_table_submitted(locator, kind)
582            }
583            _ => TableDdlEvent::lifecycle(TableDdlEventType::AlterTable, [locator]),
584        };
585
586        Some(Box::new(event))
587    }
588}
589
590#[derive(Debug, Serialize, Deserialize, AsRefStr)]
591enum AlterTableState {
592    /// Prepares to alter the table.
593    Prepare,
594    /// Sends alter region requests to Datanode.
595    SubmitAlterRegionRequests,
596    /// Updates table metadata.
597    UpdateMetadata,
598    /// Broadcasts the invalidating table cache instruction.
599    InvalidateTableCache,
600}
601
602// The serialized data of alter table.
603#[derive(Debug, Serialize, Deserialize)]
604pub struct AlterTableData {
605    state: AlterTableState,
606    task: AlterTableTask,
607    table_id: TableId,
608    #[serde(default)]
609    column_metadatas: Vec<ColumnMetadata>,
610    /// Table info value before alteration.
611    table_info_value: Option<DeserializedValueWithBytes<TableInfoValue>>,
612    /// Region distribution for table in case we need to update region options.
613    region_distribution: Option<RegionDistribution>,
614    /// Region locks held by irreversible region-option alters.
615    #[serde(default)]
616    region_locks: Vec<RegionId>,
617}
618
619impl AlterTableData {
620    pub fn new(task: AlterTableTask, table_id: TableId, region_locks: Vec<RegionId>) -> Self {
621        Self {
622            state: AlterTableState::Prepare,
623            task,
624            table_id,
625            column_metadatas: vec![],
626            table_info_value: None,
627            region_distribution: None,
628            region_locks,
629        }
630    }
631
632    fn table_ref(&self) -> TableReference<'_> {
633        self.task.table_ref()
634    }
635
636    fn table_id(&self) -> TableId {
637        self.table_id
638    }
639
640    fn table_info(&self) -> Option<&TableInfo> {
641        self.table_info_value
642            .as_ref()
643            .map(|value| &value.table_info)
644    }
645
646    fn flow(&self) -> Result<AlterTableFlow> {
647        // Safety: Checked in `AlterTableProcedure::new`.
648        AlterTableFlow::from_kind(self.task.alter_table.kind.as_ref().unwrap())
649    }
650
651    #[cfg(test)]
652    pub(crate) fn column_metadatas(&self) -> &[ColumnMetadata] {
653        &self.column_metadatas
654    }
655
656    #[cfg(test)]
657    pub(crate) fn set_column_metadatas(&mut self, column_metadatas: Vec<ColumnMetadata>) {
658        self.column_metadatas = column_metadatas;
659    }
660}