From 045441e3ccce8d474e7e34d9a4e6e5b49f16a74d Mon Sep 17 00:00:00 2001 From: fys <40801205+fengys1996@users.noreply.github.com> Date: Tue, 22 Sep 2026 12:56:26 +0000 Subject: [PATCH] feat(json2): support altering JSON2 settings (#9029) * feat(sql): support alter syntax for JSON2 columns Signed-off-by: fys * fix(json2): preserve rows on type hint mismatch during compaction * refactor(json2): simplify alter settings handling * fix(json2): preserve coerced values during compaction * chore: remove unnecessary clone * chor: reduce memory allocations * fix: cargo clippy * chore: update greptime-proto to main branch * refactor(datatypes): unify string handling with other JSON type hints * fix: cr --------- Signed-off-by: fys --- Cargo.lock | 3 +- Cargo.toml | 2 +- src/common/grpc-expr/src/alter.rs | 51 ++- src/common/grpc-expr/src/error.rs | 8 + .../meta/src/ddl/alter_table/executor.rs | 1 + .../src/ddl/alter_table/region_request.rs | 1 + src/common/meta/src/ddl/event/table.rs | 1 + src/datatypes/src/extension/json.rs | 55 ++- src/datatypes/src/json.rs | 384 +++++++++++++++--- src/datatypes/src/vectors/json/array.rs | 105 ++--- src/metric-engine/src/data_region.rs | 1 + src/mito2/src/compaction/json2.rs | 13 +- src/mito2/src/sst/parquet/json_align.rs | 78 +++- src/operator/src/expr_helper.rs | 109 +++-- src/sql/src/parsers/alter_parser.rs | 28 +- src/sql/src/statements/alter.rs | 38 +- src/sql/src/statements/create.rs | 39 +- src/store-api/Cargo.toml | 1 + src/store-api/src/metadata.rs | 36 ++ src/store-api/src/region_request.rs | 166 ++++++++ src/table/src/metadata.rs | 166 +++++++- src/table/src/requests.rs | 11 + .../common/types/json/json2_alter.result | 131 +++++- .../common/types/json/json2_alter.sql | 58 ++- .../json2_alter_type_hints_compaction.result | 83 ++++ .../json2_alter_type_hints_compaction.sql | 36 ++ ...er_type_hints_compaction_conversion.result | 78 ++++ ...alter_type_hints_compaction_conversion.sql | 35 ++ ..._compaction_null_type_hint_mismatch.result | 80 ++++ ...on2_compaction_null_type_hint_mismatch.sql | 35 ++ 30 files changed, 1602 insertions(+), 231 deletions(-) create mode 100644 tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction.result create mode 100644 tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction.sql create mode 100644 tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction_conversion.result create mode 100644 tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction_conversion.sql create mode 100644 tests/cases/standalone/common/types/json/json2_compaction_null_type_hint_mismatch.result create mode 100644 tests/cases/standalone/common/types/json/json2_compaction_null_type_hint_mismatch.sql diff --git a/Cargo.lock b/Cargo.lock index 2c7f3333329..ea55f1b297d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6130,7 +6130,7 @@ dependencies = [ [[package]] name = "greptime-proto" version = "0.1.0" -source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=0c4beffaa2b3968034915867593a717c6ab0c6b4#0c4beffaa2b3968034915867593a717c6ab0c6b4" +source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=3d232f406fad574e8c630fb1c304fe2ebb013c9c#3d232f406fad574e8c630fb1c304fe2ebb013c9c" dependencies = [ "prost 0.14.1", "prost-types 0.14.1", @@ -14457,6 +14457,7 @@ version = "1.3.0-alpha.1" dependencies = [ "api", "aquamarine", + "arrow-schema 59.2.0", "async-stream", "async-trait", "bytes", diff --git a/Cargo.toml b/Cargo.toml index b4c4c9d7f54..7c5c53d91c8 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -161,7 +161,7 @@ fs2 = "0.4" fst = "0.4.7" futures = "0.3" futures-util = "0.3" -greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "0c4beffaa2b3968034915867593a717c6ab0c6b4" } +greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "3d232f406fad574e8c630fb1c304fe2ebb013c9c" } hex = "0.4" hostname = "0.4.0" http = "1" diff --git a/src/common/grpc-expr/src/alter.rs b/src/common/grpc-expr/src/alter.rs index d5214ae5ddb..ffa148c7d0b 100644 --- a/src/common/grpc-expr/src/alter.rs +++ b/src/common/grpc-expr/src/alter.rs @@ -24,21 +24,23 @@ use api::v1::{ SkippingIndexType as PbSkippingIndexType, column_def, }; use common_query::AddColumnLocation; +use datatypes::json::{JsonSettings, JsonTypeHint}; +use datatypes::prelude::ConcreteDataType; use datatypes::schema::{ColumnSchema, FulltextOptions, Schema, SkippingIndexOptions}; use snafu::{OptionExt, ResultExt, ensure}; use store_api::region_request::{SetRegionOption, UnsetRegionOption}; use table::metadata::{TableId, TableMeta}; use table::requests::{ AddColumnRequest, AlterKind, AlterTableRequest, AnnotationFamily, ModifyColumnTypeRequest, - SetDefaultRequest, SetIndexOption, UnsetIndexOption, + SetDefaultRequest, SetIndexOption, SetJsonSettingsRequest, UnsetIndexOption, }; use crate::error::{ self, ColumnNotFoundSnafu, InvalidColumnDefSnafu, InvalidIndexOptionSnafu, - InvalidSetFulltextOptionRequestSnafu, InvalidSetSkippingIndexOptionRequestSnafu, - InvalidSetTableOptionRequestSnafu, InvalidUnsetTableOptionRequestSnafu, - MissingAlterIndexOptionSnafu, MissingFieldSnafu, MissingTableMetaSnafu, - MissingTimestampColumnSnafu, Result, UnknownLocationTypeSnafu, + InvalidJsonSettingsSnafu, InvalidSetFulltextOptionRequestSnafu, + InvalidSetSkippingIndexOptionRequestSnafu, InvalidSetTableOptionRequestSnafu, + InvalidUnsetTableOptionRequestSnafu, MissingAlterIndexOptionSnafu, MissingFieldSnafu, + MissingTableMetaSnafu, MissingTimestampColumnSnafu, Result, UnknownLocationTypeSnafu, }; const LOCATION_TYPE_FIRST: i32 = LocationType::First as i32; @@ -59,6 +61,34 @@ fn annotation_family_of_keys<'a>( }) } +fn json_settings_from_proto(settings: api::v1::JsonSettings) -> Result { + let type_hints = settings + .type_hints + .into_iter() + .map(|hint| { + let data_type = ConcreteDataType::from( + ColumnDataTypeWrapper::try_new(hint.data_type, hint.datatype_extension) + .context(error::ColumnDataTypeSnafu)?, + ); + + Ok(JsonTypeHint { + path: hint.path, + data_type, + // Index configuration is not supported yet, so this is temporarily + // hardcoded to false. + inverted_index: false, + }) + }) + .collect::>>()?; + + JsonSettings::try_new(type_hints, settings.max_auto_expanded_paths).map_err(|err| { + InvalidJsonSettingsSnafu { + err: err.to_string(), + } + .build() + }) +} + /// Returns the annotation family when `kind` is a SET/UNSET whose keys all /// belong to one family — the alters that only rewrite table metadata and skip /// region dispatch. A mixed batch is an error; never interpret it as "not an @@ -189,6 +219,17 @@ pub fn alter_expr_to_request( columns: modify_column_type_requests, } } + Kind::SetJsonSettings(set_json_settings) => { + let settings = set_json_settings + .settings + .context(MissingFieldSnafu { field: "settings" })?; + AlterKind::SetJsonSettings { + request: SetJsonSettingsRequest { + column_name: set_json_settings.column_name, + settings: json_settings_from_proto(settings)?, + }, + } + } Kind::DropColumns(DropColumns { drop_columns }) => AlterKind::DropColumns { names: drop_columns.into_iter().map(|c| c.name).collect(), }, diff --git a/src/common/grpc-expr/src/error.rs b/src/common/grpc-expr/src/error.rs index ccc40efb5cb..6c9302c39f4 100644 --- a/src/common/grpc-expr/src/error.rs +++ b/src/common/grpc-expr/src/error.rs @@ -73,6 +73,13 @@ pub enum Error { source: api::error::Error, }, + #[snafu(display("Invalid JSON settings: {}", err))] + InvalidJsonSettings { + err: String, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display("Unknown location type: {}", location_type))] UnknownLocationType { location_type: i32, @@ -188,6 +195,7 @@ impl ErrorExt for Error { | Error::MissingTimestampColumn { .. } => StatusCode::InvalidArguments, Error::MissingField { .. } => StatusCode::InvalidArguments, Error::InvalidColumnDef { source, .. } => source.status_code(), + Error::InvalidJsonSettings { .. } => StatusCode::InvalidArguments, Error::UnknownLocationType { .. } => StatusCode::InvalidArguments, Error::UnknownColumnDataType { .. } | Error::InvalidFulltextIndexColumnType { .. } => { diff --git a/src/common/meta/src/ddl/alter_table/executor.rs b/src/common/meta/src/ddl/alter_table/executor.rs index 53bc4b3d05b..419aa9ae12d 100644 --- a/src/common/meta/src/ddl/alter_table/executor.rs +++ b/src/common/meta/src/ddl/alter_table/executor.rs @@ -333,6 +333,7 @@ fn build_new_table_info( } AlterKind::DropColumns { .. } | AlterKind::ModifyColumnTypes { .. } + | AlterKind::SetJsonSettings { .. } | AlterKind::SetTableOptions { .. } | AlterKind::UnsetTableOptions { .. } | AlterKind::SetAnnotations { .. } diff --git a/src/common/meta/src/ddl/alter_table/region_request.rs b/src/common/meta/src/ddl/alter_table/region_request.rs index 73edd4c3615..b4f7b679422 100644 --- a/src/common/meta/src/ddl/alter_table/region_request.rs +++ b/src/common/meta/src/ddl/alter_table/region_request.rs @@ -90,6 +90,7 @@ fn create_proto_alter_kind( }))) } Kind::ModifyColumnTypes(x) => Ok(Some(alter_request::Kind::ModifyColumnTypes(x.clone()))), + Kind::SetJsonSettings(x) => Ok(Some(alter_request::Kind::SetJsonSettings(x.clone()))), Kind::DropColumns(x) => { let drop_columns = x .drop_columns diff --git a/src/common/meta/src/ddl/event/table.rs b/src/common/meta/src/ddl/event/table.rs index 670f3c7c1dc..c1ab0e5e5a2 100644 --- a/src/common/meta/src/ddl/event/table.rs +++ b/src/common/meta/src/ddl/event/table.rs @@ -193,6 +193,7 @@ pub(crate) fn alter_table_kind_name(kind: &AlterTableKind) -> Option<&'static st AlterTableKind::DropColumns(_) => Some("drop_columns"), AlterTableKind::RenameTable(_) => Some("rename_table"), AlterTableKind::ModifyColumnTypes(_) => Some("modify_column_types"), + AlterTableKind::SetJsonSettings(_) => Some("set_json_settings"), AlterTableKind::SetTableOptions(_) => Some("set_table_options"), AlterTableKind::UnsetTableOptions(_) => Some("unset_table_options"), AlterTableKind::SetIndex(_) => Some("set_index"), diff --git a/src/datatypes/src/extension/json.rs b/src/datatypes/src/extension/json.rs index 3aa0a4c9792..2e86d7791c5 100644 --- a/src/datatypes/src/extension/json.rs +++ b/src/datatypes/src/extension/json.rs @@ -24,9 +24,10 @@ use parquet_variant_compute::VariantType; use serde::{Deserialize, Serialize}; use snafu::{ResultExt, ensure}; -use crate::error::InvalidJson2LayoutSnafu; +use crate::error::{InvalidJson2LayoutSnafu, SerializeSnafu}; pub use crate::json::JSON2_REMAINDER_FIELD_NAME; use crate::json::JsonSettings; +use crate::schema::Metadata; /// Aligns JSON2 field types with their built arrays while preserving field and schema metadata. pub fn align_schema_with_json_array(schema: SchemaRef, columns: &[ArrayRef]) -> SchemaRef { @@ -151,7 +152,7 @@ pub struct JsonMetadata { } impl JsonMetadata { - /// Creates metadata for the JSON2 layout version 2. + /// Creates metadata for the latest JSON2 layout (currently V2). pub fn new(json_settings: JsonSettings) -> Self { Self { json_settings, @@ -300,6 +301,27 @@ impl ExtensionType for Json2ExtensionType { } } +/// Returns JSON2 column metadata with updated settings and the latest layout. +/// +/// Existing SSTs retain their own physical layout metadata. +pub fn json2_metadata_with_updated_settings( + current_metadata: &Metadata, + settings: JsonSettings, +) -> crate::error::Result { + let json_metadata = JsonMetadata::new(settings); + + let mut metadata = current_metadata.clone(); + metadata.insert( + EXTENSION_TYPE_NAME_KEY.to_string(), + Json2ExtensionType::NAME.to_string(), + ); + metadata.insert( + EXTENSION_TYPE_METADATA_KEY.to_string(), + serde_json::to_string(&json_metadata).context(SerializeSnafu)?, + ); + Ok(metadata) +} + /// Checks whether this field is either a legacy JSONB or JSON2 extension type. pub fn is_any_json_extension_type>(field: T) -> bool { let name = field.as_ref().extension_type_name(); @@ -442,6 +464,35 @@ mod tests { Ok(()) } + #[test] + fn test_json2_metadata_with_updated_settings_upgrades_layout() -> crate::error::Result<()> { + for (name, json) in [ + (JsonExtensionType::NAME, r#"{"json_settings":{}}"#), + (Json2ExtensionType::NAME, r#"{"json_settings":{}}"#), + ( + Json2ExtensionType::NAME, + r#"{"json_settings":{},"layout_version":2}"#, + ), + ] { + let metadata = HashMap::from([ + (EXTENSION_TYPE_NAME_KEY.to_string(), name.to_string()), + (EXTENSION_TYPE_METADATA_KEY.to_string(), json.to_string()), + ("other".to_string(), "kept".to_string()), + ]); + let settings = JsonSettings::try_new(vec![], Some(10))?; + let updated = json2_metadata_with_updated_settings(&metadata, settings.clone())?; + assert_eq!(Some("kept"), updated.get("other").map(String::as_str)); + assert_eq!( + Some(Json2ExtensionType::NAME), + updated.get(EXTENSION_TYPE_NAME_KEY).map(String::as_str) + ); + let json_metadata: JsonMetadata = + serde_json::from_str(updated.get(EXTENSION_TYPE_METADATA_KEY).unwrap()).unwrap(); + assert_eq!(JsonMetadata::new(settings), json_metadata); + } + Ok(()) + } + #[test] fn test_parse_json2_physical_layout() -> crate::error::Result<()> { let legacy = Field::new("data", DataType::Struct(Fields::empty()), true) diff --git a/src/datatypes/src/json.rs b/src/datatypes/src/json.rs index a5b4940e048..cd73e202bf3 100644 --- a/src/datatypes/src/json.rs +++ b/src/datatypes/src/json.rs @@ -23,6 +23,7 @@ pub mod value; use std::collections::BTreeMap; use std::collections::btree_map::Entry; +use std::sync::Arc; use serde::{Deserialize, Serialize}; use serde_json::{Map, Value as Json}; @@ -31,6 +32,7 @@ use snafu::ResultExt; use crate::data_type::ConcreteDataType; use crate::error::{self, InvalidJson2SettingsSnafu, Result, UnsupportedJsonTypeSnafu}; use crate::json::value::{JsonValue, JsonVariant, encode_serde_json_as_jsonb}; +use crate::prelude::DataType as _; use crate::types::json_type::{JsonNativeType, JsonObjectType}; use crate::value::{ListValue, StructValue, Value}; @@ -93,6 +95,15 @@ pub struct JsonTypeHint { pub inverted_index: bool, } +/// Specifies how encoding handles a value that does not match its JSON2 type hint. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TypeHintMismatchPolicy { + /// Returns an error for a type hint mismatch. + Reject, + /// Coerces to the hint type, encoding unconvertible values as null. + CoerceOrNull, +} + /// Context for JSON encoding/decoding that tracks the current key path. #[derive(Clone, Debug)] pub struct JsonContext<'a> { @@ -135,6 +146,30 @@ impl JsonSettings { (self.type_hints, self.max_auto_expanded_paths) } + /// Returns whether the settings are equivalent, ignoring type hint order. + pub fn equivalent(&self, other: &Self) -> bool { + let Self { + type_hints, + max_auto_expanded_paths, + } = self; + let Self { + type_hints: other_hints, + max_auto_expanded_paths: other_max_auto_expanded_paths, + } = other; + + if max_auto_expanded_paths != other_max_auto_expanded_paths + || type_hints.len() != other_hints.len() + { + return false; + } + + let mut hints = type_hints.iter().collect::>(); + let mut other_hints = other_hints.iter().collect::>(); + hints.sort_unstable_by(|a, b| a.path.cmp(&b.path)); + other_hints.sort_unstable_by(|a, b| a.path.cmp(&b.path)); + hints == other_hints + } + /// Decode an encoded StructValue back into a serde_json::Value. pub fn decode(&self, value: Value) -> Result { let mut context = JsonContext { @@ -146,6 +181,15 @@ impl JsonSettings { /// Encode a serde_json::Value into a Value::Json using current settings. pub fn encode(&self, json: Json) -> Result { + self.encode_with_type_hint_mismatch_policy(json, TypeHintMismatchPolicy::Reject) + } + + /// Encodes a serde_json::Value using the given type hint mismatch policy. + pub fn encode_with_type_hint_mismatch_policy( + &self, + json: Json, + policy: TypeHintMismatchPolicy, + ) -> Result { if let Json::Object(object) = &json && object.contains_key(JSON2_REMAINDER_FIELD_NAME) { @@ -160,7 +204,7 @@ impl JsonSettings { path: Vec::new(), settings: self, }; - encode_json_with_context(json, &mut context).map(|v| Value::Json(Box::new(v))) + encode_json_with_context(json, &mut context, policy).map(|v| Value::Json(Box::new(v))) } } @@ -272,36 +316,41 @@ fn with_key_context( } /// Main encoding function with key path tracking -fn encode_json_with_context(json: Json, context: &mut JsonContext) -> Result { +fn encode_json_with_context( + json: Json, + context: &mut JsonContext, + policy: TypeHintMismatchPolicy, +) -> Result { if context.path.is_empty() && !matches!(json, Json::Object(_)) { return UnsupportedJsonTypeSnafu.fail(); } match json { - Json::Object(json_object) => encode_json_object_with_context(json_object, context), - Json::Array(json_array) => encode_json_array_with_context(json_array, context), - _ => encode_json_value_with_context(json, context), + Json::Object(json_object) => encode_json_object_with_context(json_object, context, policy), + Json::Array(json_array) => encode_json_array_with_context(json_array, context, policy), + _ => encode_json_value_with_context(json, context, policy), } } fn encode_json_object_with_context<'a>( json_object: Map, context: &mut JsonContext<'a>, + policy: TypeHintMismatchPolicy, ) -> Result { let mut object = BTreeMap::new(); for (key, value) in json_object { let value = with_key_context(context, &key, |context| { if let Some(hint) = context.type_hint() { - encode_json_value_with_hint(value, hint, context) + encode_json_value_with_hint(value, hint, context, policy) } else { - encode_json_value_with_context(value, context) + encode_json_value_with_context(value, context, policy) } })?; object.insert(key, value.into_variant()); } - fill_missing_type_hints(&mut object, context)?; + fill_missing_type_hints(&mut object, context, policy)?; Ok(JsonValue::new(JsonVariant::Object(object))) } @@ -309,13 +358,14 @@ fn encode_json_object_with_context<'a>( fn fill_missing_type_hints( object: &mut BTreeMap, context: &mut JsonContext, + policy: TypeHintMismatchPolicy, ) -> Result<()> { for hint in &context.settings.type_hints { if hint.path.len() > context.path.len() && hint.path.starts_with(&context.path) { let depth = context.path.len(); let key = &hint.path[depth]; with_key_context(context, key, |context| { - insert_missing_type_hint(object, context, hint, depth) + insert_missing_type_hint(object, context, hint, depth, policy) })?; } } @@ -327,6 +377,7 @@ fn insert_missing_type_hint( field_context: &mut JsonContext, hint: &JsonTypeHint, depth: usize, + policy: TypeHintMismatchPolicy, ) -> Result<()> { let key = &hint.path[depth]; let is_leaf = depth + 1 == hint.path.len(); @@ -341,7 +392,11 @@ fn insert_missing_type_hint( match object.entry(key.clone()) { Entry::Occupied(mut entry) => match entry.get_mut() { JsonVariant::Object(child) => { - insert_missing_type_hint(child, field_context, hint, depth + 1) + insert_missing_type_hint(child, field_context, hint, depth + 1, policy) + } + _ if policy == TypeHintMismatchPolicy::CoerceOrNull => { + entry.insert(JsonValue::null().into_variant()); + Ok(()) } _ => error::InvalidJsonSnafu { value: format!( @@ -354,7 +409,7 @@ fn insert_missing_type_hint( }, Entry::Vacant(entry) => { let mut child = BTreeMap::new(); - insert_missing_type_hint(&mut child, field_context, hint, depth + 1)?; + insert_missing_type_hint(&mut child, field_context, hint, depth + 1, policy)?; entry.insert(JsonVariant::Object(child)); Ok(()) } @@ -365,65 +420,152 @@ fn encode_json_value_with_hint( json: Json, hint: &JsonTypeHint, context: &mut JsonContext, + policy: TypeHintMismatchPolicy, ) -> Result { if json.is_null() { return Ok(JsonValue::null()); } - let invalid_type = || { - error::InvalidJsonSnafu { + let json = match (&hint.data_type, json) { + (ConcreteDataType::String(_), Json::String(value)) => { + return Ok(value.into()); + } + ( + ConcreteDataType::Int8(_) + | ConcreteDataType::Int16(_) + | ConcreteDataType::Int32(_) + | ConcreteDataType::Int64(_), + Json::Number(value), + ) => { + if let Some(value) = value.as_i64() { + return Ok(value.into()); + } + Json::Number(value) + } + ( + ConcreteDataType::UInt8(_) + | ConcreteDataType::UInt16(_) + | ConcreteDataType::UInt32(_) + | ConcreteDataType::UInt64(_), + Json::Number(value), + ) => { + if let Some(value) = value.as_u64() { + return Ok(value.into()); + } + Json::Number(value) + } + (ConcreteDataType::Float32(_) | ConcreteDataType::Float64(_), Json::Number(value)) => { + if let Some(value) = value.as_f64() { + return Ok(value.into()); + } + Json::Number(value) + } + (ConcreteDataType::Boolean(_), Json::Bool(value)) => { + return Ok(value.into()); + } + (_, json) => json, + }; + + match policy { + TypeHintMismatchPolicy::Reject => error::InvalidJsonSnafu { value: format!( "JSON value at {} does not match JSON2 type hint {}", context.path.join("."), hint.data_type ), } - .fail() - }; + .fail(), + TypeHintMismatchPolicy::CoerceOrNull => { + let value = coerce_json_value_to_type(json, &hint.data_type); + let value = Json::try_from(value).map_err(|error| { + error::InvalidJsonSnafu { + value: error.to_string(), + } + .build() + })?; + Ok(JsonValue::new(value.into())) + } + } +} - match (&hint.data_type, json) { - (ConcreteDataType::String(_), Json::String(v)) => Ok(v.into()), - ( - ConcreteDataType::Int8(_) - | ConcreteDataType::Int16(_) - | ConcreteDataType::Int32(_) - | ConcreteDataType::Int64(_), - Json::Number(v), - ) => match v.as_i64() { - Some(v) => Ok(v.into()), - None => invalid_type(), - }, - ( - ConcreteDataType::UInt8(_) - | ConcreteDataType::UInt16(_) - | ConcreteDataType::UInt32(_) - | ConcreteDataType::UInt64(_), - Json::Number(v), - ) => match v.as_u64() { - Some(v) => Ok(v.into()), - None => invalid_type(), - }, - (ConcreteDataType::Float32(_) | ConcreteDataType::Float64(_), Json::Number(v)) => { - match v.as_f64() { - Some(v) => Ok(v.into()), - None => invalid_type(), +/// Coerces a JSON value using the same semantics as JSON2 query projection. +/// Values that cannot be converted become null. +pub(crate) fn coerce_json_value_to_type(value: Json, to_type: &ConcreteDataType) -> Value { + if value.is_null() { + return Value::Null; + } + + if to_type.is_string() { + let value = match value { + Json::String(value) => value, + value => value.to_string(), + }; + return Value::String(value.into()); + } + + if matches!(to_type, ConcreteDataType::Binary(_)) { + return Value::Binary(encode_serde_json_as_jsonb(value).into()); + } + + if let Some(struct_type) = to_type.as_struct() { + let Json::Object(mut object) = value else { + return Value::Null; + }; + let values = struct_type + .fields() + .iter() + .map(|field| { + object + .remove(field.name()) + .map(|value| coerce_json_value_to_type(value, field.data_type())) + .unwrap_or(Value::Null) + }) + .collect::>(); + return Value::Struct(StructValue::new(values, struct_type.clone())); + } + + if let Some(list_type) = to_type.as_list() { + let Json::Array(values) = value else { + return Value::Null; + }; + let item_type = list_type.item_type().clone(); + let values = values + .into_iter() + .map(|value| coerce_json_value_to_type(value, &item_type)) + .collect::>(); + return Value::List(ListValue::new(values, Arc::new(item_type))); + } + + let value = match value { + Json::Bool(value) => Value::Boolean(value), + Json::Number(value) => { + if let Some(value) = value.as_i64() { + Value::Int64(value) + } else if let Some(value) = value.as_u64() { + Value::UInt64(value) + } else if let Some(value) = value.as_f64() { + Value::Float64(value.into()) + } else { + Value::Null } } - (ConcreteDataType::Boolean(_), Json::Bool(v)) => Ok(v.into()), - _ => invalid_type(), - } + Json::String(value) => Value::String(value.into()), + Json::Array(_) | Json::Object(_) | Json::Null => Value::Null, + }; + to_type.try_cast(value).unwrap_or(Value::Null) } fn encode_json_array_with_context<'a>( json_array: Vec, context: &mut JsonContext<'a>, + policy: TypeHintMismatchPolicy, ) -> Result { let json_array_len = json_array.len(); let mut items = Vec::with_capacity(json_array_len); for (index, value) in json_array.into_iter().enumerate() { let item_value = with_key_context(context, &index.to_string(), |context| { - encode_json_value_with_context(value, context) + encode_json_value_with_context(value, context, policy) })?; items.push(item_value); } @@ -458,7 +600,11 @@ fn encode_json_array_with_context<'a>( } /// Helper function to encode a JSON value to a Value and determine its ConcreteDataType with context -fn encode_json_value_with_context(json: Json, context: &mut JsonContext) -> Result { +fn encode_json_value_with_context( + json: Json, + context: &mut JsonContext, + policy: TypeHintMismatchPolicy, +) -> Result { if context.path.len() >= JSON2_MAX_STRUCTURED_DEPTH && matches!(&json, Json::Object(_) | Json::Array(_)) { @@ -487,8 +633,8 @@ fn encode_json_value_with_context(json: Json, context: &mut JsonContext) -> Resu } } Json::String(s) => Ok(s.into()), - Json::Array(arr) => encode_json_array_with_context(arr, context), - Json::Object(obj) => encode_json_object_with_context(obj, context), + Json::Array(arr) => encode_json_array_with_context(arr, context, policy), + Json::Object(obj) => encode_json_object_with_context(obj, context, policy), } } @@ -597,6 +743,50 @@ mod tests { &struct_value.items()[index] } + #[test] + fn test_json_settings_equivalent() -> Result<()> { + let hints = vec![ + JsonTypeHint { + path: vec!["user".to_string(), "name".to_string()], + data_type: ConcreteDataType::string_datatype(), + inverted_index: false, + }, + JsonTypeHint { + path: vec!["count".to_string()], + data_type: ConcreteDataType::int64_datatype(), + inverted_index: false, + }, + ]; + let settings = JsonSettings::try_new(hints.clone(), Some(10))?; + let mut reversed = hints.clone(); + reversed.reverse(); + let reordered = JsonSettings::try_new(reversed, Some(10))?; + assert_ne!(settings, reordered); + assert!(settings.equivalent(&reordered)); + assert!(reordered.equivalent(&settings)); + assert!(settings.equivalent(&settings)); + assert_eq!(hints, settings.type_hints()); + assert!(JsonSettings::default().equivalent(&JsonSettings::default())); + + for limit in [None, Some(0), Some(11)] { + assert!(!settings.equivalent(&JsonSettings::try_new(hints.clone(), limit)?)); + } + assert!(!settings.equivalent(&JsonSettings::try_new(vec![hints[0].clone()], Some(10))?)); + + let mut changed = hints.clone(); + changed[0].path = vec!["user".to_string(), "id".to_string()]; + assert!(!settings.equivalent(&JsonSettings::try_new(changed, Some(10))?)); + + let mut changed = hints.clone(); + changed[0].data_type = ConcreteDataType::int64_datatype(); + assert!(!settings.equivalent(&JsonSettings::try_new(changed, Some(10))?)); + + let mut changed = hints; + changed[0].inverted_index = true; + assert!(!settings.equivalent(&JsonSettings::try_new(changed, Some(10))?)); + Ok(()) + } + #[test] fn test_json_settings_deserializes_legacy_type_hint_constraints() -> std::result::Result<(), Box> { @@ -755,6 +945,104 @@ mod tests { } } + #[test] + fn test_coerce_or_null_on_type_hint_mismatch() -> Result<()> { + let settings = JsonSettings::try_new( + vec![JsonTypeHint { + path: vec!["kind".to_string()], + data_type: ConcreteDataType::int64_datatype(), + inverted_index: false, + }], + None, + )?; + + let result = settings + .encode_with_type_hint_mismatch_policy( + json!({"kind": "invalid", "message": "kept"}), + TypeHintMismatchPolicy::CoerceOrNull, + )? + .into_json_inner() + .unwrap(); + let Value::Struct(root) = result else { + panic!("Expected Struct value"); + }; + assert_eq!(struct_field_value(&root, "kind"), &Value::Null); + assert_eq!( + struct_field_value(&root, "message"), + &Value::String("kept".into()) + ); + Ok(()) + } + + #[test] + fn test_coerce_or_null_converts_type_hint_values() -> Result<()> { + let settings = JsonSettings::try_new( + vec![ + JsonTypeHint { + path: vec!["to_string".to_string()], + data_type: ConcreteDataType::string_datatype(), + inverted_index: false, + }, + JsonTypeHint { + path: vec!["to_int".to_string()], + data_type: ConcreteDataType::int64_datatype(), + inverted_index: false, + }, + ], + None, + )?; + + let result = settings + .encode_with_type_hint_mismatch_policy( + json!({"to_string": 123, "to_int": "456", "invalid": "value"}), + TypeHintMismatchPolicy::CoerceOrNull, + )? + .into_json_inner() + .unwrap(); + let Value::Struct(root) = result else { + panic!("Expected Struct value"); + }; + assert_eq!( + struct_field_value(&root, "to_string"), + &Value::String("123".into()) + ); + assert_eq!(struct_field_value(&root, "to_int"), &Value::Int64(456)); + assert_eq!( + struct_field_value(&root, "invalid"), + &Value::String("value".into()) + ); + Ok(()) + } + + #[test] + fn test_coerce_or_null_on_nested_type_hint_structure_mismatch() -> Result<()> { + let settings = JsonSettings::try_new( + vec![JsonTypeHint { + path: vec!["user".to_string(), "age".to_string()], + data_type: ConcreteDataType::int64_datatype(), + inverted_index: false, + }], + None, + )?; + + let result = settings + .encode_with_type_hint_mismatch_policy( + json!({"user": "invalid", "message": "kept"}), + TypeHintMismatchPolicy::CoerceOrNull, + )? + .into_json_inner() + .unwrap(); + let Value::Struct(root) = result else { + panic!("Expected Struct value"); + }; + assert_eq!(struct_field_value(&root, "user"), &Value::Null); + assert_eq!( + struct_field_value(&root, "message"), + &Value::String("kept".into()) + ); + Ok(()) + } + #[test] fn test_encode_rejects_reserved_remainder_field() -> Result<()> { let settings = JsonSettings::default(); diff --git a/src/datatypes/src/vectors/json/array.rs b/src/datatypes/src/vectors/json/array.rs index 08a2fd86555..5b6265c7b37 100644 --- a/src/datatypes/src/vectors/json/array.rs +++ b/src/datatypes/src/vectors/json/array.rs @@ -26,15 +26,13 @@ use serde_json::Value; use snafu::{OptionExt, ResultExt}; use crate::arrow_array::{binary_array_value, string_array_value}; -use crate::data_type::ConcreteDataType; +use crate::data_type::{ConcreteDataType, DataType as _}; use crate::error::{ AlignJsonArraySnafu, ArrowComputeSnafu, InvalidJsonSnafu, InvalidJsonbSnafu, Result, }; use crate::extension::json::{JSON2_REMAINDER_FIELD_NAME, json2_remainder_field}; -use crate::json::JsonSettings; -use crate::json::value::{decode_json_variant, encode_serde_json_as_jsonb}; -use crate::prelude::{DataType as _, Value as GreptimeValue}; -use crate::value::{ListValue, StructValue}; +use crate::json::value::decode_json_variant; +use crate::json::{JsonSettings, TypeHintMismatchPolicy, coerce_json_value_to_type}; use crate::vectors::MutableVector; use crate::vectors::json::builder::{JsonVectorBuilder, json2_physical_data_type}; use crate::vectors::json::variant::variant_to_json_values; @@ -123,6 +121,23 @@ impl JsonArray<'_> { field: &Field, logical_settings: &JsonSettings, target_layout: &JsonSettings, + ) -> Result { + self.rewrite_to_v2_with_type_hint_mismatch_policy( + field, + logical_settings, + target_layout, + TypeHintMismatchPolicy::Reject, + ) + } + + /// Rewrites a JSON2 array to the specified v2 physical layout using the + /// given type hint mismatch policy. + pub fn rewrite_to_v2_with_type_hint_mismatch_policy( + &self, + field: &Field, + logical_settings: &JsonSettings, + target_layout: &JsonSettings, + policy: TypeHintMismatchPolicy, ) -> Result { let is_v2 = json2_remainder_field(field)?.is_some(); if is_v2 && self.inner.data_type() == &json2_physical_data_type(target_layout) { @@ -141,7 +156,8 @@ impl JsonArray<'_> { if value.is_null() { builder.push_null(); } else { - let value = logical_settings.encode(value)?; + let value = + logical_settings.encode_with_type_hint_mismatch_policy(value, policy)?; builder.try_push_value_ref(&value.as_value_ref())?; } } @@ -319,87 +335,12 @@ fn project_json_values(values: Vec, to_type: &DataType) -> Result Result { - if value.is_null() { - return Ok(GreptimeValue::Null); - } - - if to_type.is_string() { - let value = match value { - Value::String(value) => value, - value => value.to_string(), - }; - return Ok(GreptimeValue::String(value.into())); - } - - if matches!(to_type, ConcreteDataType::Binary(_)) { - return Ok(GreptimeValue::Binary( - encode_serde_json_as_jsonb(value).into(), - )); - } - - if let Some(struct_type) = to_type.as_struct() { - let Value::Object(mut object) = value else { - return Ok(GreptimeValue::Null); - }; - let values = struct_type - .fields() - .iter() - .map(|field| { - object - .remove(field.name()) - .map(|value| project_json_value_to_type(value, field.data_type())) - .transpose() - .map(|value| value.unwrap_or(GreptimeValue::Null)) - }) - .collect::>>()?; - return Ok(GreptimeValue::Struct(StructValue::new( - values, - struct_type.clone(), - ))); - } - - if let Some(list_type) = to_type.as_list() { - let Value::Array(values) = value else { - return Ok(GreptimeValue::Null); - }; - let item_type = list_type.item_type().clone(); - let values = values - .into_iter() - .map(|value| project_json_value_to_type(value, &item_type)) - .collect::>>()?; - return Ok(GreptimeValue::List(ListValue::new( - values, - Arc::new(item_type), - ))); - } - - let value = match value { - Value::Bool(value) => GreptimeValue::Boolean(value), - Value::Number(value) => { - if let Some(value) = value.as_i64() { - GreptimeValue::Int64(value) - } else if let Some(value) = value.as_u64() { - GreptimeValue::UInt64(value) - } else if let Some(value) = value.as_f64() { - GreptimeValue::Float64(value.into()) - } else { - GreptimeValue::Null - } - } - Value::String(value) => GreptimeValue::String(value.into()), - Value::Array(_) | Value::Object(_) => GreptimeValue::Null, - Value::Null => GreptimeValue::Null, - }; - Ok(to_type.try_cast(value).unwrap_or(GreptimeValue::Null)) -} - impl<'a> From<&'a ArrayRef> for JsonArray<'a> { fn from(inner: &'a ArrayRef) -> Self { Self { inner } diff --git a/src/metric-engine/src/data_region.rs b/src/metric-engine/src/data_region.rs index 55b9eaaf683..a11b3e57b54 100644 --- a/src/metric-engine/src/data_region.rs +++ b/src/metric-engine/src/data_region.rs @@ -227,6 +227,7 @@ impl DataRegion { | AlterKind::UnsetRegionOptions { keys: _ } | AlterKind::SetIndexes { options: _ } | AlterKind::UnsetIndexes { options: _ } + | AlterKind::SetJsonSettings { .. } | AlterKind::SyncColumns { column_metadatas: _, } => { diff --git a/src/mito2/src/compaction/json2.rs b/src/mito2/src/compaction/json2.rs index f381f80722e..cb2a9d34ec8 100644 --- a/src/mito2/src/compaction/json2.rs +++ b/src/mito2/src/compaction/json2.rs @@ -19,7 +19,9 @@ use arrow_schema::extension::ExtensionType; use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, Schema, SchemaRef}; use datatypes::arrow::record_batch::RecordBatch; use datatypes::extension::json::{JSON2_REMAINDER_FIELD_NAME, Json2ExtensionType, JsonMetadata}; -use datatypes::json::{JSON2_DEFAULT_MAX_AUTO_EXPANDED_PATHS, JsonSettings, JsonTypeHint}; +use datatypes::json::{ + JSON2_DEFAULT_MAX_AUTO_EXPANDED_PATHS, JsonSettings, JsonTypeHint, TypeHintMismatchPolicy, +}; use datatypes::prelude::ConcreteDataType; use datatypes::types::json_type::JsonNativeType; use datatypes::vectors::json::array::JsonArray; @@ -260,7 +262,7 @@ fn select_dynamic_hints( .filter(|(path, stat)| { !stat.is_type_conflicted // TODO(LFC): Instead of "primitive only", consider retaining stable compound types - // that are safe to write to Parquet, as flush does. Or better, unite the two + // that are safe to write to Parquet, as flush does. Or better, unite the two // selection process. && stat.data_type.is_primitive() && !has_ancestor_path(path) @@ -327,7 +329,12 @@ pub(crate) fn rewrite_json2_batch( }; let array = JsonArray::from(array) - .rewrite_to_v2(field, &plan.logical_settings, &plan.target_layout) + .rewrite_to_v2_with_type_hint_mismatch_policy( + field, + &plan.logical_settings, + &plan.target_layout, + TypeHintMismatchPolicy::CoerceOrNull, + ) .context(ConvertValueSnafu)?; debug_assert_eq!( &json2_physical_data_type(&plan.target_layout), diff --git a/src/mito2/src/sst/parquet/json_align.rs b/src/mito2/src/sst/parquet/json_align.rs index 8af88f6ac0f..93ecffa072f 100644 --- a/src/mito2/src/sst/parquet/json_align.rs +++ b/src/mito2/src/sst/parquet/json_align.rs @@ -22,7 +22,7 @@ use datatypes::arrow::array::{ArrayRef, new_null_array}; use datatypes::arrow::datatypes::{DataType, Field, FieldRef, Schema, SchemaRef}; use datatypes::arrow::record_batch::RecordBatch; use datatypes::extension::json::{JsonMetadata, is_json2_extension_type}; -use datatypes::json::JsonSettings; +use datatypes::json::{JsonSettings, TypeHintMismatchPolicy}; use datatypes::vectors::json::array::JsonArray; use datatypes::vectors::json::json2_physical_data_type; use futures::Stream; @@ -251,10 +251,11 @@ fn rewrite_array( settings: &RewriteSettings, ) -> Result { JsonArray::from(source_array) - .rewrite_to_v2( + .rewrite_to_v2_with_type_hint_mismatch_policy( source_field, &settings.logical_settings, &settings.target_layout, + TypeHintMismatchPolicy::CoerceOrNull, ) .context(DataTypeMismatchSnafu) } @@ -330,7 +331,9 @@ mod tests { Array, ArrayRef, BinaryArray, Int64Array, StringArray, StringViewArray, StructArray, }; use datatypes::arrow::datatypes::{DataType, Field, Fields, Schema}; - use datatypes::extension::json::Json2ExtensionType; + use datatypes::extension::json::{Json2ExtensionType, JsonMetadata}; + use datatypes::json::JsonTypeHint; + use datatypes::prelude::ConcreteDataType; use datatypes::types::parse_string_to_jsonb; use futures::{StreamExt, stream}; @@ -783,6 +786,75 @@ mod tests { assert_eq!(int_array([10, 20]).as_ref(), output.column(3).as_ref()); } + #[tokio::test] + async fn test_rewrite_keeps_rows_with_invalid_json2_settings() { + let settings = JsonSettings::try_new( + vec![JsonTypeHint { + path: vec!["kind".to_string()], + data_type: ConcreteDataType::string_datatype(), + inverted_index: false, + }], + Some(0), + ) + .unwrap(); + let target_type = json2_physical_data_type(&settings); + let output_schema = schema([ + Field::new("j", target_type.clone(), true).with_extension_type( + Json2ExtensionType::new(Arc::new(JsonMetadata::new(settings.clone()))), + ), + Field::new("value", DataType::Int64, true), + ]); + let source = Arc::new(BinaryArray::from_iter([ + Some(parse_string_to_jsonb(r#"{"kind":"valid"}"#).unwrap()), + Some(parse_string_to_jsonb(r#"{"kind":1}"#).unwrap()), + ])) as ArrayRef; + let input = RecordBatch::try_new( + schema([ + Field::new("j", DataType::Binary, true) + .with_extension_type(Json2ExtensionType::default()), + Field::new("value", DataType::Int64, true), + ]), + vec![source, int_array([10, 20])], + ) + .unwrap(); + let columns = HashMap::from([( + "j".to_string(), + Json2TargetLayout { + extension_metadata: serde_json::to_string(&JsonMetadata::new(settings.clone())) + .unwrap(), + target_layout: settings, + }, + )]); + let mut aligner = JsonSchemaAligner::new( + stream::iter([Ok(input)]), + vec![true, true], + output_schema, + AlignMode::Rewrite { columns }, + ) + .unwrap(); + + let output = aligner.next().await.unwrap().unwrap(); + assert_eq!(2, output.num_rows()); + assert_eq!( + 10, + output + .column(1) + .as_any() + .downcast_ref::() + .unwrap() + .value(0) + ); + assert_eq!( + 20, + output + .column(1) + .as_any() + .downcast_ref::() + .unwrap() + .value(1) + ); + } + #[test] fn test_rewrite_rejects_mismatched_output_layout() { let columns = HashMap::from([( diff --git a/src/operator/src/expr_helper.rs b/src/operator/src/expr_helper.rs index b28af78d383..1d06c24ff7f 100644 --- a/src/operator/src/expr_helper.rs +++ b/src/operator/src/expr_helper.rs @@ -24,9 +24,10 @@ use api::v1::column_def::{options_from_column_schema, try_as_column_schema}; use api::v1::{ AddColumn, AddColumns, AlterDatabaseExpr, AlterTableExpr, Analyzer, ColumnDataType, ColumnDataTypeExtension, CreateFlowExpr, CreateTableExpr, CreateViewExpr, DropColumn, - DropColumns, DropDefaults, ExpireAfter, FulltextBackend as PbFulltextBackend, ModifyColumnType, + DropColumns, DropDefaults, ExpireAfter, FulltextBackend as PbFulltextBackend, + JsonSettings as PbJsonSettings, JsonTypeHint as PbJsonTypeHint, ModifyColumnType, ModifyColumnTypes, RenameTable, SemanticType, SetDatabaseOptions, SetDefaults, SetFulltext, - SetIndex, SetIndexes, SetInverted, SetSkipping, SetTableOptions, + SetIndex, SetIndexes, SetInverted, SetJsonSettings, SetSkipping, SetTableOptions, SkippingIndexType as PbSkippingIndexType, TableName, UnsetDatabaseOptions, UnsetFulltext, UnsetIndex, UnsetIndexes, UnsetInverted, UnsetSkipping, UnsetTableOptions, set_index, unset_index, @@ -36,6 +37,7 @@ use common_error::ext::BoxedError; use common_grpc_expr::util::ColumnExpr; use common_time::Timezone; use datafusion::sql::planner::object_name_to_table_reference; +use datatypes::json::JsonSettings; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{ COLUMN_FULLTEXT_OPT_KEY_ANALYZER, COLUMN_FULLTEXT_OPT_KEY_BACKEND, @@ -782,6 +784,29 @@ pub(crate) fn to_repartition_request( }) } +fn json_settings_to_proto(settings: JsonSettings) -> Result { + let (type_hints, max_auto_expanded_paths) = settings.into_parts(); + let type_hints = type_hints + .into_iter() + .map(|hint| { + let (data_type, datatype_extension) = ColumnDataTypeWrapper::try_from(hint.data_type) + .map(|w| w.to_parts()) + .context(ColumnDataTypeSnafu)?; + + Ok(PbJsonTypeHint { + path: hint.path, + data_type: data_type as i32, + datatype_extension, + }) + }) + .collect::>>()?; + + Ok(PbJsonSettings { + type_hints, + max_auto_expanded_paths, + }) +} + /// Converts a SQL alter table statement into a gRPC alter table expression. pub(crate) fn to_alter_table_expr( alter_table: AlterTable, @@ -830,12 +855,15 @@ pub(crate) fn to_alter_table_expr( } => { let target_type = sql_data_type_to_concrete_data_type(&target_type).context(ParseSqlSnafu)?; + + // Currently disallow modify column type to json2. if target_type.is_json2() { return NotSupportedSnafu { feat: "ALTER TABLE MODIFY COLUMN to JSON2 type", } .fail(); } + let (target_type, target_type_extension) = ColumnDataTypeWrapper::try_from(target_type) .map(|w| w.to_parts()) .context(ColumnDataTypeSnafu)?; @@ -853,6 +881,20 @@ pub(crate) fn to_alter_table_expr( }], }) } + AlterTableOperation::SetJsonSettings { + column_name, + json2_options, + } => { + let settings = match json2_options { + Some(options) => options.build_json_settings().context(ParseSqlSnafu)?, + None => datatypes::json::JsonSettings::new_v2(), + }; + + AlterTableKind::SetJsonSettings(SetJsonSettings { + column_name: column_name.value, + settings: Some(json_settings_to_proto(settings)?), + }) + } AlterTableOperation::DropColumn { name } => AlterTableKind::DropColumns(DropColumns { drop_columns: vec![DropColumn { name: name.value.clone(), @@ -1752,31 +1794,48 @@ SELECT max(c1), min(c2) FROM schema_2.table_2;"; } #[test] - fn test_alter_json2_is_not_supported() { - for sql in [ - "ALTER TABLE monitor MODIFY COLUMN payload JSON2;", - "ALTER TABLE monitor MODIFY COLUMN payload JSON2 (service STRING);", - ] { - let stmt = ParserContext::create_with_dialect( - sql, - &GreptimeDbDialect {}, - ParseOptions::default(), - ) - .unwrap() - .pop() - .unwrap(); + fn test_to_alter_set_json_settings_expr() { + let sql = "ALTER TABLE monitor MODIFY COLUMN payload JSON2 (service STRING);"; + let stmt = + ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default()) + .unwrap() + .pop() + .unwrap(); - let Statement::AlterTable(alter_table) = stmt else { - unreachable!() - }; - let err = to_alter_table_expr(alter_table, &QueryContext::arc()).unwrap_err(); + let Statement::AlterTable(alter_table) = stmt else { + unreachable!() + }; + let expr = to_alter_table_expr(alter_table, &QueryContext::arc()).unwrap(); + let kind = expr.kind.unwrap(); - assert!(matches!(err, crate::error::Error::NotSupported { .. })); - assert_eq!( - "Not supported: ALTER TABLE MODIFY COLUMN to JSON2 type", - err.to_string() - ); - } + let AlterTableKind::SetJsonSettings(modify) = kind else { + unreachable!() + }; + assert_eq!("payload", modify.column_name); + let settings = modify.settings.as_ref().unwrap(); + assert_eq!(Some(100), settings.max_auto_expanded_paths); + assert_eq!(1, settings.type_hints.len()); + assert_eq!(["service"], &settings.type_hints[0].path[..]); + + let sql = "ALTER TABLE monitor MODIFY COLUMN payload JSON2;"; + let stmt = + ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default()) + .unwrap() + .pop() + .unwrap(); + + let Statement::AlterTable(alter_table) = stmt else { + unreachable!() + }; + let expr = to_alter_table_expr(alter_table, &QueryContext::arc()).unwrap(); + let kind = expr.kind.unwrap(); + + let AlterTableKind::SetJsonSettings(modify) = kind else { + unreachable!() + }; + let settings = modify.settings.as_ref().unwrap(); + assert_eq!(Some(100), settings.max_auto_expanded_paths); + assert!(settings.type_hints.is_empty()); } #[test] diff --git a/src/sql/src/parsers/alter_parser.rs b/src/sql/src/parsers/alter_parser.rs index ca38ec0d65e..3e5c27fda73 100644 --- a/src/sql/src/parsers/alter_parser.rs +++ b/src/sql/src/parsers/alter_parser.rs @@ -462,19 +462,19 @@ impl ParserContext<'_> { .context(error::SyntaxSnafu)?; self.parse_alter_table_drop_default(column_name) } else { - let (data_type, json2_options) = - if let Some(json2) = parse_json2_type_and_options(&mut self.parser)? { - json2 - } else { - ( - self.parser.parse_data_type().context(error::SyntaxSnafu)?, - None, - ) - }; + if let Some((_, json2_options)) = + parse_json2_type_and_options(&mut self.parser)? + { + return Ok(AlterTableOperation::SetJsonSettings { + column_name, + json2_options, + }); + } + Ok(AlterTableOperation::ModifyColumnType { column_name, - target_type: data_type, - json2_options, + target_type: self.parser.parse_data_type().context(error::SyntaxSnafu)?, + json2_options: None, }) } } @@ -1576,9 +1576,8 @@ MODIFY COLUMN attrs JSON2 ( let Statement::AlterTable(alter_table) = statements.remove(0) else { unreachable!() }; - let AlterTableOperation::ModifyColumnType { + let AlterTableOperation::SetJsonSettings { column_name, - target_type, json2_options: Some(options), } = alter_table.alter_operation() else { @@ -1586,7 +1585,6 @@ MODIFY COLUMN attrs JSON2 ( }; assert_eq!("attrs", column_name.value); - assert_eq!("JSON2", target_type.to_string()); assert_eq!(Some(2000), options.max_auto_expanded_paths); assert_eq!(4, options.type_hints.len()); assert_eq!(vec!["user", "id"], options.type_hints[1].path); @@ -1610,7 +1608,7 @@ MODIFY COLUMN attrs JSON2 ( let Statement::AlterTable(empty) = empty.remove(0) else { unreachable!() }; - let AlterTableOperation::ModifyColumnType { json2_options, .. } = empty.alter_operation() + let AlterTableOperation::SetJsonSettings { json2_options, .. } = empty.alter_operation() else { unreachable!() }; diff --git a/src/sql/src/statements/alter.rs b/src/sql/src/statements/alter.rs index 2f98d8737ae..4fc45c2e69b 100644 --- a/src/sql/src/statements/alter.rs +++ b/src/sql/src/statements/alter.rs @@ -88,6 +88,11 @@ pub enum AlterTableOperation { target_type: DataType, json2_options: Option, }, + /// `MODIFY JSON2 [json2_options]` + SetJsonSettings { + column_name: Ident, + json2_options: Option, + }, /// `SET =
` SetTableOptions { options: Vec, @@ -265,6 +270,16 @@ impl Display for AlterTableOperation { } Ok(()) } + AlterTableOperation::SetJsonSettings { + column_name, + json2_options, + } => { + write!(f, r#"MODIFY COLUMN {column_name} JSON2"#)?; + if let Some(options) = json2_options { + write!(f, "{options}")?; + } + Ok(()) + } AlterTableOperation::SetTableOptions { options } => { let kvs = options .iter() @@ -433,8 +448,10 @@ impl Display for AlterDatabaseOperation { mod tests { use std::assert_matches; + use super::AlterTableOperation; use crate::dialect::GreptimeDbDialect; use crate::parser::{ParseOptions, ParserContext}; + use crate::statements::create::Json2Options; use crate::statements::statement::Statement; #[test] @@ -503,7 +520,7 @@ ALTER TABLE monitor ADD COLUMN app STRING DEFAULT 'shop' PRIMARY KEY, ADD COLUMN } let sql = r"alter table monitor modify column load_15 string;"; - let stmts = + let mut stmts = ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default()) .unwrap(); assert_eq!(1, stmts.len()); @@ -523,6 +540,25 @@ ALTER TABLE monitor MODIFY COLUMN load_15 STRING"#, } } + let Statement::AlterTable(alter_table) = &mut stmts[0] else { + unreachable!(); + }; + let AlterTableOperation::ModifyColumnType { json2_options, .. } = + alter_table.alter_operation_mut() + else { + unreachable!(); + }; + *json2_options = Some(Json2Options { + max_auto_expanded_paths: Some(1), + type_hints: vec![], + }); + assert_eq!( + r#"ALTER TABLE monitor MODIFY COLUMN load_15 STRING( + max_auto_expanded_paths = 1 + )"#, + alter_table.to_string() + ); + let sql = r"alter table monitor drop column load_15;"; let stmts = ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default()) diff --git a/src/sql/src/statements/create.rs b/src/sql/src/statements/create.rs index 46b5b07c9f7..2310531841c 100644 --- a/src/sql/src/statements/create.rs +++ b/src/sql/src/statements/create.rs @@ -157,6 +157,26 @@ impl Display for Json2Options { } } +impl Json2Options { + pub fn build_json_settings(&self) -> Result { + let type_hints = self + .type_hints + .iter() + .map(|hint| { + Ok(datatypes::json::JsonTypeHint { + path: hint.path.clone(), + data_type: json_type_hint_concrete_data_type(&hint.data_type)?, + inverted_index: hint.inverted_index, + }) + }) + .collect::>>()?; + let max_auto_expanded_paths = self + .max_auto_expanded_paths + .or(Some(JSON2_DEFAULT_MAX_AUTO_EXPANDED_PATHS)); + JsonSettings::try_new(type_hints, max_auto_expanded_paths).map_err(Into::into) + } +} + #[derive(Debug, PartialEq, Eq, Clone, Visit, VisitMut, Serialize)] pub struct JsonTypeHint { pub path: Vec, @@ -351,24 +371,7 @@ impl ColumnExtensions { return Ok(None); }; - let type_hints = options - .type_hints - .iter() - .map(|hint| { - Ok(datatypes::json::JsonTypeHint { - path: hint.path.clone(), - data_type: json_type_hint_concrete_data_type(&hint.data_type)?, - inverted_index: hint.inverted_index, - }) - }) - .collect::>>()?; - let settings = JsonSettings::try_new( - type_hints, - options - .max_auto_expanded_paths - .or(Some(JSON2_DEFAULT_MAX_AUTO_EXPANDED_PATHS)), - )?; - Ok(Some(settings)) + options.build_json_settings().map(Some) } pub fn set_json_settings(&mut self, settings: JsonSettings) -> Result<()> { diff --git a/src/store-api/Cargo.toml b/src/store-api/Cargo.toml index 61ce914fe13..2e5a203912e 100644 --- a/src/store-api/Cargo.toml +++ b/src/store-api/Cargo.toml @@ -10,6 +10,7 @@ workspace = true [dependencies] api.workspace = true aquamarine.workspace = true +arrow-schema.workspace = true async-trait.workspace = true bytes.workspace = true common-base.workspace = true diff --git a/src/store-api/src/metadata.rs b/src/store-api/src/metadata.rs index 0c663bccc05..21fc3295c44 100644 --- a/src/store-api/src/metadata.rs +++ b/src/store-api/src/metadata.rs @@ -30,6 +30,7 @@ use common_error::status_code::StatusCode; use common_macro::stack_trace_debug; use datatypes::arrow; use datatypes::arrow::datatypes::FieldRef; +use datatypes::extension::json::json2_metadata_with_updated_settings; use datatypes::schema::{ColumnSchema, FulltextOptions, Schema, SchemaRef, VectorIndexOptions}; use datatypes::types::TimestampType; use itertools::Itertools; @@ -663,6 +664,10 @@ impl RegionMetadataBuilder { AlterKind::AddColumns { columns } => self.add_columns(columns)?, AlterKind::DropColumns { names } => self.drop_columns(&names), AlterKind::ModifyColumnTypes { columns } => self.modify_column_types(columns)?, + AlterKind::SetJsonSettings { + column_name, + settings, + } => self.set_json_settings(column_name, settings)?, AlterKind::SetIndexes { options } => self.set_indexes(options)?, AlterKind::UnsetIndexes { options } => self.unset_indexes(options)?, AlterKind::SetRegionOptions { options: _ } => { @@ -831,6 +836,37 @@ impl RegionMetadataBuilder { Ok(()) } + fn set_json_settings( + &mut self, + col_name: String, + settings: datatypes::json::JsonSettings, + ) -> Result<()> { + let Some(col_meta) = self + .column_metadatas + .iter_mut() + .find(|col| col.column_schema.name == col_name) + else { + return InvalidRegionRequestSnafu { + region_id: self.region_id, + err: format!("column {col_name} not found"), + } + .fail(); + }; + + let old_metadata = col_meta.column_schema.metadata(); + let new_metadata = + json2_metadata_with_updated_settings(old_metadata, settings).map_err(|err| { + InvalidRegionRequestSnafu { + region_id: self.region_id, + err: err.to_string(), + } + .build() + })?; + + *col_meta.column_schema.mut_metadata() = new_metadata; + Ok(()) + } + fn set_indexes(&mut self, options: Vec) -> Result<()> { let mut set_index_map: HashMap<_, Vec<_>> = HashMap::new(); for option in &options { diff --git a/src/store-api/src/region_request.rs b/src/store-api/src/region_request.rs index 826ef7d6bc9..113a2393ccc 100644 --- a/src/store-api/src/region_request.rs +++ b/src/store-api/src/region_request.rs @@ -32,6 +32,7 @@ use api::v1::{ self, Analyzer, ArrowIpc, FulltextBackend as PbFulltextBackend, Option as PbOption, Rows, SemanticType, SkippingIndexType as PbSkippingIndexType, WriteHint, }; +use arrow_schema::extension::ExtensionType; pub use common_base::AffectedRows; use common_base::readable_size::ReadableSize; use common_grpc::flight::FlightDecoder; @@ -39,6 +40,8 @@ use common_recordbatch::DfRecordBatch; use common_time::range::TimestampRange; use common_time::{TimeToLive, Timestamp}; use datatypes::error::time_index_not_widening_error; +use datatypes::extension::json::Json2ExtensionType; +use datatypes::json::{JsonSettings, JsonTypeHint}; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{FulltextOptions, SkippingIndexOptions}; use num_enum::TryFromPrimitive; @@ -795,6 +798,13 @@ pub enum AlterKind { /// Columns to change. columns: Vec, }, + /// Set JSON2 settings of a region column. + SetJsonSettings { + /// Column name. + column_name: String, + /// Target JSON2 settings. + settings: JsonSettings, + }, /// Set region options. SetRegionOptions { options: Vec }, /// Unset region options. @@ -975,6 +985,9 @@ impl AlterKind { col_to_change.validate(metadata)?; } } + AlterKind::SetJsonSettings { column_name, .. } => { + Self::validate_set_json_settings(column_name, metadata)? + } AlterKind::SetRegionOptions { .. } => {} AlterKind::UnsetRegionOptions { .. } => {} AlterKind::SetIndexes { options } => { @@ -1084,6 +1097,18 @@ impl AlterKind { AlterKind::ModifyColumnTypes { columns } => columns .iter() .any(|col_to_change| col_to_change.need_alter(metadata)), + AlterKind::SetJsonSettings { + column_name, + settings, + } => metadata.column_by_name(column_name).is_some_and(|col| { + col.column_schema + .extension_type::() + .ok() + .flatten() + .is_none_or(|extension| { + !extension.metadata().json_settings().equivalent(settings) + }) + }), AlterKind::SetRegionOptions { .. } => true, AlterKind::UnsetRegionOptions { .. } => true, AlterKind::SetIndexes { options, .. } => options @@ -1160,6 +1185,34 @@ impl AlterKind { Ok(()) } + + fn validate_set_json_settings(col_name: &String, metadata: &RegionMetadata) -> Result<()> { + let region_id = metadata.region_id; + + let col = metadata + .column_by_name(col_name) + .with_context(|| InvalidRegionRequestSnafu { + region_id, + err: format!("column {} not found", col_name), + })?; + + ensure!( + col.semantic_type == SemanticType::Field, + InvalidRegionRequestSnafu { + region_id, + err: format!("column {} is not a field column", col_name), + } + ); + ensure!( + col.column_schema.data_type.is_json2(), + InvalidRegionRequestSnafu { + region_id, + err: format!("column {} is not a JSON2 column", col_name), + } + ); + + Ok(()) + } } impl TryFrom for AlterKind { @@ -1183,6 +1236,15 @@ impl TryFrom for AlterKind { .collect::>(); AlterKind::ModifyColumnTypes { columns } } + alter_request::Kind::SetJsonSettings(x) => { + let settings = x.settings.context(InvalidRawRegionRequestSnafu { + err: "missing settings in SetJsonSettings", + })?; + AlterKind::SetJsonSettings { + column_name: x.column_name, + settings: json_settings_from_proto(settings)?, + } + } alter_request::Kind::DropColumns(x) => { let names = x.drop_columns.into_iter().map(|x| x.name).collect(); AlterKind::DropColumns { names } @@ -1484,6 +1546,36 @@ impl From for ModifyColumnType { } } +fn json_settings_from_proto(settings: v1::JsonSettings) -> Result { + let type_hints = settings + .type_hints + .into_iter() + .map(|hint| { + let wrapper = ColumnDataTypeWrapper::try_new(hint.data_type, hint.datatype_extension) + .map_err(|err| { + InvalidRawRegionRequestSnafu { + err: err.to_string(), + } + .build() + })?; + let data_type = ConcreteDataType::from(wrapper); + + Ok(JsonTypeHint { + path: hint.path, + data_type, + inverted_index: false, + }) + }) + .collect::>>()?; + + JsonSettings::try_new(type_hints, settings.max_auto_expanded_paths).map_err(|err| { + InvalidRawRegionRequestSnafu { + err: err.to_string(), + } + .build() + }) +} + /// Region option changes used by ALTER requests. /// /// This type is serialized for request persistence. Keep future changes backward @@ -1953,8 +2045,11 @@ mod tests { use api::v1::region::RegionColumnDef; use api::v1::{ColumnDataType, ColumnDef}; use common_time::range::TimestampRange; + use datatypes::extension::json::{Json2ExtensionType, JsonMetadata}; + use datatypes::json::JsonSettings; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{ColumnSchema, FulltextAnalyzer, FulltextBackend}; + use datatypes::types::JsonType; use super::*; use crate::metadata::RegionMetadataBuilder; @@ -2364,6 +2459,34 @@ mod tests { }, } ); + + let request = RegionAlterRequest::try_from(AlterRequest { + region_id: 0, + schema_version: 1, + kind: Some(alter_request::Kind::SetJsonSettings(v1::SetJsonSettings { + column_name: "payload".to_string(), + settings: Some(v1::JsonSettings { + type_hints: vec![v1::JsonTypeHint { + path: vec!["service".to_string()], + data_type: ColumnDataType::String as i32, + datatype_extension: None, + }], + max_auto_expanded_paths: Some(10), + }), + })), + }) + .unwrap(); + + let AlterKind::SetJsonSettings { + column_name, + settings, + } = request.kind + else { + unreachable!() + }; + assert_eq!("payload", column_name); + assert_eq!(Some(10), settings.max_auto_expanded_paths()); + assert_eq!(1, settings.type_hints().len()); } #[test] @@ -2446,6 +2569,15 @@ mod tests { builder.build().unwrap() } + fn json2_column_schema(name: &str, settings: JsonSettings) -> ColumnSchema { + let mut column_schema = + ColumnSchema::new(name, ConcreteDataType::Json(JsonType::null()), true); + column_schema.with_extension_type(&Json2ExtensionType::new(std::sync::Arc::new( + JsonMetadata::new(settings), + ))); + column_schema + } + #[test] fn test_add_column_validate() { let metadata = new_metadata(); @@ -2730,6 +2862,40 @@ mod tests { assert!(kind.need_alter(&metadata)); } + #[test] + fn test_validate_set_json_settings() { + let mut metadata = new_metadata(); + let current_schema = json2_column_schema("field_0", JsonSettings::new_v2()); + let target_settings = JsonSettings::try_new(vec![], Some(10)).unwrap(); + metadata + .column_metadatas + .iter_mut() + .find(|column| column.column_schema.name == "field_0") + .unwrap() + .column_schema = current_schema.clone(); + + let kind = AlterKind::SetJsonSettings { + column_name: "field_0".to_string(), + settings: target_settings.clone(), + }; + kind.validate(&metadata).unwrap(); + assert!(kind.need_alter(&metadata)); + + let no_op = AlterKind::SetJsonSettings { + column_name: "field_0".to_string(), + settings: JsonSettings::new_v2(), + }; + no_op.validate(&metadata).unwrap(); + assert!(!no_op.need_alter(&metadata)); + + AlterKind::SetJsonSettings { + column_name: "tag_0".to_string(), + settings: target_settings, + } + .validate(&metadata) + .unwrap_err(); + } + #[test] fn test_validate_add_columns() { let kind = AlterKind::AddColumns { diff --git a/src/table/src/metadata.rs b/src/table/src/metadata.rs index 5916dd7ab49..d91771b9c5f 100644 --- a/src/table/src/metadata.rs +++ b/src/table/src/metadata.rs @@ -22,6 +22,7 @@ use common_query::AddColumnLocation; use datafusion_expr::TableProviderFilterPushDown; use datatypes::error::time_index_not_widening_error; pub use datatypes::error::{Error as ConvertError, Result as ConvertResult}; +use datatypes::extension::json::json2_metadata_with_updated_settings; use datatypes::schema::{ ColumnSchema, FulltextOptions, Schema, SchemaBuilder, SchemaRef, SkippingIndexOptions, }; @@ -41,9 +42,9 @@ use crate::error::{self, Result}; use crate::requests::{ AddColumnRequest, AlterKind, AnnotationContext, AnnotationFamily, AnnotationValidationError, ModifyColumnTypeRequest, REPARTITION_COLUMN_HINT_KEY, REPARTITION_PARTITION_NUM_HINT_KEY, - SetDefaultRequest, SetIndexOption, TableOptions, UnsetIndexOption, has_stable_string_form, - parse_entity_columns, parse_entity_option_key, validate_and_normalize_annotation, - validate_annotation_keys, + SetDefaultRequest, SetIndexOption, SetJsonSettingsRequest, TableOptions, UnsetIndexOption, + has_stable_string_form, parse_entity_columns, parse_entity_option_key, + validate_and_normalize_annotation, validate_annotation_keys, }; use crate::table_reference::TableReference; @@ -332,6 +333,7 @@ impl TableMeta { AlterKind::ModifyColumnTypes { columns } => { self.modify_column_types(table_name, columns) } + AlterKind::SetJsonSettings { request } => self.set_json_settings(table_name, request), // No need to rebuild table meta when renaming tables. AlterKind::RenameTable { .. } => Ok(self.new_meta_builder()), AlterKind::SetTableOptions { options } => self.set_table_options(options), @@ -1187,6 +1189,77 @@ impl TableMeta { Ok(meta_builder) } + fn set_json_settings( + &self, + table_name: &str, + req: &SetJsonSettingsRequest, + ) -> Result { + let table_schema = &self.schema; + let idx = table_schema + .column_index_by_name(&req.column_name) + .with_context(|| error::ColumnNotExistsSnafu { + column_name: &req.column_name, + table_name, + })?; + let col = &table_schema.column_schemas()[idx]; + + ensure!( + !self.primary_key_indices.contains(&idx) && table_schema.timestamp_index() != Some(idx), + error::InvalidAlterRequestSnafu { + table: table_name, + err: format!( + "Not allowed to change JSON settings for key or timestamp column '{}'", + col.name + ), + } + ); + ensure!( + col.data_type.is_json2(), + error::InvalidAlterRequestSnafu { + table: table_name, + err: format!("column '{}' is not a JSON2 column", col.name), + } + ); + + let target_metadata = + json2_metadata_with_updated_settings(col.metadata(), req.settings.clone()).map_err( + |err| { + error::InvalidAlterRequestSnafu { + table: table_name, + err: err.to_string(), + } + .build() + }, + )?; + + let mut cols = table_schema.column_schemas().to_vec(); + *cols[idx].mut_metadata() = target_metadata; + + let mut builder = SchemaBuilder::try_from_columns(cols) + .with_context(|_| error::SchemaBuildSnafu { + msg: format!("Failed to convert column schemas into schema for table {table_name}"), + })? + .version(table_schema.version() + 1); + + for (k, v) in table_schema.metadata().iter() { + builder = builder.add_metadata(k, v); + } + + let new_schema = builder.build().with_context(|_| error::SchemaBuildSnafu { + msg: format!( + "Table {table_name} cannot change JSON settings for column {}", + req.column_name + ), + })?; + + let mut meta_builder = self.new_meta_builder(); + let _ = meta_builder + .schema(Arc::new(new_schema)) + .primary_key_indices(self.primary_key_indices.clone()); + + Ok(meta_builder) + } + /// Split requests into different groups using column location info. fn split_requests_by_column_location<'a>( &self, @@ -1603,9 +1676,12 @@ mod tests { use common_error::ext::ErrorExt; use common_error::status_code::StatusCode; use datatypes::data_type::ConcreteDataType; + use datatypes::extension::json::{Json2ExtensionType, JsonMetadata}; + use datatypes::json::JsonSettings; use datatypes::schema::{ ColumnSchema, FulltextAnalyzer, FulltextBackend, Schema, SchemaBuilder, }; + use datatypes::types::JsonType; use super::*; use crate::Error; @@ -1700,6 +1776,90 @@ mod tests { builder.build().unwrap() } + #[test] + fn test_set_json_settings() { + let schema = Arc::new( + SchemaBuilder::try_from_columns(vec![ + ColumnSchema::new("col1", ConcreteDataType::int32_datatype(), true), + ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ) + .with_time_index(true), + json2_column_schema_v1("payload", JsonSettings::new_v2()), + ]) + .unwrap() + .version(123) + .build() + .unwrap(), + ); + let meta = TableMetaBuilder::empty() + .schema(schema) + .primary_key_indices(vec![0]) + .engine("engine") + .next_column_id(3) + .build() + .unwrap(); + + let alter_kind = AlterKind::SetJsonSettings { + request: SetJsonSettingsRequest { + column_name: "payload".to_string(), + settings: JsonSettings::try_new(vec![], Some(10)).unwrap(), + }, + }; + let new_meta = meta + .builder_with_alter_kind("my_table", &alter_kind) + .unwrap() + .build() + .unwrap(); + + let payload = new_meta.schema.column_schema_by_name("payload").unwrap(); + let json_metadata: JsonMetadata = + serde_json::from_str(payload.metadata().get("ARROW:extension:metadata").unwrap()) + .unwrap(); + assert!(json_metadata.is_version_2()); + assert_eq!( + Some(10), + json_metadata.json_settings().max_auto_expanded_paths() + ); + assert_eq!(124, new_meta.schema.version()); + assert_eq!(meta.primary_key_indices, new_meta.primary_key_indices); + } + + fn json2_column_schema_v1(name: &str, settings: JsonSettings) -> ColumnSchema { + let mut column_schema = + ColumnSchema::new(name, ConcreteDataType::Json(JsonType::null()), true); + column_schema.with_extension_type(&Json2ExtensionType::new(Arc::new( + JsonMetadata::new_v1(settings), + ))); + column_schema + } + + #[test] + fn test_set_json_settings_rejects_non_json2_column() { + let schema = Arc::new(new_test_schema()); + let meta = TableMetaBuilder::empty() + .schema(schema) + .primary_key_indices(vec![0]) + .engine("engine") + .next_column_id(3) + .build() + .unwrap(); + let alter_kind = AlterKind::SetJsonSettings { + request: SetJsonSettingsRequest { + column_name: "col2".to_string(), + settings: JsonSettings::new_v2(), + }, + }; + + let err = match meta.builder_with_alter_kind("my_table", &alter_kind) { + Ok(_) => panic!("expected modifying JSON settings on a non-JSON2 column to fail"), + Err(err) => err, + }; + assert!(err.to_string().contains("is not a JSON2 column")); + } + #[test] fn test_modify_time_index_column_type() { let schema = Arc::new(new_test_schema()); diff --git a/src/table/src/requests.rs b/src/table/src/requests.rs index 1ca99268d67..2c519b32a5e 100644 --- a/src/table/src/requests.rs +++ b/src/table/src/requests.rs @@ -25,6 +25,7 @@ use common_query::AddColumnLocation; use common_time::TimeToLive; use common_time::range::TimestampRange; use datatypes::data_type::ConcreteDataType; +use datatypes::json::JsonSettings; use datatypes::prelude::VectorRef; use datatypes::schema::{ ColumnDefaultConstraint, ColumnSchema, FulltextOptions, Schema, SkippingIndexOptions, @@ -371,6 +372,13 @@ pub struct ModifyColumnTypeRequest { pub target_type: ConcreteDataType, } +/// Set JSON2 settings request. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SetJsonSettingsRequest { + pub column_name: String, + pub settings: JsonSettings, +} + /// A family of annotation table options: pure metadata markers that no region /// consumes. Setting or unsetting them only rewrites the table's /// `extra_options`, so the alter skips region dispatch entirely. @@ -661,6 +669,9 @@ pub enum AlterKind { ModifyColumnTypes { columns: Vec, }, + SetJsonSettings { + request: SetJsonSettingsRequest, + }, RenameTable { new_table_name: String, }, diff --git a/tests/cases/standalone/common/types/json/json2_alter.result b/tests/cases/standalone/common/types/json/json2_alter.result index c23f9e2533f..ab9e39dc777 100644 --- a/tests/cases/standalone/common/types/json/json2_alter.result +++ b/tests/cases/standalone/common/types/json/json2_alter.result @@ -2,37 +2,110 @@ CREATE TABLE application_logs ( ts TIMESTAMP TIME INDEX, attrs JSON2 ) WITH ( - 'append_mode' = 'true' + 'append_mode' = 'true', + 'sst_format' = 'flat' ); Affected Rows: 0 -ALTER TABLE application_logs - MODIFY COLUMN attrs JSON2; +INSERT INTO application_logs VALUES + (1, '{"service":"frontend","duration":12,"trace_id":"before-1"}'), + (2, '{"service":"ingester","duration":34,"trace_id":"before-2"}'); -Error: 1001(Unsupported), Not supported: ALTER TABLE MODIFY COLUMN to JSON2 type - -ALTER TABLE application_logs - MODIFY COLUMN attrs JSON2 (); - -Error: 1001(Unsupported), Not supported: ALTER TABLE MODIFY COLUMN to JSON2 type +Affected Rows: 2 ALTER TABLE application_logs MODIFY COLUMN attrs JSON2 ( - max_auto_expanded_paths = 2000 + max_auto_expanded_paths = 2000, + service STRING, + duration INT64, + trace_id STRING ); -Error: 1001(Unsupported), Not supported: ALTER TABLE MODIFY COLUMN to JSON2 type +Affected Rows: 0 + +SELECT COUNT(*) AS alter_flushed_sst_num +FROM information_schema.ssts_manifest +WHERE visible + AND table_id IN ( + SELECT table_id FROM information_schema.tables + WHERE table_schema = 'public' AND table_name = 'application_logs' + ); + ++-----------------------+ +| alter_flushed_sst_num | ++-----------------------+ +| 1 | ++-----------------------+ + +INSERT INTO application_logs VALUES + (3, '{"service":"frontend","duration":56,"trace_id":"after-1"}'), + (4, '{"service":"compactor","trace_id":"after-2"}'); + +Affected Rows: 2 + +SELECT ts, attrs.service, attrs.duration, attrs.trace_id +FROM application_logs +ORDER BY ts; + ++-------------------------+-------------------------------------------------------------------+-----------------------------------------------------------------+--------------------------------------------------------------------+ +| ts | json_get(application_logs.attrs,Utf8("$.service"),Utf8View(NULL)) | json_get(application_logs.attrs,Utf8("$.duration"),Int64(NULL)) | json_get(application_logs.attrs,Utf8("$.trace_id"),Utf8View(NULL)) | ++-------------------------+-------------------------------------------------------------------+-----------------------------------------------------------------+--------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | frontend | 12 | before-1 | +| 1970-01-01T00:00:00.002 | ingester | 34 | before-2 | +| 1970-01-01T00:00:00.003 | frontend | 56 | after-1 | +| 1970-01-01T00:00:00.004 | compactor | | after-2 | ++-------------------------+-------------------------------------------------------------------+-----------------------------------------------------------------+--------------------------------------------------------------------+ + +ADMIN FLUSH_TABLE('application_logs'); + ++---------------------------------------+ +| ADMIN FLUSH_TABLE('application_logs') | ++---------------------------------------+ +| 0 | ++---------------------------------------+ + +SELECT ts, attrs.service, attrs.duration, attrs.trace_id +FROM application_logs +ORDER BY ts; + ++-------------------------+-------------------------------------------------------------------+-----------------------------------------------------------------+--------------------------------------------------------------------+ +| ts | json_get(application_logs.attrs,Utf8("$.service"),Utf8View(NULL)) | json_get(application_logs.attrs,Utf8("$.duration"),Int64(NULL)) | json_get(application_logs.attrs,Utf8("$.trace_id"),Utf8View(NULL)) | ++-------------------------+-------------------------------------------------------------------+-----------------------------------------------------------------+--------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | frontend | 12 | before-1 | +| 1970-01-01T00:00:00.002 | ingester | 34 | before-2 | +| 1970-01-01T00:00:00.003 | frontend | 56 | after-1 | +| 1970-01-01T00:00:00.004 | compactor | | after-2 | ++-------------------------+-------------------------------------------------------------------+-----------------------------------------------------------------+--------------------------------------------------------------------+ ALTER TABLE application_logs MODIFY COLUMN attrs JSON2 ( trace_id STRING, user.id STRING, user.name STRING, - request_id STRING INVERTED INDEX + request_id STRING ); -Error: 1001(Unsupported), Not supported: ALTER TABLE MODIFY COLUMN to JSON2 type +Affected Rows: 0 + +INSERT INTO application_logs VALUES + (5, '{"trace_id":"after-3","user":{"id":"u1","name":"alice"},"request_id":"r1"}'); + +Affected Rows: 1 + +SELECT ts, attrs.trace_id, attrs.user.id, attrs.user.name, attrs.request_id +FROM application_logs +ORDER BY ts; + ++-------------------------+--------------------------------------------------------------------+-------------------------------------------------------------------+---------------------------------------------------------------------+----------------------------------------------------------------------+ +| ts | json_get(application_logs.attrs,Utf8("$.trace_id"),Utf8View(NULL)) | json_get(application_logs.attrs,Utf8("$.user.id"),Utf8View(NULL)) | json_get(application_logs.attrs,Utf8("$.user.name"),Utf8View(NULL)) | json_get(application_logs.attrs,Utf8("$.request_id"),Utf8View(NULL)) | ++-------------------------+--------------------------------------------------------------------+-------------------------------------------------------------------+---------------------------------------------------------------------+----------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | before-1 | | | | +| 1970-01-01T00:00:00.002 | before-2 | | | | +| 1970-01-01T00:00:00.003 | after-1 | | | | +| 1970-01-01T00:00:00.004 | after-2 | | | | +| 1970-01-01T00:00:00.005 | after-3 | u1 | alice | r1 | ++-------------------------+--------------------------------------------------------------------+-------------------------------------------------------------------+---------------------------------------------------------------------+----------------------------------------------------------------------+ ALTER TABLE application_logs MODIFY COLUMN attrs JSON2 ( @@ -40,10 +113,38 @@ ALTER TABLE application_logs trace_id STRING, user.id STRING, user.name STRING, - request_id STRING INVERTED INDEX + request_id STRING ); -Error: 1001(Unsupported), Not supported: ALTER TABLE MODIFY COLUMN to JSON2 type +Affected Rows: 0 + +INSERT INTO application_logs VALUES + (6, '{"trace_id":"after-4","user":{"id":"u2"},"request_id":"r2"}'); + +Affected Rows: 1 + +SELECT ts, attrs.trace_id, attrs.user.id, attrs.user.name, attrs.request_id +FROM application_logs +ORDER BY ts; + ++-------------------------+--------------------------------------------------------------------+-------------------------------------------------------------------+---------------------------------------------------------------------+----------------------------------------------------------------------+ +| ts | json_get(application_logs.attrs,Utf8("$.trace_id"),Utf8View(NULL)) | json_get(application_logs.attrs,Utf8("$.user.id"),Utf8View(NULL)) | json_get(application_logs.attrs,Utf8("$.user.name"),Utf8View(NULL)) | json_get(application_logs.attrs,Utf8("$.request_id"),Utf8View(NULL)) | ++-------------------------+--------------------------------------------------------------------+-------------------------------------------------------------------+---------------------------------------------------------------------+----------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | before-1 | | | | +| 1970-01-01T00:00:00.002 | before-2 | | | | +| 1970-01-01T00:00:00.003 | after-1 | | | | +| 1970-01-01T00:00:00.004 | after-2 | | | | +| 1970-01-01T00:00:00.005 | after-3 | u1 | alice | r1 | +| 1970-01-01T00:00:00.006 | after-4 | u2 | | r2 | ++-------------------------+--------------------------------------------------------------------+-------------------------------------------------------------------+---------------------------------------------------------------------+----------------------------------------------------------------------+ + +ADMIN FLUSH_TABLE('application_logs'); + ++---------------------------------------+ +| ADMIN FLUSH_TABLE('application_logs') | ++---------------------------------------+ +| 0 | ++---------------------------------------+ DROP TABLE application_logs; diff --git a/tests/cases/standalone/common/types/json/json2_alter.sql b/tests/cases/standalone/common/types/json/json2_alter.sql index c89556c3485..7fce02b76ee 100644 --- a/tests/cases/standalone/common/types/json/json2_alter.sql +++ b/tests/cases/standalone/common/types/json/json2_alter.sql @@ -2,35 +2,75 @@ CREATE TABLE application_logs ( ts TIMESTAMP TIME INDEX, attrs JSON2 ) WITH ( - 'append_mode' = 'true' + 'append_mode' = 'true', + 'sst_format' = 'flat' ); -ALTER TABLE application_logs - MODIFY COLUMN attrs JSON2; - -ALTER TABLE application_logs - MODIFY COLUMN attrs JSON2 (); +INSERT INTO application_logs VALUES + (1, '{"service":"frontend","duration":12,"trace_id":"before-1"}'), + (2, '{"service":"ingester","duration":34,"trace_id":"before-2"}'); ALTER TABLE application_logs MODIFY COLUMN attrs JSON2 ( - max_auto_expanded_paths = 2000 + max_auto_expanded_paths = 2000, + service STRING, + duration INT64, + trace_id STRING ); +SELECT COUNT(*) AS alter_flushed_sst_num +FROM information_schema.ssts_manifest +WHERE visible + AND table_id IN ( + SELECT table_id FROM information_schema.tables + WHERE table_schema = 'public' AND table_name = 'application_logs' + ); + +INSERT INTO application_logs VALUES + (3, '{"service":"frontend","duration":56,"trace_id":"after-1"}'), + (4, '{"service":"compactor","trace_id":"after-2"}'); + +SELECT ts, attrs.service, attrs.duration, attrs.trace_id +FROM application_logs +ORDER BY ts; + +ADMIN FLUSH_TABLE('application_logs'); + +SELECT ts, attrs.service, attrs.duration, attrs.trace_id +FROM application_logs +ORDER BY ts; + ALTER TABLE application_logs MODIFY COLUMN attrs JSON2 ( trace_id STRING, user.id STRING, user.name STRING, - request_id STRING INVERTED INDEX + request_id STRING ); +INSERT INTO application_logs VALUES + (5, '{"trace_id":"after-3","user":{"id":"u1","name":"alice"},"request_id":"r1"}'); + +SELECT ts, attrs.trace_id, attrs.user.id, attrs.user.name, attrs.request_id +FROM application_logs +ORDER BY ts; + ALTER TABLE application_logs MODIFY COLUMN attrs JSON2 ( max_auto_expanded_paths = 2000, trace_id STRING, user.id STRING, user.name STRING, - request_id STRING INVERTED INDEX + request_id STRING ); +INSERT INTO application_logs VALUES + (6, '{"trace_id":"after-4","user":{"id":"u2"},"request_id":"r2"}'); + +SELECT ts, attrs.trace_id, attrs.user.id, attrs.user.name, attrs.request_id +FROM application_logs +ORDER BY ts; + +ADMIN FLUSH_TABLE('application_logs'); + DROP TABLE application_logs; diff --git a/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction.result b/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction.result new file mode 100644 index 00000000000..284b7e6f850 --- /dev/null +++ b/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction.result @@ -0,0 +1,83 @@ +CREATE TABLE json2_alter_type_hints_compaction ( + ts TIMESTAMP TIME INDEX, + attrs JSON2 +) WITH ( + 'append_mode' = 'true', + 'sst_format' = 'flat' +); + +Affected Rows: 0 + +INSERT INTO json2_alter_type_hints_compaction VALUES + (1, '{"kind":1,"message":"before-has"}'), + (2, '{"message":"before-missing"}'); + +Affected Rows: 2 + +ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction'); + ++--------------------------------------------------------+ +| ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction') | ++--------------------------------------------------------+ +| 0 | ++--------------------------------------------------------+ + +ALTER TABLE json2_alter_type_hints_compaction + MODIFY COLUMN attrs JSON2 ( + kind INT64 + ); + +Affected Rows: 0 + +INSERT INTO json2_alter_type_hints_compaction VALUES + (3, '{"kind":3,"message":"after-has"}'), + (4, '{"message":"after-missing"}'); + +Affected Rows: 2 + +ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction'); + ++--------------------------------------------------------+ +| ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction') | ++--------------------------------------------------------+ +| 0 | ++--------------------------------------------------------+ + +SELECT ts, attrs.kind, attrs.message +FROM json2_alter_type_hints_compaction +ORDER BY ts; + ++-------------------------+------------------------------------------------------------------------------+---------------------------------------------------------------------+ +| ts | json_get(json2_alter_type_hints_compaction.attrs,Utf8("$.kind"),Int64(NULL)) | json_get(json2_alter_type_hints_compaction.attrs,Utf8("$.message")) | ++-------------------------+------------------------------------------------------------------------------+---------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | 1 | before-has | +| 1970-01-01T00:00:00.002 | | before-missing | +| 1970-01-01T00:00:00.003 | 3 | after-has | +| 1970-01-01T00:00:00.004 | | after-missing | ++-------------------------+------------------------------------------------------------------------------+---------------------------------------------------------------------+ + +ADMIN COMPACT_TABLE('json2_alter_type_hints_compaction', 'swcs', '86400'); + ++---------------------------------------------------------------------------+ +| ADMIN COMPACT_TABLE('json2_alter_type_hints_compaction', 'swcs', '86400') | ++---------------------------------------------------------------------------+ +| 0 | ++---------------------------------------------------------------------------+ + +SELECT ts, attrs.kind, attrs.message +FROM json2_alter_type_hints_compaction +ORDER BY ts; + ++-------------------------+------------------------------------------------------------------------------+---------------------------------------------------------------------+ +| ts | json_get(json2_alter_type_hints_compaction.attrs,Utf8("$.kind"),Int64(NULL)) | json_get(json2_alter_type_hints_compaction.attrs,Utf8("$.message")) | ++-------------------------+------------------------------------------------------------------------------+---------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | 1 | before-has | +| 1970-01-01T00:00:00.002 | | before-missing | +| 1970-01-01T00:00:00.003 | 3 | after-has | +| 1970-01-01T00:00:00.004 | | after-missing | ++-------------------------+------------------------------------------------------------------------------+---------------------------------------------------------------------+ + +DROP TABLE json2_alter_type_hints_compaction; + +Affected Rows: 0 + diff --git a/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction.sql b/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction.sql new file mode 100644 index 00000000000..45d859d4ece --- /dev/null +++ b/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction.sql @@ -0,0 +1,36 @@ +CREATE TABLE json2_alter_type_hints_compaction ( + ts TIMESTAMP TIME INDEX, + attrs JSON2 +) WITH ( + 'append_mode' = 'true', + 'sst_format' = 'flat' +); + +INSERT INTO json2_alter_type_hints_compaction VALUES + (1, '{"kind":1,"message":"before-has"}'), + (2, '{"message":"before-missing"}'); + +ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction'); + +ALTER TABLE json2_alter_type_hints_compaction + MODIFY COLUMN attrs JSON2 ( + kind INT64 + ); + +INSERT INTO json2_alter_type_hints_compaction VALUES + (3, '{"kind":3,"message":"after-has"}'), + (4, '{"message":"after-missing"}'); + +ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction'); + +SELECT ts, attrs.kind, attrs.message +FROM json2_alter_type_hints_compaction +ORDER BY ts; + +ADMIN COMPACT_TABLE('json2_alter_type_hints_compaction', 'swcs', '86400'); + +SELECT ts, attrs.kind, attrs.message +FROM json2_alter_type_hints_compaction +ORDER BY ts; + +DROP TABLE json2_alter_type_hints_compaction; diff --git a/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction_conversion.result b/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction_conversion.result new file mode 100644 index 00000000000..5ba2cf4700a --- /dev/null +++ b/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction_conversion.result @@ -0,0 +1,78 @@ +CREATE TABLE json2_alter_type_hints_compaction_conversion ( + ts TIMESTAMP TIME INDEX, + attrs JSON2 +) WITH ( + 'append_mode' = 'true', + 'sst_format' = 'flat' +); + +Affected Rows: 0 + +INSERT INTO json2_alter_type_hints_compaction_conversion VALUES + (1, '{"to_string":123,"to_int":"456"}'); + +Affected Rows: 1 + +ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction_conversion'); + ++-------------------------------------------------------------------+ +| ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction_conversion') | ++-------------------------------------------------------------------+ +| 0 | ++-------------------------------------------------------------------+ + +ALTER TABLE json2_alter_type_hints_compaction_conversion + MODIFY COLUMN attrs JSON2 ( + to_string STRING, + to_int INT64 + ); + +Affected Rows: 0 + +INSERT INTO json2_alter_type_hints_compaction_conversion VALUES + (2, '{"to_string":"after","to_int":789}'); + +Affected Rows: 1 + +ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction_conversion'); + ++-------------------------------------------------------------------+ +| ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction_conversion') | ++-------------------------------------------------------------------+ +| 0 | ++-------------------------------------------------------------------+ + +SELECT ts, attrs.to_string, attrs.to_int +FROM json2_alter_type_hints_compaction_conversion +ORDER BY ts; + ++-------------------------+-------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------+ +| ts | json_get(json2_alter_type_hints_compaction_conversion.attrs,Utf8("$.to_string"),Utf8View(NULL)) | json_get(json2_alter_type_hints_compaction_conversion.attrs,Utf8("$.to_int"),Int64(NULL)) | ++-------------------------+-------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | 123 | 456 | +| 1970-01-01T00:00:00.002 | after | 789 | ++-------------------------+-------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------+ + +ADMIN COMPACT_TABLE('json2_alter_type_hints_compaction_conversion', 'swcs', '86400'); + ++--------------------------------------------------------------------------------------+ +| ADMIN COMPACT_TABLE('json2_alter_type_hints_compaction_conversion', 'swcs', '86400') | ++--------------------------------------------------------------------------------------+ +| 0 | ++--------------------------------------------------------------------------------------+ + +SELECT ts, attrs.to_string, attrs.to_int +FROM json2_alter_type_hints_compaction_conversion +ORDER BY ts; + ++-------------------------+-------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------+ +| ts | json_get(json2_alter_type_hints_compaction_conversion.attrs,Utf8("$.to_string"),Utf8View(NULL)) | json_get(json2_alter_type_hints_compaction_conversion.attrs,Utf8("$.to_int"),Int64(NULL)) | ++-------------------------+-------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | 123 | 456 | +| 1970-01-01T00:00:00.002 | after | 789 | ++-------------------------+-------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------------+ + +DROP TABLE json2_alter_type_hints_compaction_conversion; + +Affected Rows: 0 + diff --git a/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction_conversion.sql b/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction_conversion.sql new file mode 100644 index 00000000000..de9f211cfb5 --- /dev/null +++ b/tests/cases/standalone/common/types/json/json2_alter_type_hints_compaction_conversion.sql @@ -0,0 +1,35 @@ +CREATE TABLE json2_alter_type_hints_compaction_conversion ( + ts TIMESTAMP TIME INDEX, + attrs JSON2 +) WITH ( + 'append_mode' = 'true', + 'sst_format' = 'flat' +); + +INSERT INTO json2_alter_type_hints_compaction_conversion VALUES + (1, '{"to_string":123,"to_int":"456"}'); + +ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction_conversion'); + +ALTER TABLE json2_alter_type_hints_compaction_conversion + MODIFY COLUMN attrs JSON2 ( + to_string STRING, + to_int INT64 + ); + +INSERT INTO json2_alter_type_hints_compaction_conversion VALUES + (2, '{"to_string":"after","to_int":789}'); + +ADMIN FLUSH_TABLE('json2_alter_type_hints_compaction_conversion'); + +SELECT ts, attrs.to_string, attrs.to_int +FROM json2_alter_type_hints_compaction_conversion +ORDER BY ts; + +ADMIN COMPACT_TABLE('json2_alter_type_hints_compaction_conversion', 'swcs', '86400'); + +SELECT ts, attrs.to_string, attrs.to_int +FROM json2_alter_type_hints_compaction_conversion +ORDER BY ts; + +DROP TABLE json2_alter_type_hints_compaction_conversion; diff --git a/tests/cases/standalone/common/types/json/json2_compaction_null_type_hint_mismatch.result b/tests/cases/standalone/common/types/json/json2_compaction_null_type_hint_mismatch.result new file mode 100644 index 00000000000..6e0f833ccd6 --- /dev/null +++ b/tests/cases/standalone/common/types/json/json2_compaction_null_type_hint_mismatch.result @@ -0,0 +1,80 @@ +CREATE TABLE json2_compaction_null_type_hint_mismatch ( + ts TIMESTAMP TIME INDEX, + attrs JSON2 +) WITH ( + 'append_mode' = 'true', + 'sst_format' = 'flat' +); + +Affected Rows: 0 + +INSERT INTO json2_compaction_null_type_hint_mismatch VALUES + (1, '{"kind":1}'), + (2, '{"kind":"invalid","message":"keep"}'); + +Affected Rows: 2 + +ADMIN FLUSH_TABLE('json2_compaction_null_type_hint_mismatch'); + ++---------------------------------------------------------------+ +| ADMIN FLUSH_TABLE('json2_compaction_null_type_hint_mismatch') | ++---------------------------------------------------------------+ +| 0 | ++---------------------------------------------------------------+ + +ALTER TABLE json2_compaction_null_type_hint_mismatch + MODIFY COLUMN attrs JSON2 ( + kind INT64 + ); + +Affected Rows: 0 + +INSERT INTO json2_compaction_null_type_hint_mismatch VALUES + (3, '{"kind":3}'); + +Affected Rows: 1 + +ADMIN FLUSH_TABLE('json2_compaction_null_type_hint_mismatch'); + ++---------------------------------------------------------------+ +| ADMIN FLUSH_TABLE('json2_compaction_null_type_hint_mismatch') | ++---------------------------------------------------------------+ +| 0 | ++---------------------------------------------------------------+ + +SELECT ts, attrs.kind, attrs.message +FROM json2_compaction_null_type_hint_mismatch +ORDER BY ts; + ++-------------------------+-------------------------------------------------------------------------------------+----------------------------------------------------------------------------+ +| ts | json_get(json2_compaction_null_type_hint_mismatch.attrs,Utf8("$.kind"),Int64(NULL)) | json_get(json2_compaction_null_type_hint_mismatch.attrs,Utf8("$.message")) | ++-------------------------+-------------------------------------------------------------------------------------+----------------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | 1 | | +| 1970-01-01T00:00:00.002 | | keep | +| 1970-01-01T00:00:00.003 | 3 | | ++-------------------------+-------------------------------------------------------------------------------------+----------------------------------------------------------------------------+ + +ADMIN COMPACT_TABLE('json2_compaction_null_type_hint_mismatch', 'swcs', '86400'); + ++----------------------------------------------------------------------------------+ +| ADMIN COMPACT_TABLE('json2_compaction_null_type_hint_mismatch', 'swcs', '86400') | ++----------------------------------------------------------------------------------+ +| 0 | ++----------------------------------------------------------------------------------+ + +SELECT ts, attrs.kind, attrs.message +FROM json2_compaction_null_type_hint_mismatch +ORDER BY ts; + ++-------------------------+-------------------------------------------------------------------------------------+----------------------------------------------------------------------------+ +| ts | json_get(json2_compaction_null_type_hint_mismatch.attrs,Utf8("$.kind"),Int64(NULL)) | json_get(json2_compaction_null_type_hint_mismatch.attrs,Utf8("$.message")) | ++-------------------------+-------------------------------------------------------------------------------------+----------------------------------------------------------------------------+ +| 1970-01-01T00:00:00.001 | 1 | | +| 1970-01-01T00:00:00.002 | | keep | +| 1970-01-01T00:00:00.003 | 3 | | ++-------------------------+-------------------------------------------------------------------------------------+----------------------------------------------------------------------------+ + +DROP TABLE json2_compaction_null_type_hint_mismatch; + +Affected Rows: 0 + diff --git a/tests/cases/standalone/common/types/json/json2_compaction_null_type_hint_mismatch.sql b/tests/cases/standalone/common/types/json/json2_compaction_null_type_hint_mismatch.sql new file mode 100644 index 00000000000..7ebff9cceb1 --- /dev/null +++ b/tests/cases/standalone/common/types/json/json2_compaction_null_type_hint_mismatch.sql @@ -0,0 +1,35 @@ +CREATE TABLE json2_compaction_null_type_hint_mismatch ( + ts TIMESTAMP TIME INDEX, + attrs JSON2 +) WITH ( + 'append_mode' = 'true', + 'sst_format' = 'flat' +); + +INSERT INTO json2_compaction_null_type_hint_mismatch VALUES + (1, '{"kind":1}'), + (2, '{"kind":"invalid","message":"keep"}'); + +ADMIN FLUSH_TABLE('json2_compaction_null_type_hint_mismatch'); + +ALTER TABLE json2_compaction_null_type_hint_mismatch + MODIFY COLUMN attrs JSON2 ( + kind INT64 + ); + +INSERT INTO json2_compaction_null_type_hint_mismatch VALUES + (3, '{"kind":3}'); + +ADMIN FLUSH_TABLE('json2_compaction_null_type_hint_mismatch'); + +SELECT ts, attrs.kind, attrs.message +FROM json2_compaction_null_type_hint_mismatch +ORDER BY ts; + +ADMIN COMPACT_TABLE('json2_compaction_null_type_hint_mismatch', 'swcs', '86400'); + +SELECT ts, attrs.kind, attrs.message +FROM json2_compaction_null_type_hint_mismatch +ORDER BY ts; + +DROP TABLE json2_compaction_null_type_hint_mismatch;