diff --git a/src/common/meta/src/reconciliation/reconcile_table/resolve_column_metadata.rs b/src/common/meta/src/reconciliation/reconcile_table/resolve_column_metadata.rs index 27950f0fd5d..8defc6bd0f2 100644 --- a/src/common/meta/src/reconciliation/reconcile_table/resolve_column_metadata.rs +++ b/src/common/meta/src/reconciliation/reconcile_table/resolve_column_metadata.rs @@ -25,8 +25,8 @@ use crate::reconciliation::reconcile_table::reconcile_regions::ReconcileRegions; use crate::reconciliation::reconcile_table::update_table_info::UpdateTableInfo; use crate::reconciliation::reconcile_table::{ReconcileTableContext, State, TableMetadataState}; use crate::reconciliation::utils::{ - ResolveColumnMetadataResult, build_column_metadata_from_table_info, - check_column_metadatas_consistent, resolve_column_metadatas_with_latest, + ResolveColumnMetadataResult, build_reconciliation_column_metadata, + check_column_metadatas_consistent, reorder_tag_columns, resolve_column_metadatas_with_latest, resolve_column_metadatas_with_metasrv, }; @@ -92,6 +92,8 @@ impl State for ResolveColumnMetadata { ctx.persistent_ctx.table_info_value = Some(table_info_value); if let Some(column_metadatas) = check_column_metadatas_consistent(&self.region_metadata) { + let column_metadatas = + reorder_tag_columns(&column_metadatas, &self.region_metadata[0].primary_key)?; // Safety: fetched in the above. let table_info_value = ctx.persistent_ctx.table_info_value.clone().unwrap(); info!( @@ -125,7 +127,7 @@ impl State for ResolveColumnMetadata { .table_info .name_to_ids() .context(MissingColumnIdsSnafu)?; - let column_metadata = build_column_metadata_from_table_info( + let column_metadata = build_reconciliation_column_metadata( table_info_value.table_info.meta.schema.column_schemas(), &table_info_value.table_info.meta.primary_key_indices, &name_to_ids, diff --git a/src/common/meta/src/reconciliation/utils.rs b/src/common/meta/src/reconciliation/utils.rs index 6ddc084596b..326516ecd80 100644 --- a/src/common/meta/src/reconciliation/utils.rs +++ b/src/common/meta/src/reconciliation/utils.rs @@ -85,6 +85,70 @@ pub(crate) fn check_column_metadatas_consistent( Some(region_metadatas[0].column_metadatas.clone()) } +/// Returns columns with tag values in primary-key order while preserving the +/// positions of non-tag columns. +pub(crate) fn reorder_tag_columns( + column_metadatas: &[ColumnMetadata], + primary_key: &[u32], +) -> Result> { + let tag_count = column_metadatas + .iter() + .filter(|column| column.semantic_type == SemanticType::Tag) + .count(); + ensure!( + primary_key.len() == tag_count, + UnexpectedSnafu { + err_msg: format!( + "Number of primary key columns {} does not match tag columns {}", + primary_key.len(), + tag_count, + ), + } + ); + + let mut primary_key_ids = HashSet::with_capacity(primary_key.len()); + let mut tags = Vec::with_capacity(primary_key.len()); + for column_id in primary_key { + let column = column_metadatas + .iter() + .find(|column| column.column_id == *column_id) + .with_context(|| UnexpectedSnafu { + err_msg: format!( + "Primary key column {} not found in column metadata", + column_id + ), + })?; + ensure!( + column.semantic_type == SemanticType::Tag, + UnexpectedSnafu { + err_msg: format!("Primary key column {} is not a tag", column_id), + } + ); + ensure!( + primary_key_ids.insert(column_id), + UnexpectedSnafu { + err_msg: format!("Primary key column {} is duplicated", column_id), + } + ); + tags.push(column); + } + + let mut tags = tags.into_iter(); + + column_metadatas + .iter() + .map(|column| { + if column.semantic_type == SemanticType::Tag { + tags.next().cloned().with_context(|| UnexpectedSnafu { + err_msg: "Primary key has fewer columns than tags".to_string(), + }) + } else { + Ok(column.clone()) + } + }) + .collect() +} + /// Resolves column metadata inconsistencies among the given region metadatas /// by using the column metadata from the metasrv as the source of truth. /// @@ -157,8 +221,13 @@ pub(crate) fn resolve_column_metadatas_with_latest( } } - // TODO(weny): verify the new column metadatas are acceptable for regions. - Ok((latest_region_metadata.column_metadatas.clone(), region_ids)) + Ok(( + reorder_tag_columns( + &latest_region_metadata.column_metadatas, + &latest_region_metadata.primary_key, + )?, + region_ids, + )) } /// Constructs a vector of [`ColumnMetadata`] from the provided table information. @@ -206,6 +275,22 @@ pub(crate) fn build_column_metadata_from_table_info( .collect::>>() } +/// Builds SyncColumns metadata with tags in the TableInfo primary-key order. +/// Non-tag columns retain their schema positions. +pub(crate) fn build_reconciliation_column_metadata( + column_schemas: &[ColumnSchema], + primary_key_indexes: &[usize], + name_to_ids: &HashMap, +) -> Result> { + let column_metadatas = + build_column_metadata_from_table_info(column_schemas, primary_key_indexes, name_to_ids)?; + let primary_key = primary_key_indexes + .iter() + .map(|index| column_metadatas[*index].column_id) + .collect::>(); + reorder_tag_columns(&column_metadatas, &primary_key) +} + /// Checks whether the schema invariants hold between the existing and new column metadata. /// /// Invariants: @@ -1180,6 +1265,104 @@ mod tests { assert_matches!(err, Error::MissingColumnInColumnMetadata { .. }); } + #[test] + fn test_reconciliation_columns_follow_primary_key_order() { + let columns = vec![ + ColumnMetadata { + column_schema: ColumnSchema::new( + "tag_b", + ConcreteDataType::string_datatype(), + true, + ), + semantic_type: SemanticType::Tag, + column_id: 1, + }, + ColumnMetadata { + column_schema: ColumnSchema::new("field", ConcreteDataType::int32_datatype(), true), + semantic_type: SemanticType::Field, + column_id: 2, + }, + ColumnMetadata { + column_schema: ColumnSchema::new( + "tag_a", + ConcreteDataType::string_datatype(), + true, + ), + semantic_type: SemanticType::Tag, + column_id: 3, + }, + ColumnMetadata { + column_schema: ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ) + .with_time_index(true), + semantic_type: SemanticType::Timestamp, + column_id: 4, + }, + ]; + let schemas = columns + .iter() + .map(|column| column.column_schema.clone()) + .collect::>(); + let ids = columns + .iter() + .map(|column| (column.column_schema.name.clone(), column.column_id)) + .collect(); + let expected = vec![3, 2, 1, 4]; + assert_eq!( + build_reconciliation_column_metadata(&schemas, &[2, 0], &ids) + .unwrap() + .iter() + .map(|column| column.column_id) + .collect::>(), + expected + ); + + let mut metadata = build_region_metadata(RegionId::new(1024, 0), &columns); + metadata.primary_key = vec![3, 1]; + metadata.schema_version = 2; + + assert_eq!( + check_column_metadatas_consistent(&[metadata.clone()]).unwrap(), + columns + ); + + assert_eq!( + resolve_column_metadatas_with_latest(&[metadata]) + .unwrap() + .0 + .iter() + .map(|column| column.column_id) + .collect::>(), + expected + ); + } + + #[test] + fn test_reorder_tag_columns_rejects_invalid_primary_key() { + let columns = new_test_column_metadatas(); + + let err = reorder_tag_columns(&columns, &[999]).unwrap_err(); + assert_matches!(err, Error::Unexpected { .. }); + assert!(err.to_string().contains("not found in column metadata")); + + let err = reorder_tag_columns(&columns, &[]).unwrap_err(); + assert_matches!(err, Error::Unexpected { .. }); + assert!(err.to_string().contains("does not match tag columns")); + + let err = reorder_tag_columns(&columns, &[2]).unwrap_err(); + assert_matches!(err, Error::Unexpected { .. }); + assert!(err.to_string().contains("is not a tag")); + + let mut duplicate_tags = columns.clone(); + duplicate_tags[2].semantic_type = SemanticType::Tag; + let err = reorder_tag_columns(&duplicate_tags, &[0, 0]).unwrap_err(); + assert_matches!(err, Error::Unexpected { .. }); + assert!(err.to_string().contains("is duplicated")); + } + #[test] fn test_check_column_metadatas_consistent() { let column_metadatas = new_test_column_metadatas(); diff --git a/src/mito2/src/engine/alter_test.rs b/src/mito2/src/engine/alter_test.rs index fd74c4b498a..cf039b52278 100644 --- a/src/mito2/src/engine/alter_test.rs +++ b/src/mito2/src/engine/alter_test.rs @@ -37,7 +37,7 @@ use datatypes::prelude::ConcreteDataType; use datatypes::schema::{ColumnSchema, FulltextAnalyzer, FulltextBackend, FulltextOptions}; use store_api::codec::PrimaryKeyEncoding; use store_api::logstore::provider::Provider; -use store_api::metadata::ColumnMetadata; +use store_api::metadata::{ColumnMetadata, RegionMetadata}; use store_api::metric_engine_consts::{PRIMARY_KEY_ENCODING, TABLE_COLUMN_METADATA_EXTENSION_KEY}; use store_api::region_engine::{RegionEngine, RegionManifestInfo, RegionRole}; use store_api::region_request::{ @@ -193,6 +193,133 @@ fn assert_column_metadatas(column_name: &[(&str, ColumnId)], column_metadatas: & } } +#[tokio::test] +async fn test_sync_columns_preserves_primary_key_order() { + test_sync_columns_preserves_primary_key_order_with_format(false).await; + test_sync_columns_preserves_primary_key_order_with_format(true).await; +} + +async fn assert_sync_columns_readable( + engine: &MitoEngine, + region_id: RegionId, + metadata: &RegionMetadata, + expected: &str, +) { + let current = engine.get_region(region_id).unwrap().metadata(); + assert_eq!(current.primary_key, metadata.primary_key); + let projection = metadata + .column_metadatas + .iter() + .map(|column| current.column_index_by_id(column.column_id).unwrap()) + .collect(); + let stream = engine + .scan_to_stream( + region_id, + ScanRequest { + projection: Some(projection), + ..Default::default() + }, + ) + .await + .unwrap(); + assert_eq!( + RecordBatches::try_collect(stream) + .await + .unwrap() + .pretty_print() + .unwrap(), + expected + ); +} + +async fn test_sync_columns_preserves_primary_key_order_with_format(flat_format: bool) { + let mut env = TestEnv::new().await; + let engine = env + .create_engine(MitoConfig { + default_flat_format: flat_format, + ..Default::default() + }) + .await; + let region_id = RegionId::new(1, 1); + let request = CreateRequestBuilder::new().build(); + let table_dir = request.table_dir.clone(); + let mut row_schema = rows_schema(&request); + row_schema.push(tag_column_schema("tag_1", ColumnDataType::String)); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + engine + .handle_request(region_id, RegionRequest::Create(request)) + .await + .unwrap(); + engine + .handle_request(region_id, RegionRequest::Alter(add_tag1())) + .await + .unwrap(); + let metadata = engine.get_region(region_id).unwrap().metadata(); + assert_eq!(metadata.primary_key, vec![0, 3]); + + put_rows( + &engine, + region_id, + Rows { + schema: row_schema, + rows: build_rows_for_tags("a", "b", 0, 2, 10), + }, + ) + .await; + flush_region(&engine, region_id, None).await; + let before = RecordBatches::try_collect( + engine + .scanner(region_id, ScanRequest::default()) + .await + .unwrap() + .scan() + .await + .unwrap(), + ) + .await + .unwrap() + .pretty_print() + .unwrap(); + + let mut columns = metadata.column_metadatas.clone(); + columns.push(ColumnMetadata { + column_schema: ColumnSchema::new("field_1", ConcreteDataType::float64_datatype(), true), + semantic_type: SemanticType::Field, + column_id: 4, + }); + let tags = metadata.primary_key_columns().cloned().collect::>(); + let mut tags = tags.into_iter(); + for column in &mut columns { + if column.semantic_type == SemanticType::Tag { + *column = tags.next().unwrap(); + } + } + engine + .handle_request( + region_id, + RegionRequest::Alter(RegionAlterRequest { + kind: AlterKind::SyncColumns { + column_metadatas: columns, + }, + }), + ) + .await + .unwrap(); + + assert_sync_columns_readable(&engine, region_id, &metadata, &before).await; + reopen_region(&engine, region_id, table_dir, true, HashMap::new()).await; + assert_sync_columns_readable(&engine, region_id, &metadata, &before).await; +} + #[tokio::test] async fn test_alter_region() { test_alter_region_with_format(false).await;