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