feat(json2): support altering JSON2 settings (#9029)

* feat(sql): support alter syntax for JSON2 columns

Signed-off-by: fys <fengys1996@gmail.com>

* 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 <fengys1996@gmail.com>
This commit is contained in:
fys
2026-09-22 12:56:26 +00:00
committed by GitHub
parent aa36f74feb
commit 045441e3cc
30 changed files with 1602 additions and 231 deletions
Generated
+2 -1
View File
@@ -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",
+1 -1
View File
@@ -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"
+46 -5
View File
@@ -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<JsonSettings> {
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::<Result<Vec<_>>>()?;
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(),
},
+8
View File
@@ -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 { .. } => {
@@ -333,6 +333,7 @@ fn build_new_table_info(
}
AlterKind::DropColumns { .. }
| AlterKind::ModifyColumnTypes { .. }
| AlterKind::SetJsonSettings { .. }
| AlterKind::SetTableOptions { .. }
| AlterKind::UnsetTableOptions { .. }
| AlterKind::SetAnnotations { .. }
@@ -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
+1
View File
@@ -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"),
+53 -2
View File
@@ -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<Metadata> {
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<T: AsRef<Field>>(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)
+336 -48
View File
@@ -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::<Vec<_>>();
let mut other_hints = other_hints.iter().collect::<Vec<_>>();
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<Json> {
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<Value> {
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<Value> {
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<T>(
}
/// Main encoding function with key path tracking
fn encode_json_with_context(json: Json, context: &mut JsonContext) -> Result<JsonValue> {
fn encode_json_with_context(
json: Json,
context: &mut JsonContext,
policy: TypeHintMismatchPolicy,
) -> Result<JsonValue> {
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<String, Json>,
context: &mut JsonContext<'a>,
policy: TypeHintMismatchPolicy,
) -> Result<JsonValue> {
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<String, JsonVariant>,
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<JsonValue> {
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::<Vec<_>>();
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::<Vec<_>>();
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<Json>,
context: &mut JsonContext<'a>,
policy: TypeHintMismatchPolicy,
) -> Result<JsonValue> {
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<JsonValue> {
fn encode_json_value_with_context(
json: Json,
context: &mut JsonContext,
policy: TypeHintMismatchPolicy,
) -> Result<JsonValue> {
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<dyn std::error::Error>> {
@@ -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();
+23 -82
View File
@@ -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<ArrayRef> {
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<ArrayRef> {
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<Value>, to_type: &DataType) -> Result<ArrayRe
let concrete_type = ConcreteDataType::from_arrow_type(to_type);
let mut builder = concrete_type.create_mutable_vector(values.len());
for value in values {
let value = project_json_value_to_type(value, &concrete_type)?;
let value = coerce_json_value_to_type(value, &concrete_type);
builder.try_push_value_ref(&value.as_value_ref())?;
}
Ok(builder.to_vector().to_arrow_array())
}
fn project_json_value_to_type(value: Value, to_type: &ConcreteDataType) -> Result<GreptimeValue> {
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::<Result<Vec<_>>>()?;
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::<Result<Vec<_>>>()?;
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 }
+1
View File
@@ -227,6 +227,7 @@ impl DataRegion {
| AlterKind::UnsetRegionOptions { keys: _ }
| AlterKind::SetIndexes { options: _ }
| AlterKind::UnsetIndexes { options: _ }
| AlterKind::SetJsonSettings { .. }
| AlterKind::SyncColumns {
column_metadatas: _,
} => {
+10 -3
View File
@@ -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),
+75 -3
View File
@@ -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<ArrayRef> {
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::<Int64Array>()
.unwrap()
.value(0)
);
assert_eq!(
20,
output
.column(1)
.as_any()
.downcast_ref::<Int64Array>()
.unwrap()
.value(1)
);
}
#[test]
fn test_rewrite_rejects_mismatched_output_layout() {
let columns = HashMap::from([(
+84 -25
View File
@@ -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<PbJsonSettings> {
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::<Result<Vec<_>>>()?;
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]
+13 -15
View File
@@ -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!()
};
+37 -1
View File
@@ -88,6 +88,11 @@ pub enum AlterTableOperation {
target_type: DataType,
json2_options: Option<Json2Options>,
},
/// `MODIFY <column_name> JSON2 [json2_options]`
SetJsonSettings {
column_name: Ident,
json2_options: Option<Json2Options>,
},
/// `SET <table attrs key> = <table attr value>`
SetTableOptions {
options: Vec<KeyValueOption>,
@@ -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())
+21 -18
View File
@@ -157,6 +157,26 @@ impl Display for Json2Options {
}
}
impl Json2Options {
pub fn build_json_settings(&self) -> Result<JsonSettings> {
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::<Result<Vec<_>>>()?;
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<String>,
@@ -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::<Result<Vec<_>>>()?;
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<()> {
+1
View File
@@ -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
+36
View File
@@ -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<SetIndexOption>) -> Result<()> {
let mut set_index_map: HashMap<_, Vec<_>> = HashMap::new();
for option in &options {
+166
View File
@@ -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<ModifyColumnType>,
},
/// Set JSON2 settings of a region column.
SetJsonSettings {
/// Column name.
column_name: String,
/// Target JSON2 settings.
settings: JsonSettings,
},
/// Set region options.
SetRegionOptions { options: Vec<SetRegionOption> },
/// 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::<Json2ExtensionType>()
.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<alter_request::Kind> for AlterKind {
@@ -1183,6 +1236,15 @@ impl TryFrom<alter_request::Kind> for AlterKind {
.collect::<Vec<_>>();
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<v1::ModifyColumnType> for ModifyColumnType {
}
}
fn json_settings_from_proto(settings: v1::JsonSettings) -> Result<JsonSettings> {
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::<Result<Vec<_>>>()?;
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 {
+163 -3
View File
@@ -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<TableMetaBuilder> {
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());
+11
View File
@@ -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<ModifyColumnTypeRequest>,
},
SetJsonSettings {
request: SetJsonSettingsRequest,
},
RenameTable {
new_table_name: String,
},
@@ -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;
@@ -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;
@@ -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
@@ -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;
@@ -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
@@ -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;
@@ -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
@@ -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;