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_reconciliation_column_metadata,
29 check_column_metadatas_consistent, reorder_tag_columns, 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 column_metadatas =
96 reorder_tag_columns(&column_metadatas, &self.region_metadata[0].primary_key)?;
97 let table_info_value = ctx.persistent_ctx.table_info_value.clone().unwrap();
99 info!(
100 "Column metadatas are consistent for table: {}, table_id: {}.",
101 table_name, table_id
102 );
103
104 ctx.volatile_ctx.result_summary.record_resolved_columns(
105 TableMetadataState::Consistent,
106 None,
107 Some(column_metadatas.len()),
108 );
109
110 ctx.mut_metrics().resolve_column_metadata_result =
112 Some(ResolveColumnMetadataResult::Consistent);
113 return Ok((
114 Box::new(UpdateTableInfo::new(table_info_value, column_metadatas)),
115 Status::executing(false),
116 ));
117 };
118
119 ctx.volatile_ctx
120 .result_summary
121 .record_metadata_state(TableMetadataState::Inconsistent);
122
123 match self.strategy {
124 ResolveStrategy::UseMetasrv => {
125 let table_info_value = ctx.persistent_ctx.table_info_value.as_ref().unwrap();
126 let name_to_ids = table_info_value
127 .table_info
128 .name_to_ids()
129 .context(MissingColumnIdsSnafu)?;
130 let column_metadata = build_reconciliation_column_metadata(
131 table_info_value.table_info.meta.schema.column_schemas(),
132 &table_info_value.table_info.meta.primary_key_indices,
133 &name_to_ids,
134 )?;
135
136 let region_ids =
137 resolve_column_metadatas_with_metasrv(&column_metadata, &self.region_metadata)?;
138
139 ctx.volatile_ctx.result_summary.record_resolved_columns(
140 TableMetadataState::Inconsistent,
141 Some(self.strategy),
142 Some(column_metadata.len()),
143 );
144
145 let metrics = ctx.mut_metrics();
147 metrics.resolve_column_metadata_result =
148 Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
149 Ok((
150 Box::new(ReconcileRegions::new(column_metadata, region_ids)),
151 Status::executing(true),
152 ))
153 }
154 ResolveStrategy::UseLatest => {
155 let (column_metadatas, region_ids) =
156 resolve_column_metadatas_with_latest(&self.region_metadata)?;
157
158 ctx.volatile_ctx.result_summary.record_resolved_columns(
159 TableMetadataState::Inconsistent,
160 Some(self.strategy),
161 Some(column_metadatas.len()),
162 );
163
164 let metrics = ctx.mut_metrics();
166 metrics.resolve_column_metadata_result =
167 Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
168 Ok((
169 Box::new(ReconcileRegions::new(column_metadatas, region_ids)),
170 Status::executing(true),
171 ))
172 }
173 ResolveStrategy::AbortOnConflict => {
174 let table_name = table_name.to_string();
175
176 ctx.volatile_ctx
177 .result_summary
178 .record_resolution_strategy(self.strategy);
179
180 let metrics = ctx.mut_metrics();
182 metrics.resolve_column_metadata_result =
183 Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
184 error::ColumnMetadataConflictsSnafu {
185 table_name,
186 table_id,
187 }
188 .fail()
189 }
190 }
191 }
192}