1use std::collections::{HashMap, HashSet};
16use std::fmt::{self, Display};
17use std::ops::AddAssign;
18use std::sync::Arc;
19use std::time::Instant;
20
21use api::v1::SemanticType;
22use common_procedure::{Context as ProcedureContext, ProcedureId, watcher};
23use common_telemetry::{error, warn};
24use datatypes::schema::{ColumnSchema, Schema};
25use futures::future::{join_all, try_join_all};
26use snafu::{OptionExt, ResultExt, ensure};
27use store_api::metadata::{ColumnMetadata, RegionMetadata};
28use store_api::storage::consts::ReservedColumnId;
29use store_api::storage::{RegionId, TableId};
30use table::metadata::{TableInfo, TableMeta};
31use table::table_name::TableName;
32use table::table_reference::TableReference;
33
34use crate::cache_invalidator::CacheInvalidatorRef;
35use crate::error::{
36 ColumnIdMismatchSnafu, ColumnNotFoundSnafu, MismatchColumnIdSnafu,
37 MissingColumnInColumnMetadataSnafu, ProcedureStateReceiverNotFoundSnafu,
38 ProcedureStateReceiverSnafu, Result, TimestampMismatchSnafu, UnexpectedSnafu,
39 WaitProcedureSnafu,
40};
41use crate::key::TableMetadataManagerRef;
42use crate::metrics;
43use crate::node_manager::NodeManagerRef;
44use crate::reconciliation::reconcile_logical_tables::ReconcileLogicalTablesProcedure;
45use crate::reconciliation::reconcile_table::ReconcileTableProcedure;
46use crate::reconciliation::reconcile_table::resolve_column_metadata::ResolveStrategy;
47
48#[derive(Debug, PartialEq, Eq)]
49pub(crate) struct PartialRegionMetadata<'a> {
50 pub(crate) column_metadatas: &'a [ColumnMetadata],
51 pub(crate) primary_key: &'a [u32],
52 pub(crate) table_id: TableId,
53}
54
55impl<'a> From<&'a RegionMetadata> for PartialRegionMetadata<'a> {
56 fn from(region_metadata: &'a RegionMetadata) -> Self {
57 Self {
58 column_metadatas: ®ion_metadata.column_metadatas,
59 primary_key: ®ion_metadata.primary_key,
60 table_id: region_metadata.region_id.table_id(),
61 }
62 }
63}
64
65pub(crate) fn check_column_metadatas_consistent(
75 region_metadatas: &[RegionMetadata],
76) -> Option<Vec<ColumnMetadata>> {
77 let is_column_metadata_consistent = region_metadatas
78 .windows(2)
79 .all(|w| PartialRegionMetadata::from(&w[0]) == PartialRegionMetadata::from(&w[1]));
80
81 if !is_column_metadata_consistent {
82 return None;
83 }
84
85 Some(region_metadatas[0].column_metadatas.clone())
86}
87
88pub(crate) fn reorder_tag_columns(
91 column_metadatas: &[ColumnMetadata],
92 primary_key: &[u32],
93) -> Result<Vec<ColumnMetadata>> {
94 let tag_count = column_metadatas
95 .iter()
96 .filter(|column| column.semantic_type == SemanticType::Tag)
97 .count();
98 ensure!(
99 primary_key.len() == tag_count,
100 UnexpectedSnafu {
101 err_msg: format!(
102 "Number of primary key columns {} does not match tag columns {}",
103 primary_key.len(),
104 tag_count,
105 ),
106 }
107 );
108
109 let mut primary_key_ids = HashSet::with_capacity(primary_key.len());
110 let mut tags = Vec::with_capacity(primary_key.len());
111 for column_id in primary_key {
112 let column = column_metadatas
113 .iter()
114 .find(|column| column.column_id == *column_id)
115 .with_context(|| UnexpectedSnafu {
116 err_msg: format!(
117 "Primary key column {} not found in column metadata",
118 column_id
119 ),
120 })?;
121 ensure!(
122 column.semantic_type == SemanticType::Tag,
123 UnexpectedSnafu {
124 err_msg: format!("Primary key column {} is not a tag", column_id),
125 }
126 );
127 ensure!(
128 primary_key_ids.insert(column_id),
129 UnexpectedSnafu {
130 err_msg: format!("Primary key column {} is duplicated", column_id),
131 }
132 );
133 tags.push(column);
134 }
135
136 let mut tags = tags.into_iter();
137
138 column_metadatas
139 .iter()
140 .map(|column| {
141 if column.semantic_type == SemanticType::Tag {
142 tags.next().cloned().with_context(|| UnexpectedSnafu {
143 err_msg: "Primary key has fewer columns than tags".to_string(),
144 })
145 } else {
146 Ok(column.clone())
147 }
148 })
149 .collect()
150}
151
152pub(crate) fn resolve_column_metadatas_with_metasrv(
160 column_metadatas: &[ColumnMetadata],
161 region_metadatas: &[RegionMetadata],
162) -> Result<Vec<RegionId>> {
163 let is_same_table = region_metadatas
164 .windows(2)
165 .all(|w| w[0].region_id.table_id() == w[1].region_id.table_id());
166
167 ensure!(
168 is_same_table,
169 UnexpectedSnafu {
170 err_msg: "Region metadatas are not from the same table"
171 }
172 );
173
174 let mut regions_ids = vec![];
175 for region_metadata in region_metadatas {
176 if region_metadata.column_metadatas != column_metadatas {
177 check_column_metadata_invariants(column_metadatas, ®ion_metadata.column_metadatas)?;
178 regions_ids.push(region_metadata.region_id);
179 }
180 }
181 Ok(regions_ids)
182}
183
184pub(crate) fn resolve_column_metadatas_with_latest(
192 region_metadatas: &[RegionMetadata],
193) -> Result<(Vec<ColumnMetadata>, Vec<RegionId>)> {
194 let is_same_table = region_metadatas
195 .windows(2)
196 .all(|w| w[0].region_id.table_id() == w[1].region_id.table_id());
197
198 ensure!(
199 is_same_table,
200 UnexpectedSnafu {
201 err_msg: "Region metadatas are not from the same table"
202 }
203 );
204
205 let latest_region_metadata = region_metadatas
206 .iter()
207 .max_by_key(|c| c.schema_version)
208 .context(UnexpectedSnafu {
209 err_msg: "All Region metadatas have the same schema version",
210 })?;
211 let latest_column_metadatas = PartialRegionMetadata::from(latest_region_metadata);
212
213 let mut region_ids = vec![];
214 for region_metadata in region_metadatas {
215 if PartialRegionMetadata::from(region_metadata) != latest_column_metadatas {
216 check_column_metadata_invariants(
217 &latest_region_metadata.column_metadatas,
218 ®ion_metadata.column_metadatas,
219 )?;
220 region_ids.push(region_metadata.region_id);
221 }
222 }
223
224 Ok((
225 reorder_tag_columns(
226 &latest_region_metadata.column_metadatas,
227 &latest_region_metadata.primary_key,
228 )?,
229 region_ids,
230 ))
231}
232
233pub(crate) fn build_column_metadata_from_table_info(
241 column_schemas: &[ColumnSchema],
242 primary_key_indexes: &[usize],
243 name_to_ids: &HashMap<String, u32>,
244) -> Result<Vec<ColumnMetadata>> {
245 let primary_names = primary_key_indexes
246 .iter()
247 .map(|i| column_schemas[*i].name.as_str())
248 .collect::<HashSet<_>>();
249
250 column_schemas
251 .iter()
252 .map(|column_schema| {
253 let column_id = *name_to_ids
254 .get(column_schema.name.as_str())
255 .with_context(|| UnexpectedSnafu {
256 err_msg: format!(
257 "Column name {} not found in name_to_ids",
258 column_schema.name
259 ),
260 })?;
261
262 let semantic_type = if primary_names.contains(&column_schema.name.as_str()) {
263 SemanticType::Tag
264 } else if column_schema.is_time_index() {
265 SemanticType::Timestamp
266 } else {
267 SemanticType::Field
268 };
269 Ok(ColumnMetadata {
270 column_schema: column_schema.clone(),
271 semantic_type,
272 column_id,
273 })
274 })
275 .collect::<Result<Vec<_>>>()
276}
277
278pub(crate) fn build_reconciliation_column_metadata(
281 column_schemas: &[ColumnSchema],
282 primary_key_indexes: &[usize],
283 name_to_ids: &HashMap<String, u32>,
284) -> Result<Vec<ColumnMetadata>> {
285 let column_metadatas =
286 build_column_metadata_from_table_info(column_schemas, primary_key_indexes, name_to_ids)?;
287 let primary_key = primary_key_indexes
288 .iter()
289 .map(|index| column_metadatas[*index].column_id)
290 .collect::<Vec<_>>();
291 reorder_tag_columns(&column_metadatas, &primary_key)
292}
293
294pub(crate) fn check_column_metadata_invariants(
300 new_column_metadatas: &[ColumnMetadata],
301 column_metadatas: &[ColumnMetadata],
302) -> Result<()> {
303 let new_primary_keys = new_column_metadatas
304 .iter()
305 .filter(|c| c.semantic_type == SemanticType::Tag)
306 .map(|c| (c.column_schema.name.as_str(), c.column_id))
307 .collect::<HashMap<_, _>>();
308
309 let old_primary_keys = column_metadatas
310 .iter()
311 .filter(|c| c.semantic_type == SemanticType::Tag)
312 .map(|c| (c.column_schema.name.as_str(), c.column_id));
313
314 for (name, id) in old_primary_keys {
315 let column_id = new_primary_keys
316 .get(name)
317 .cloned()
318 .context(ColumnNotFoundSnafu {
319 column_name: name,
320 column_id: id,
321 })?;
322
323 ensure!(
324 column_id == id,
325 ColumnIdMismatchSnafu {
326 column_name: name,
327 expected_column_id: id,
328 actual_column_id: column_id,
329 }
330 );
331 }
332
333 let new_ts_column = new_column_metadatas
334 .iter()
335 .find(|c| c.semantic_type == SemanticType::Timestamp)
336 .map(|c| (c.column_schema.name.as_str(), c.column_id))
337 .context(UnexpectedSnafu {
338 err_msg: "Timestamp column not found in new column metadata",
339 })?;
340
341 let old_ts_column = column_metadatas
342 .iter()
343 .find(|c| c.semantic_type == SemanticType::Timestamp)
344 .map(|c| (c.column_schema.name.as_str(), c.column_id))
345 .context(UnexpectedSnafu {
346 err_msg: "Timestamp column not found in column metadata",
347 })?;
348 ensure!(
349 new_ts_column == old_ts_column,
350 TimestampMismatchSnafu {
351 expected_column_name: old_ts_column.0,
352 expected_column_id: old_ts_column.1,
353 actual_column_name: new_ts_column.0,
354 actual_column_id: new_ts_column.1,
355 }
356 );
357
358 Ok(())
359}
360
361pub(crate) fn build_table_meta_from_column_metadatas(
371 table_id: TableId,
372 table_ref: TableReference,
373 table_meta: &TableMeta,
374 name_to_ids: Option<HashMap<String, u32>>,
375 column_metadata: &[ColumnMetadata],
376) -> Result<TableMeta> {
377 let column_in_column_metadata = column_metadata
378 .iter()
379 .map(|c| (c.column_schema.name.as_str(), c))
380 .collect::<HashMap<_, _>>();
381 let column_schemas = table_meta.schema.column_schemas();
382 let primary_key_names = table_meta
383 .primary_key_indices
384 .iter()
385 .map(|i| column_schemas[*i].name.as_str())
386 .collect::<HashSet<_>>();
387 let partition_key_names = table_meta
388 .partition_key_indices
389 .iter()
390 .map(|i| column_schemas[*i].name.as_str())
391 .collect::<HashSet<_>>();
392 ensure!(
393 column_metadata
394 .iter()
395 .any(|c| c.semantic_type == SemanticType::Timestamp),
396 UnexpectedSnafu {
397 err_msg: format!(
398 "Missing table index in column metadata, table: {}, table_id: {}",
399 table_ref, table_id
400 ),
401 }
402 );
403
404 if let Some(name_to_ids) = &name_to_ids {
405 for column_name in primary_key_names.iter().chain(partition_key_names.iter()) {
407 let column_in_column_metadata = column_in_column_metadata
408 .get(column_name)
409 .with_context(|| MissingColumnInColumnMetadataSnafu {
410 column_name: column_name.to_string(),
411 table_name: table_ref.to_string(),
412 table_id,
413 })?;
414
415 let column_id = *name_to_ids
416 .get(*column_name)
417 .with_context(|| UnexpectedSnafu {
418 err_msg: format!("column id not found in name_to_ids: {}", column_name),
419 })?;
420 ensure!(
421 column_id == column_in_column_metadata.column_id,
422 MismatchColumnIdSnafu {
423 column_name: column_name.to_string(),
424 column_id,
425 table_name: table_ref.to_string(),
426 table_id,
427 }
428 );
429 }
430 } else {
431 warn!(
432 "`name_to_ids` is not provided, table: {}, table_id: {}",
433 table_ref, table_id
434 );
435 }
436
437 let mut new_raw_table_meta = table_meta.clone();
438 let primary_key_indices = &mut new_raw_table_meta.primary_key_indices;
439 let partition_key_indices = &mut new_raw_table_meta.partition_key_indices;
440 let value_indices = &mut new_raw_table_meta.value_indices;
441 let mut columns = Vec::with_capacity(column_metadata.len());
442 let column_ids = &mut new_raw_table_meta.column_ids;
443 let next_column_id = &mut new_raw_table_meta.next_column_id;
444
445 column_ids.clear();
446 value_indices.clear();
447 primary_key_indices.clear();
448 partition_key_indices.clear();
449
450 for (idx, col) in column_metadata.iter().enumerate() {
451 if partition_key_names.contains(&col.column_schema.name.as_str()) {
452 partition_key_indices.push(idx);
453 }
454 match col.semantic_type {
455 SemanticType::Tag => {
456 primary_key_indices.push(idx);
457 }
458 SemanticType::Field => {
459 value_indices.push(idx);
460 }
461 SemanticType::Timestamp => {
462 value_indices.push(idx);
463 }
464 }
465
466 columns.push(col.column_schema.clone());
467 column_ids.push(col.column_id);
468 }
469
470 *next_column_id = column_ids
471 .iter()
472 .filter(|id| !ReservedColumnId::is_reserved(**id))
473 .max()
474 .map(|max| max + 1)
475 .unwrap_or(*next_column_id)
476 .max(*next_column_id);
477
478 new_raw_table_meta.schema = Arc::new(Schema::new_with_version(
479 columns,
480 table_meta.schema.version(),
481 ));
482 Ok(new_raw_table_meta)
483}
484
485pub(crate) fn need_update_logical_table_info(
490 table_info: &TableInfo,
491 column_metadatas: &[ColumnMetadata],
492) -> bool {
493 table_info.meta.schema.column_schemas().len() != column_metadatas.len()
494}
495
496pub struct PartialSuccessResult<'a> {
498 pub failed_procedures: Vec<&'a SubprocedureMeta>,
499 pub success_procedures: Vec<&'a SubprocedureMeta>,
500}
501
502pub enum WaitForInflightSubproceduresResult<'a> {
504 Success(Vec<&'a SubprocedureMeta>),
505 PartialSuccess(PartialSuccessResult<'a>),
506}
507
508pub(crate) async fn wait_for_inflight_subprocedures<'a>(
513 procedure_ctx: &ProcedureContext,
514 subprocedures: &'a [SubprocedureMeta],
515 fail_fast: bool,
516) -> Result<WaitForInflightSubproceduresResult<'a>> {
517 let mut receivers = Vec::with_capacity(subprocedures.len());
518 for subprocedure in subprocedures {
519 let procedure_id = subprocedure.procedure_id();
520 let receiver = procedure_ctx
521 .provider
522 .procedure_state_receiver(procedure_id)
523 .await
524 .context(ProcedureStateReceiverSnafu { procedure_id })?
525 .context(ProcedureStateReceiverNotFoundSnafu { procedure_id })?;
526 receivers.push((receiver, subprocedure));
527 }
528
529 let mut tasks = Vec::with_capacity(receivers.len());
530 for (receiver, subprocedure) in receivers.iter_mut() {
531 tasks.push(async move {
532 watcher::wait(receiver).await.inspect_err(|e| {
533 error!(e; "inflight subprocedure failed, parent procedure_id: {}, procedure: {}", procedure_ctx.procedure_id, subprocedure);
534 })
535 });
536 }
537
538 if fail_fast {
539 try_join_all(tasks).await.context(WaitProcedureSnafu)?;
540 return Ok(WaitForInflightSubproceduresResult::Success(
541 subprocedures.iter().collect(),
542 ));
543 }
544
545 let results = join_all(tasks).await;
547 let failed_procedures_num = results.iter().filter(|r| r.is_err()).count();
548 if failed_procedures_num == 0 {
549 return Ok(WaitForInflightSubproceduresResult::Success(
550 subprocedures.iter().collect(),
551 ));
552 }
553 warn!(
554 "{} inflight subprocedures failed, total: {}, parent procedure_id: {}",
555 failed_procedures_num,
556 subprocedures.len(),
557 procedure_ctx.procedure_id
558 );
559
560 let mut failed_procedures = Vec::with_capacity(failed_procedures_num);
561 let mut success_procedures = Vec::with_capacity(subprocedures.len() - failed_procedures_num);
562 for (result, subprocedure) in results.into_iter().zip(subprocedures) {
563 if result.is_err() {
564 failed_procedures.push(subprocedure);
565 } else {
566 success_procedures.push(subprocedure);
567 }
568 }
569
570 Ok(WaitForInflightSubproceduresResult::PartialSuccess(
571 PartialSuccessResult {
572 failed_procedures,
573 success_procedures,
574 },
575 ))
576}
577
578#[derive(Clone)]
579pub struct Context {
580 pub node_manager: NodeManagerRef,
581 pub table_metadata_manager: TableMetadataManagerRef,
582 pub cache_invalidator: CacheInvalidatorRef,
583}
584
585pub struct PhysicalTableMeta {
587 pub procedure_id: ProcedureId,
588 pub table_id: TableId,
589 pub table_name: TableName,
590}
591
592pub struct LogicalTableMeta {
594 pub procedure_id: ProcedureId,
595 pub physical_table_id: TableId,
596 pub physical_table_name: TableName,
597 pub logical_tables: Vec<(TableId, TableName)>,
598}
599
600pub struct ReconcileDatabaseMeta {
602 pub procedure_id: ProcedureId,
603 pub catalog: String,
604 pub schema: String,
605}
606
607pub enum SubprocedureMeta {
609 PhysicalTable(PhysicalTableMeta),
610 LogicalTable(LogicalTableMeta),
611 Database(ReconcileDatabaseMeta),
612}
613
614impl Display for SubprocedureMeta {
615 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
616 match self {
617 SubprocedureMeta::PhysicalTable(meta) => {
618 write!(
619 f,
620 "ReconcilePhysicalTable(procedure_id: {}, table_id: {}, table_name: {})",
621 meta.procedure_id, meta.table_id, meta.table_name
622 )
623 }
624 SubprocedureMeta::LogicalTable(meta) => {
625 write!(
626 f,
627 "ReconcileLogicalTable(procedure_id: {}, physical_table_id: {}, physical_table_name: {}, logical_tables: {:?})",
628 meta.procedure_id,
629 meta.physical_table_id,
630 meta.physical_table_name,
631 meta.logical_tables
632 )
633 }
634 SubprocedureMeta::Database(meta) => {
635 write!(
636 f,
637 "ReconcileDatabase(procedure_id: {}, catalog: {}, schema: {})",
638 meta.procedure_id, meta.catalog, meta.schema
639 )
640 }
641 }
642 }
643}
644
645impl SubprocedureMeta {
646 pub fn new_logical_table(
648 procedure_id: ProcedureId,
649 physical_table_id: TableId,
650 physical_table_name: TableName,
651 logical_tables: Vec<(TableId, TableName)>,
652 ) -> Self {
653 Self::LogicalTable(LogicalTableMeta {
654 procedure_id,
655 physical_table_id,
656 physical_table_name,
657 logical_tables,
658 })
659 }
660
661 pub fn new_physical_table(
663 procedure_id: ProcedureId,
664 table_id: TableId,
665 table_name: TableName,
666 ) -> Self {
667 Self::PhysicalTable(PhysicalTableMeta {
668 procedure_id,
669 table_id,
670 table_name,
671 })
672 }
673
674 pub fn new_reconcile_database(
676 procedure_id: ProcedureId,
677 catalog: String,
678 schema: String,
679 ) -> Self {
680 Self::Database(ReconcileDatabaseMeta {
681 procedure_id,
682 catalog,
683 schema,
684 })
685 }
686
687 pub fn procedure_id(&self) -> ProcedureId {
689 match self {
690 SubprocedureMeta::PhysicalTable(meta) => meta.procedure_id,
691 SubprocedureMeta::LogicalTable(meta) => meta.procedure_id,
692 SubprocedureMeta::Database(meta) => meta.procedure_id,
693 }
694 }
695
696 pub fn table_num(&self) -> usize {
698 match self {
699 SubprocedureMeta::PhysicalTable(_) => 1,
700 SubprocedureMeta::LogicalTable(meta) => meta.logical_tables.len(),
701 SubprocedureMeta::Database(_) => 0,
702 }
703 }
704
705 pub fn database_num(&self) -> usize {
707 match self {
708 SubprocedureMeta::Database(_) => 1,
709 _ => 0,
710 }
711 }
712}
713
714#[derive(Clone, Default)]
716pub struct ReconcileCatalogMetrics {
717 pub succeeded_databases: usize,
718 pub failed_databases: usize,
719}
720
721impl AddAssign for ReconcileCatalogMetrics {
722 fn add_assign(&mut self, other: Self) {
723 self.succeeded_databases += other.succeeded_databases;
724 self.failed_databases += other.failed_databases;
725 }
726}
727
728impl Display for ReconcileCatalogMetrics {
729 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
730 write!(
731 f,
732 "succeeded_databases: {}, failed_databases: {}",
733 self.succeeded_databases, self.failed_databases
734 )
735 }
736}
737
738impl From<WaitForInflightSubproceduresResult<'_>> for ReconcileCatalogMetrics {
739 fn from(result: WaitForInflightSubproceduresResult<'_>) -> Self {
740 match result {
741 WaitForInflightSubproceduresResult::Success(subprocedures) => ReconcileCatalogMetrics {
742 succeeded_databases: subprocedures.len(),
743 failed_databases: 0,
744 },
745 WaitForInflightSubproceduresResult::PartialSuccess(PartialSuccessResult {
746 failed_procedures,
747 success_procedures,
748 }) => {
749 let succeeded_databases = success_procedures
750 .iter()
751 .map(|subprocedure| subprocedure.database_num())
752 .sum();
753 let failed_databases = failed_procedures
754 .iter()
755 .map(|subprocedure| subprocedure.database_num())
756 .sum();
757 ReconcileCatalogMetrics {
758 succeeded_databases,
759 failed_databases,
760 }
761 }
762 }
763 }
764}
765
766#[derive(Clone, Default)]
768pub struct ReconcileDatabaseMetrics {
769 pub succeeded_tables: usize,
770 pub failed_tables: usize,
771 pub succeeded_procedures: usize,
772 pub failed_procedures: usize,
773}
774
775impl Display for ReconcileDatabaseMetrics {
776 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
777 write!(
778 f,
779 "succeeded_tables: {}, failed_tables: {}, succeeded_procedures: {}, failed_procedures: {}",
780 self.succeeded_tables,
781 self.failed_tables,
782 self.succeeded_procedures,
783 self.failed_procedures
784 )
785 }
786}
787
788impl AddAssign for ReconcileDatabaseMetrics {
789 fn add_assign(&mut self, other: Self) {
790 self.succeeded_tables += other.succeeded_tables;
791 self.failed_tables += other.failed_tables;
792 self.succeeded_procedures += other.succeeded_procedures;
793 self.failed_procedures += other.failed_procedures;
794 }
795}
796
797impl From<WaitForInflightSubproceduresResult<'_>> for ReconcileDatabaseMetrics {
798 fn from(result: WaitForInflightSubproceduresResult<'_>) -> Self {
799 match result {
800 WaitForInflightSubproceduresResult::Success(subprocedures) => {
801 let table_num = subprocedures
802 .iter()
803 .map(|subprocedure| subprocedure.table_num())
804 .sum();
805 ReconcileDatabaseMetrics {
806 succeeded_procedures: subprocedures.len(),
807 failed_procedures: 0,
808 succeeded_tables: table_num,
809 failed_tables: 0,
810 }
811 }
812 WaitForInflightSubproceduresResult::PartialSuccess(PartialSuccessResult {
813 failed_procedures,
814 success_procedures,
815 }) => {
816 let succeeded_tables = success_procedures
817 .iter()
818 .map(|subprocedure| subprocedure.table_num())
819 .sum();
820 let failed_tables = failed_procedures
821 .iter()
822 .map(|subprocedure| subprocedure.table_num())
823 .sum();
824 ReconcileDatabaseMetrics {
825 succeeded_procedures: success_procedures.len(),
826 failed_procedures: failed_procedures.len(),
827 succeeded_tables,
828 failed_tables,
829 }
830 }
831 }
832 }
833}
834
835#[derive(Clone)]
837pub struct ReconcileLogicalTableMetrics {
838 pub start_time: Instant,
839 pub update_table_info_count: usize,
840 pub create_tables_count: usize,
841 pub column_metadata_consistent_count: usize,
842 pub column_metadata_inconsistent_count: usize,
843}
844
845impl Default for ReconcileLogicalTableMetrics {
846 fn default() -> Self {
847 Self {
848 start_time: Instant::now(),
849 update_table_info_count: 0,
850 create_tables_count: 0,
851 column_metadata_consistent_count: 0,
852 column_metadata_inconsistent_count: 0,
853 }
854 }
855}
856
857const CREATE_TABLES: &str = "create_tables";
858const UPDATE_TABLE_INFO: &str = "update_table_info";
859const COLUMN_METADATA_CONSISTENT: &str = "column_metadata_consistent";
860const COLUMN_METADATA_INCONSISTENT: &str = "column_metadata_inconsistent";
861
862impl ReconcileLogicalTableMetrics {
863 pub fn total_table_count(&self) -> usize {
865 self.create_tables_count
866 + self.column_metadata_consistent_count
867 + self.column_metadata_inconsistent_count
868 }
869}
870
871impl Drop for ReconcileLogicalTableMetrics {
872 fn drop(&mut self) {
873 let procedure_name = ReconcileLogicalTablesProcedure::TYPE_NAME;
874 metrics::METRIC_META_RECONCILIATION_STATS
875 .with_label_values(&[procedure_name, metrics::TABLE_TYPE_LOGICAL, CREATE_TABLES])
876 .inc_by(self.create_tables_count as u64);
877 metrics::METRIC_META_RECONCILIATION_STATS
878 .with_label_values(&[
879 procedure_name,
880 metrics::TABLE_TYPE_LOGICAL,
881 UPDATE_TABLE_INFO,
882 ])
883 .inc_by(self.update_table_info_count as u64);
884 metrics::METRIC_META_RECONCILIATION_STATS
885 .with_label_values(&[
886 procedure_name,
887 metrics::TABLE_TYPE_LOGICAL,
888 COLUMN_METADATA_CONSISTENT,
889 ])
890 .inc_by(self.column_metadata_consistent_count as u64);
891 metrics::METRIC_META_RECONCILIATION_STATS
892 .with_label_values(&[
893 procedure_name,
894 metrics::TABLE_TYPE_LOGICAL,
895 COLUMN_METADATA_INCONSISTENT,
896 ])
897 .inc_by(self.column_metadata_inconsistent_count as u64);
898 }
899}
900
901impl Display for ReconcileLogicalTableMetrics {
902 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
903 let elapsed = self.start_time.elapsed();
904 if self.create_tables_count > 0 {
905 write!(f, "create_tables_count: {}, ", self.create_tables_count)?;
906 }
907 if self.update_table_info_count > 0 {
908 write!(
909 f,
910 "update_table_info_count: {}, ",
911 self.update_table_info_count
912 )?;
913 }
914 if self.column_metadata_consistent_count > 0 {
915 write!(
916 f,
917 "column_metadata_consistent_count: {}, ",
918 self.column_metadata_consistent_count
919 )?;
920 }
921 if self.column_metadata_inconsistent_count > 0 {
922 write!(
923 f,
924 "column_metadata_inconsistent_count: {}, ",
925 self.column_metadata_inconsistent_count
926 )?;
927 }
928
929 write!(
930 f,
931 "total_table_count: {}, elapsed: {:?}",
932 self.total_table_count(),
933 elapsed
934 )
935 }
936}
937
938#[derive(Clone, Copy)]
940pub enum ResolveColumnMetadataResult {
941 Consistent,
942 Inconsistent(ResolveStrategy),
943}
944
945impl Display for ResolveColumnMetadataResult {
946 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
947 match self {
948 ResolveColumnMetadataResult::Consistent => write!(f, "Consistent"),
949 ResolveColumnMetadataResult::Inconsistent(strategy) => {
950 let strategy_str = strategy.as_ref();
951 write!(f, "Inconsistent({})", strategy_str)
952 }
953 }
954 }
955}
956
957#[derive(Clone)]
959pub struct ReconcileTableMetrics {
960 pub start_time: Instant,
962 pub resolve_column_metadata_result: Option<ResolveColumnMetadataResult>,
964 pub update_table_info: bool,
966}
967
968impl Drop for ReconcileTableMetrics {
969 fn drop(&mut self) {
970 if let Some(resolve_column_metadata_result) = self.resolve_column_metadata_result {
971 match resolve_column_metadata_result {
972 ResolveColumnMetadataResult::Consistent => {
973 metrics::METRIC_META_RECONCILIATION_STATS
974 .with_label_values(&[
975 ReconcileTableProcedure::TYPE_NAME,
976 metrics::TABLE_TYPE_PHYSICAL,
977 COLUMN_METADATA_CONSISTENT,
978 ])
979 .inc();
980 }
981 ResolveColumnMetadataResult::Inconsistent(strategy) => {
982 metrics::METRIC_META_RECONCILIATION_STATS
983 .with_label_values(&[
984 ReconcileTableProcedure::TYPE_NAME,
985 metrics::TABLE_TYPE_PHYSICAL,
986 COLUMN_METADATA_INCONSISTENT,
987 ])
988 .inc();
989 metrics::METRIC_META_RECONCILIATION_RESOLVED_COLUMN_METADATA
990 .with_label_values(&[strategy.as_ref()])
991 .inc();
992 }
993 }
994 }
995 if self.update_table_info {
996 metrics::METRIC_META_RECONCILIATION_STATS
997 .with_label_values(&[
998 ReconcileTableProcedure::TYPE_NAME,
999 metrics::TABLE_TYPE_PHYSICAL,
1000 UPDATE_TABLE_INFO,
1001 ])
1002 .inc();
1003 }
1004 }
1005}
1006
1007impl Default for ReconcileTableMetrics {
1008 fn default() -> Self {
1009 Self {
1010 start_time: Instant::now(),
1011 resolve_column_metadata_result: None,
1012 update_table_info: false,
1013 }
1014 }
1015}
1016
1017impl Display for ReconcileTableMetrics {
1018 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1019 let elapsed = self.start_time.elapsed();
1020 if let Some(resolve_column_metadata_result) = self.resolve_column_metadata_result {
1021 write!(
1022 f,
1023 "resolve_column_metadata_result: {}, ",
1024 resolve_column_metadata_result
1025 )?;
1026 }
1027 write!(
1028 f,
1029 "update_table_info: {}, elapsed: {:?}",
1030 self.update_table_info, elapsed
1031 )
1032 }
1033}
1034
1035#[cfg(test)]
1036mod tests {
1037 use std::assert_matches;
1038 use std::collections::HashMap;
1039 use std::sync::Arc;
1040
1041 use api::v1::SemanticType;
1042 use datatypes::prelude::ConcreteDataType;
1043 use datatypes::schema::{ColumnSchema, Schema, SchemaBuilder};
1044 use store_api::metadata::ColumnMetadata;
1045 use store_api::storage::RegionId;
1046 use table::metadata::TableMetaBuilder;
1047 use table::table_reference::TableReference;
1048
1049 use super::*;
1050 use crate::ddl::test_util::region_metadata::build_region_metadata;
1051 use crate::error::Error;
1052 use crate::reconciliation::utils::check_column_metadatas_consistent;
1053
1054 fn new_test_schema() -> Schema {
1055 let column_schemas = vec![
1056 ColumnSchema::new("col1", ConcreteDataType::int32_datatype(), true),
1057 ColumnSchema::new(
1058 "ts",
1059 ConcreteDataType::timestamp_millisecond_datatype(),
1060 false,
1061 )
1062 .with_time_index(true),
1063 ColumnSchema::new("col2", ConcreteDataType::int32_datatype(), true),
1064 ];
1065 SchemaBuilder::try_from(column_schemas)
1066 .unwrap()
1067 .version(123)
1068 .build()
1069 .unwrap()
1070 }
1071
1072 fn new_test_column_metadatas() -> Vec<ColumnMetadata> {
1073 vec![
1074 ColumnMetadata {
1075 column_schema: ColumnSchema::new("col1", ConcreteDataType::int32_datatype(), true),
1076 semantic_type: SemanticType::Tag,
1077 column_id: 0,
1078 },
1079 ColumnMetadata {
1080 column_schema: ColumnSchema::new(
1081 "ts",
1082 ConcreteDataType::timestamp_millisecond_datatype(),
1083 false,
1084 )
1085 .with_time_index(true),
1086 semantic_type: SemanticType::Timestamp,
1087 column_id: 1,
1088 },
1089 ColumnMetadata {
1090 column_schema: ColumnSchema::new("col2", ConcreteDataType::int32_datatype(), true),
1091 semantic_type: SemanticType::Field,
1092 column_id: 2,
1093 },
1094 ]
1095 }
1096
1097 fn new_test_raw_table_info() -> TableMeta {
1098 let mut table_meta_builder = TableMetaBuilder::empty();
1099 table_meta_builder
1100 .schema(Arc::new(new_test_schema()))
1101 .primary_key_indices(vec![0])
1102 .partition_key_indices(vec![2])
1103 .next_column_id(4)
1104 .build()
1105 .unwrap()
1106 }
1107
1108 #[test]
1109 fn test_build_table_info_from_column_metadatas_identical() {
1110 let column_metadatas = new_test_column_metadatas();
1111 let table_id = 1;
1112 let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1113 let mut table_meta = new_test_raw_table_info();
1114 table_meta.column_ids = vec![0, 1, 2];
1115 let name_to_ids = HashMap::from([
1116 ("col1".to_string(), 0),
1117 ("ts".to_string(), 1),
1118 ("col2".to_string(), 2),
1119 ]);
1120
1121 let new_table_meta = build_table_meta_from_column_metadatas(
1122 table_id,
1123 table_ref,
1124 &table_meta,
1125 Some(name_to_ids),
1126 &column_metadatas,
1127 )
1128 .unwrap();
1129 assert_eq!(new_table_meta, table_meta);
1130 }
1131
1132 #[test]
1133 fn test_build_table_info_from_column_metadatas() {
1134 let mut column_metadatas = new_test_column_metadatas();
1135 column_metadatas.push(ColumnMetadata {
1136 column_schema: ColumnSchema::new(
1137 "__table_id",
1138 ConcreteDataType::string_datatype(),
1139 true,
1140 ),
1141 semantic_type: SemanticType::Tag,
1142 column_id: ReservedColumnId::table_id(),
1143 });
1144
1145 let table_id = 1;
1146 let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1147 let table_meta = new_test_raw_table_info();
1148 let name_to_ids = HashMap::from([
1149 ("col1".to_string(), 0),
1150 ("ts".to_string(), 1),
1151 ("col2".to_string(), 2),
1152 ]);
1153
1154 let new_table_meta = build_table_meta_from_column_metadatas(
1155 table_id,
1156 table_ref,
1157 &table_meta,
1158 Some(name_to_ids),
1159 &column_metadatas,
1160 )
1161 .unwrap();
1162
1163 assert_eq!(new_table_meta.primary_key_indices, vec![0, 3]);
1164 assert_eq!(new_table_meta.partition_key_indices, vec![2]);
1165 assert_eq!(new_table_meta.value_indices, vec![1, 2]);
1166 assert_eq!(new_table_meta.schema.timestamp_index(), Some(1));
1167 assert_eq!(
1168 new_table_meta.column_ids,
1169 vec![0, 1, 2, ReservedColumnId::table_id()]
1170 );
1171 assert_eq!(new_table_meta.next_column_id, table_meta.next_column_id);
1172 }
1173
1174 #[test]
1175 fn test_build_table_info_from_column_metadatas_with_incorrect_name_to_ids() {
1176 let column_metadatas = new_test_column_metadatas();
1177 let table_id = 1;
1178 let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1179 let table_meta = new_test_raw_table_info();
1180 let name_to_ids = HashMap::from([
1181 ("col1".to_string(), 0),
1182 ("ts".to_string(), 1),
1183 ("col2".to_string(), 3),
1185 ]);
1186
1187 let err = build_table_meta_from_column_metadatas(
1188 table_id,
1189 table_ref,
1190 &table_meta,
1191 Some(name_to_ids),
1192 &column_metadatas,
1193 )
1194 .unwrap_err();
1195
1196 assert_matches!(err, Error::MismatchColumnId { .. });
1197 }
1198
1199 #[test]
1200 fn test_build_table_info_from_column_metadatas_with_missing_time_index() {
1201 let mut column_metadatas = new_test_column_metadatas();
1202 column_metadatas.retain(|c| c.semantic_type != SemanticType::Timestamp);
1203 let table_id = 1;
1204 let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1205 let table_meta = new_test_raw_table_info();
1206 let name_to_ids = HashMap::from([
1207 ("col1".to_string(), 0),
1208 ("ts".to_string(), 1),
1209 ("col2".to_string(), 2),
1210 ]);
1211
1212 let err = build_table_meta_from_column_metadatas(
1213 table_id,
1214 table_ref,
1215 &table_meta,
1216 Some(name_to_ids),
1217 &column_metadatas,
1218 )
1219 .unwrap_err();
1220
1221 assert!(
1222 err.to_string()
1223 .contains("Missing table index in column metadata"),
1224 "err: {}",
1225 err
1226 );
1227 }
1228
1229 #[test]
1230 fn test_build_table_info_from_column_metadatas_with_missing_column() {
1231 let mut column_metadatas = new_test_column_metadatas();
1232 column_metadatas.retain(|c| c.column_id != 0);
1234 let table_id = 1;
1235 let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1236 let table_meta = new_test_raw_table_info();
1237 let name_to_ids = HashMap::from([
1238 ("col1".to_string(), 0),
1239 ("ts".to_string(), 1),
1240 ("col2".to_string(), 2),
1241 ]);
1242
1243 let err = build_table_meta_from_column_metadatas(
1244 table_id,
1245 table_ref,
1246 &table_meta,
1247 Some(name_to_ids.clone()),
1248 &column_metadatas,
1249 )
1250 .unwrap_err();
1251 assert_matches!(err, Error::MissingColumnInColumnMetadata { .. });
1252
1253 let mut column_metadatas = new_test_column_metadatas();
1254 column_metadatas.retain(|c| c.column_id != 2);
1256
1257 let err = build_table_meta_from_column_metadatas(
1258 table_id,
1259 table_ref,
1260 &table_meta,
1261 Some(name_to_ids),
1262 &column_metadatas,
1263 )
1264 .unwrap_err();
1265 assert_matches!(err, Error::MissingColumnInColumnMetadata { .. });
1266 }
1267
1268 #[test]
1269 fn test_reconciliation_columns_follow_primary_key_order() {
1270 let columns = vec![
1271 ColumnMetadata {
1272 column_schema: ColumnSchema::new(
1273 "tag_b",
1274 ConcreteDataType::string_datatype(),
1275 true,
1276 ),
1277 semantic_type: SemanticType::Tag,
1278 column_id: 1,
1279 },
1280 ColumnMetadata {
1281 column_schema: ColumnSchema::new("field", ConcreteDataType::int32_datatype(), true),
1282 semantic_type: SemanticType::Field,
1283 column_id: 2,
1284 },
1285 ColumnMetadata {
1286 column_schema: ColumnSchema::new(
1287 "tag_a",
1288 ConcreteDataType::string_datatype(),
1289 true,
1290 ),
1291 semantic_type: SemanticType::Tag,
1292 column_id: 3,
1293 },
1294 ColumnMetadata {
1295 column_schema: ColumnSchema::new(
1296 "ts",
1297 ConcreteDataType::timestamp_millisecond_datatype(),
1298 false,
1299 )
1300 .with_time_index(true),
1301 semantic_type: SemanticType::Timestamp,
1302 column_id: 4,
1303 },
1304 ];
1305 let schemas = columns
1306 .iter()
1307 .map(|column| column.column_schema.clone())
1308 .collect::<Vec<_>>();
1309 let ids = columns
1310 .iter()
1311 .map(|column| (column.column_schema.name.clone(), column.column_id))
1312 .collect();
1313 let expected = vec![3, 2, 1, 4];
1314 assert_eq!(
1315 build_reconciliation_column_metadata(&schemas, &[2, 0], &ids)
1316 .unwrap()
1317 .iter()
1318 .map(|column| column.column_id)
1319 .collect::<Vec<_>>(),
1320 expected
1321 );
1322
1323 let mut metadata = build_region_metadata(RegionId::new(1024, 0), &columns);
1324 metadata.primary_key = vec![3, 1];
1325 metadata.schema_version = 2;
1326
1327 assert_eq!(
1328 check_column_metadatas_consistent(&[metadata.clone()]).unwrap(),
1329 columns
1330 );
1331
1332 assert_eq!(
1333 resolve_column_metadatas_with_latest(&[metadata])
1334 .unwrap()
1335 .0
1336 .iter()
1337 .map(|column| column.column_id)
1338 .collect::<Vec<_>>(),
1339 expected
1340 );
1341 }
1342
1343 #[test]
1344 fn test_reorder_tag_columns_rejects_invalid_primary_key() {
1345 let columns = new_test_column_metadatas();
1346
1347 let err = reorder_tag_columns(&columns, &[999]).unwrap_err();
1348 assert_matches!(err, Error::Unexpected { .. });
1349 assert!(err.to_string().contains("not found in column metadata"));
1350
1351 let err = reorder_tag_columns(&columns, &[]).unwrap_err();
1352 assert_matches!(err, Error::Unexpected { .. });
1353 assert!(err.to_string().contains("does not match tag columns"));
1354
1355 let err = reorder_tag_columns(&columns, &[2]).unwrap_err();
1356 assert_matches!(err, Error::Unexpected { .. });
1357 assert!(err.to_string().contains("is not a tag"));
1358
1359 let mut duplicate_tags = columns.clone();
1360 duplicate_tags[2].semantic_type = SemanticType::Tag;
1361 let err = reorder_tag_columns(&duplicate_tags, &[0, 0]).unwrap_err();
1362 assert_matches!(err, Error::Unexpected { .. });
1363 assert!(err.to_string().contains("is duplicated"));
1364 }
1365
1366 #[test]
1367 fn test_check_column_metadatas_consistent() {
1368 let column_metadatas = new_test_column_metadatas();
1369 let region_metadata1 = build_region_metadata(RegionId::new(1024, 0), &column_metadatas);
1370 let region_metadata2 = build_region_metadata(RegionId::new(1024, 1), &column_metadatas);
1371 let result =
1372 check_column_metadatas_consistent(&[region_metadata1, region_metadata2]).unwrap();
1373 assert_eq!(result, column_metadatas);
1374
1375 let region_metadata1 = build_region_metadata(RegionId::new(1025, 0), &column_metadatas);
1376 let region_metadata2 = build_region_metadata(RegionId::new(1024, 1), &column_metadatas);
1377 let result = check_column_metadatas_consistent(&[region_metadata1, region_metadata2]);
1378 assert!(result.is_none());
1379 }
1380
1381 #[test]
1382 fn test_check_column_metadata_invariants() {
1383 let column_metadatas = new_test_column_metadatas();
1384 let mut new_column_metadatas = column_metadatas.clone();
1385 new_column_metadatas.push(ColumnMetadata {
1386 column_schema: ColumnSchema::new("col3", ConcreteDataType::int32_datatype(), true),
1387 semantic_type: SemanticType::Field,
1388 column_id: 3,
1389 });
1390 check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap();
1391 }
1392
1393 #[test]
1394 fn test_check_column_metadata_invariants_missing_primary_key_column_or_ts_column() {
1395 let column_metadatas = new_test_column_metadatas();
1396 let mut new_column_metadatas = column_metadatas.clone();
1397 new_column_metadatas.retain(|c| c.semantic_type != SemanticType::Timestamp);
1398 check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap_err();
1399
1400 let column_metadatas = new_test_column_metadatas();
1401 let mut new_column_metadatas = column_metadatas.clone();
1402 new_column_metadatas.retain(|c| c.semantic_type != SemanticType::Tag);
1403 check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap_err();
1404 }
1405
1406 #[test]
1407 fn test_check_column_metadata_invariants_mismatch_column_id() {
1408 let column_metadatas = new_test_column_metadatas();
1409 let mut new_column_metadatas = column_metadatas.clone();
1410 if let Some(col) = new_column_metadatas
1411 .iter_mut()
1412 .find(|c| c.semantic_type == SemanticType::Timestamp)
1413 {
1414 col.column_id = 100;
1415 }
1416 check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap_err();
1417
1418 let column_metadatas = new_test_column_metadatas();
1419 let mut new_column_metadatas = column_metadatas.clone();
1420 if let Some(col) = new_column_metadatas
1421 .iter_mut()
1422 .find(|c| c.semantic_type == SemanticType::Tag)
1423 {
1424 col.column_id = 100;
1425 }
1426 check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap_err();
1427 }
1428
1429 #[test]
1430 fn test_resolve_column_metadatas_with_use_metasrv_strategy() {
1431 let column_metadatas = new_test_column_metadatas();
1432 let region_metadata1 = build_region_metadata(RegionId::new(1024, 0), &column_metadatas);
1433 let mut metasrv_column_metadatas = region_metadata1.column_metadatas.clone();
1434 metasrv_column_metadatas.push(ColumnMetadata {
1435 column_schema: ColumnSchema::new("col3", ConcreteDataType::int32_datatype(), true),
1436 semantic_type: SemanticType::Field,
1437 column_id: 3,
1438 });
1439 let result =
1440 resolve_column_metadatas_with_metasrv(&metasrv_column_metadatas, &[region_metadata1])
1441 .unwrap();
1442
1443 assert_eq!(result, vec![RegionId::new(1024, 0)]);
1444 }
1445
1446 #[test]
1447 fn test_resolve_column_metadatas_with_use_latest_strategy() {
1448 let column_metadatas = new_test_column_metadatas();
1449 let region_metadata1 = build_region_metadata(RegionId::new(1024, 0), &column_metadatas);
1450 let mut new_column_metadatas = column_metadatas.clone();
1451 new_column_metadatas.push(ColumnMetadata {
1452 column_schema: ColumnSchema::new("col3", ConcreteDataType::int32_datatype(), true),
1453 semantic_type: SemanticType::Field,
1454 column_id: 3,
1455 });
1456
1457 let mut region_metadata2 =
1458 build_region_metadata(RegionId::new(1024, 1), &new_column_metadatas);
1459 region_metadata2.schema_version = 2;
1460
1461 let (resolved_column_metadatas, region_ids) =
1462 resolve_column_metadatas_with_latest(&[region_metadata1, region_metadata2]).unwrap();
1463 assert_eq!(region_ids, vec![RegionId::new(1024, 0)]);
1464 assert_eq!(resolved_column_metadatas, new_column_metadatas);
1465 }
1466}