common_meta/reconciliation/reconcile_table/
resolve_column_metadata.rs1use async_trait::async_trait;
16use common_procedure::{Context as ProcedureContext, Status};
17use common_telemetry::info;
18use serde::{Deserialize, Serialize};
19use snafu::OptionExt;
20use store_api::metadata::RegionMetadata;
21use strum::AsRefStr;
22
23use crate::error::{self, MissingColumnIdsSnafu, Result};
24use crate::reconciliation::reconcile_table::reconcile_regions::ReconcileRegions;
25use crate::reconciliation::reconcile_table::update_table_info::UpdateTableInfo;
26use crate::reconciliation::reconcile_table::{ReconcileTableContext, State, TableMetadataState};
27use crate::reconciliation::utils::{
28 ResolveColumnMetadataResult, build_column_metadata_from_table_info,
29 check_column_metadatas_consistent, resolve_column_metadatas_with_latest,
30 resolve_column_metadatas_with_metasrv,
31};
32
33#[derive(Debug, Serialize, Deserialize, Clone, Copy, Default, AsRefStr)]
35pub enum ResolveStrategy {
36 #[default]
37 UseLatest,
39
40 UseMetasrv,
42
43 AbortOnConflict,
45}
46
47impl From<api::v1::meta::ResolveStrategy> for ResolveStrategy {
48 fn from(strategy: api::v1::meta::ResolveStrategy) -> Self {
49 match strategy {
50 api::v1::meta::ResolveStrategy::UseMetasrv => Self::UseMetasrv,
51 api::v1::meta::ResolveStrategy::UseLatest => Self::UseLatest,
52 api::v1::meta::ResolveStrategy::AbortOnConflict => Self::AbortOnConflict,
53 }
54 }
55}
56
57#[derive(Debug, Serialize, Deserialize)]
59pub struct ResolveColumnMetadata {
60 strategy: ResolveStrategy,
61 region_metadata: Vec<RegionMetadata>,
62}
63
64impl ResolveColumnMetadata {
65 pub fn new(strategy: ResolveStrategy, region_metadata: Vec<RegionMetadata>) -> Self {
66 Self {
67 strategy,
68 region_metadata,
69 }
70 }
71}
72
73#[async_trait]
74#[typetag::serde]
75impl State for ResolveColumnMetadata {
76 async fn next(
77 &mut self,
78 ctx: &mut ReconcileTableContext,
79 _procedure_ctx: &ProcedureContext,
80 ) -> Result<(Box<dyn State>, Status)> {
81 let table_id = ctx.persistent_ctx.table_id;
82 let table_name = &ctx.persistent_ctx.table_name;
83
84 let table_info_value = ctx
85 .table_metadata_manager
86 .table_info_manager()
87 .get(table_id)
88 .await?
89 .with_context(|| error::TableNotFoundSnafu {
90 table_name: table_name.to_string(),
91 })?;
92 ctx.persistent_ctx.table_info_value = Some(table_info_value);
93
94 if let Some(column_metadatas) = check_column_metadatas_consistent(&self.region_metadata) {
95 let table_info_value = ctx.persistent_ctx.table_info_value.clone().unwrap();
97 info!(
98 "Column metadatas are consistent for table: {}, table_id: {}.",
99 table_name, table_id
100 );
101
102 ctx.volatile_ctx.result_summary.record_resolved_columns(
103 TableMetadataState::Consistent,
104 None,
105 Some(column_metadatas.len()),
106 );
107
108 ctx.mut_metrics().resolve_column_metadata_result =
110 Some(ResolveColumnMetadataResult::Consistent);
111 return Ok((
112 Box::new(UpdateTableInfo::new(table_info_value, column_metadatas)),
113 Status::executing(false),
114 ));
115 };
116
117 ctx.volatile_ctx
118 .result_summary
119 .record_metadata_state(TableMetadataState::Inconsistent);
120
121 match self.strategy {
122 ResolveStrategy::UseMetasrv => {
123 let table_info_value = ctx.persistent_ctx.table_info_value.as_ref().unwrap();
124 let name_to_ids = table_info_value
125 .table_info
126 .name_to_ids()
127 .context(MissingColumnIdsSnafu)?;
128 let column_metadata = build_column_metadata_from_table_info(
129 table_info_value.table_info.meta.schema.column_schemas(),
130 &table_info_value.table_info.meta.primary_key_indices,
131 &name_to_ids,
132 )?;
133
134 let region_ids =
135 resolve_column_metadatas_with_metasrv(&column_metadata, &self.region_metadata)?;
136
137 ctx.volatile_ctx.result_summary.record_resolved_columns(
138 TableMetadataState::Inconsistent,
139 Some(self.strategy),
140 Some(column_metadata.len()),
141 );
142
143 let metrics = ctx.mut_metrics();
145 metrics.resolve_column_metadata_result =
146 Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
147 Ok((
148 Box::new(ReconcileRegions::new(column_metadata, region_ids)),
149 Status::executing(true),
150 ))
151 }
152 ResolveStrategy::UseLatest => {
153 let (column_metadatas, region_ids) =
154 resolve_column_metadatas_with_latest(&self.region_metadata)?;
155
156 ctx.volatile_ctx.result_summary.record_resolved_columns(
157 TableMetadataState::Inconsistent,
158 Some(self.strategy),
159 Some(column_metadatas.len()),
160 );
161
162 let metrics = ctx.mut_metrics();
164 metrics.resolve_column_metadata_result =
165 Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
166 Ok((
167 Box::new(ReconcileRegions::new(column_metadatas, region_ids)),
168 Status::executing(true),
169 ))
170 }
171 ResolveStrategy::AbortOnConflict => {
172 let table_name = table_name.to_string();
173
174 ctx.volatile_ctx
175 .result_summary
176 .record_resolution_strategy(self.strategy);
177
178 let metrics = ctx.mut_metrics();
180 metrics.resolve_column_metadata_result =
181 Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
182 error::ColumnMetadataConflictsSnafu {
183 table_name,
184 table_id,
185 }
186 .fail()
187 }
188 }
189 }
190}