diff --git a/fluss-rust/crates/fluss/src/metadata/json_serde.rs b/fluss-rust/crates/fluss/src/metadata/json_serde.rs index c504660e92b..cc636a4ad6b 100644 --- a/fluss-rust/crates/fluss/src/metadata/json_serde.rs +++ b/fluss-rust/crates/fluss/src/metadata/json_serde.rs @@ -543,6 +543,23 @@ impl JsonSerde for Schema { } } + if let Some(highest_field_id_node) = node.get(Self::HIGHEST_FIELD_ID) { + let highest_field_id = + highest_field_id_node + .as_i64() + .ok_or_else(|| Error::JsonSerdeError { + message: format!("{} must be an integer", Self::HIGHEST_FIELD_ID), + })?; + let highest_field_id = + i32::try_from(highest_field_id).map_err(|_| Error::JsonSerdeError { + message: format!( + "{} value {highest_field_id} does not fit in i32", + Self::HIGHEST_FIELD_ID + ), + })?; + schema_builder = schema_builder.highest_field_id(highest_field_id); + } + schema_builder.build() } } @@ -984,6 +1001,44 @@ mod tests { assert_eq!(round_tripped, original); } + #[test] + fn schema_deserialization_preserves_highest_field_id_above_maximum() { + let json = json!({ + "version": 1, + "columns": [ + {"name": "id", "data_type": {"type": "INTEGER"}, "id": 0}, + {"name": "name", "data_type": {"type": "STRING"}, "id": 1} + ], + "highest_field_id": 10 + }); + + let schema = Schema::deserialize_json(&json).unwrap(); + assert_eq!(schema.highest_field_id(), 10); + } + + #[test] + fn schema_deserialization_clamps_highest_field_id_to_computed_minimum() { + let json = json!({ + "version": 1, + "columns": [ + {"name": "id", "data_type": {"type": "INTEGER"}, "id": 0}, + {"name": "name", "data_type": {"type": "STRING"}, "id": 1} + ], + "highest_field_id": 0 + }); + + let schema = Schema::deserialize_json(&json).unwrap(); + assert_eq!(schema.highest_field_id(), 1); + + let empty_json = json!({ + "version": 1, + "columns": [], + "highest_field_id": -3 + }); + let empty_schema = Schema::deserialize_json(&empty_json).unwrap(); + assert_eq!(empty_schema.highest_field_id(), -1); + } + #[test] fn schema_rejects_duplicate_nested_field_ids() { let nested = DataTypes::row(vec![ diff --git a/fluss-rust/crates/fluss/src/metadata/table.rs b/fluss-rust/crates/fluss/src/metadata/table.rs index c03a60065b2..405f54872fe 100644 --- a/fluss-rust/crates/fluss/src/metadata/table.rs +++ b/fluss-rust/crates/fluss/src/metadata/table.rs @@ -253,6 +253,7 @@ pub struct SchemaBuilder { columns: Vec, primary_key: Option, auto_increment_col_names: Vec, + highest_field_id: Option, } impl SchemaBuilder { @@ -330,9 +331,18 @@ impl SchemaBuilder { Ok(self) } + pub(crate) fn highest_field_id(mut self, highest_field_id: i32) -> Self { + self.highest_field_id = Some(highest_field_id); + self + } + pub fn build(&self) -> Result { let columns = Self::normalize_columns(&self.columns, self.primary_key.as_ref())?; - let (columns_with_ids, highest_field_id) = Self::assign_all_field_ids(columns)?; + let (columns_with_ids, maximum_field_id) = Self::assign_all_field_ids(columns)?; + let highest_field_id = self + .highest_field_id + .unwrap_or(maximum_field_id) + .max(maximum_field_id); if !self.auto_increment_col_names.is_empty() && self.primary_key.is_none() { return Err(IllegalArgument {