fix: preserve primary key order when syncing columns (#9189)

* fix: preserve primary key order when syncing columns

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix: reject invalid SyncColumns metadata

Signed-off-by: WenyXu <wenymedia@gmail.com>

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-09-27 02:15:47 +00:00
committed by GitHub
parent 194bc2fb3c
commit 5ffd01a70a
3 changed files with 318 additions and 6 deletions
@@ -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,
+185 -2
View File
@@ -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<Vec<ColumnMetadata>> {
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::<Result<Vec<_>>>()
}
/// 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<String, u32>,
) -> Result<Vec<ColumnMetadata>> {
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::<Vec<_>>();
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::<Vec<_>>();
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::<Vec<_>>(),
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::<Vec<_>>(),
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();
+128 -1
View File
@@ -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::<Vec<_>>();
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;