diff --git a/src/datatypes/src/error.rs b/src/datatypes/src/error.rs index b34a5ffa9b..1def0d6be8 100644 --- a/src/datatypes/src/error.rs +++ b/src/datatypes/src/error.rs @@ -77,6 +77,13 @@ pub enum Error { location: Location, }, + #[snafu(display("Unimplemented: {feat}"))] + Unimplemented { + feat: String, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display("Failed to parse version in schema meta, value: {}", value))] ParseSchemaVersion { value: String, @@ -312,6 +319,7 @@ impl ErrorExt for Error { use Error::*; match self { UnsupportedOperation { .. } + | Unimplemented { .. } | UnsupportedArrowType { .. } | UnsupportedJsonType { .. } | UnsupportedDefaultExpr { .. } => StatusCode::Unsupported, diff --git a/src/datatypes/src/json.rs b/src/datatypes/src/json.rs index ec9ca4ac78..79fc10b12e 100644 --- a/src/datatypes/src/json.rs +++ b/src/datatypes/src/json.rs @@ -140,6 +140,16 @@ impl JsonSettings { /// Encode a serde_json::Value into a Value::Json using current settings. pub fn encode(&self, json: Json) -> Result { + if let Json::Object(object) = &json + && object.contains_key(JSON2_REMAINDER_FIELD_NAME) + { + return error::InvalidJsonSnafu { + value: format!( + "root object cannot contain reserved field '{JSON2_REMAINDER_FIELD_NAME}'" + ), + } + .fail(); + } let mut context = JsonContext { path: Vec::new(), settings: self, @@ -777,6 +787,21 @@ mod tests { } } + #[test] + fn test_encode_rejects_reserved_remainder_field() -> Result<()> { + let settings = JsonSettings::default(); + let err = settings + .encode(json!({"!__remainder__!": "user-value"})) + .unwrap_err(); + assert!( + err.to_string() + .contains("root object cannot contain reserved field") + ); + + settings.encode(json!({"nested": {"!__remainder__!": "user-value"}}))?; + Ok(()) + } + #[test] fn test_encode_json_object() { let json = json!({ diff --git a/src/datatypes/src/vectors/json/builder.rs b/src/datatypes/src/vectors/json/builder.rs index 125dbdad58..ab3fd6d043 100644 --- a/src/datatypes/src/vectors/json/builder.rs +++ b/src/datatypes/src/vectors/json/builder.rs @@ -24,17 +24,18 @@ use snafu::{ResultExt, ensure}; use crate::data_type::ConcreteDataType; use crate::error::{ - ArrowComputeSnafu, Result, TryFromValueSnafu, UnexpectedSnafu, UnsupportedOperationSnafu, + ArrowComputeSnafu, Result, TryFromValueSnafu, UnexpectedSnafu, UnimplementedSnafu, + UnsupportedOperationSnafu, }; use crate::extension::json::JSON2_REMAINDER_FIELD_NAME; -use crate::json::value::{JsonNumber, JsonVariant, encode_json_variant}; +use crate::json::value::{JsonNumber, JsonVariant, JsonVariantRef, encode_json_variant}; use crate::json::{JSON2_MAX_STRUCTURED_DEPTH, JsonSettings}; use crate::prelude::{ValueRef, Vector, VectorRef}; use crate::types::StructType; use crate::types::json_type::{JsonNativeType, is_include}; use crate::value::{ListValue, StructValue, StructValueRef, Value}; -use crate::vectors::json::variant::{append_json_variant, variant_field}; -use crate::vectors::{Helper, MutableVector, StructVectorBuilder}; +use crate::vectors::json::variant::{append_json_variant, append_json_variant_ref, variant_field}; +use crate::vectors::{Helper, MutableVector, NullVector, StructVectorBuilder}; type JsonObjectValue = BTreeMap; @@ -44,17 +45,25 @@ type JsonObjectValue = BTreeMap; /// Auto-expanding mode always materializes type-hinted paths, selects up to /// `max_auto_expanded_paths` compatible unhinted leaf paths by frequency, and stores /// conflicting or unselected paths in the Variant remainder field. -#[derive(Clone)] pub(crate) struct JsonVectorBuilder { state: JsonVectorBuilderState, } -#[derive(Clone)] enum JsonVectorBuilderState { Legacy { merged_type: JsonNativeType, values: Vec, }, + ExplicitOnly { + /// Paths declared by type hints and stored as dedicated Struct fields. + explicit_type: JsonNativeType, + /// Concrete Struct type used to append explicit values without buffering rows. + struct_type: StructType, + /// Builder for values selected by the explicit type hints. + explicit: StructVectorBuilder, + /// Builder for all values outside the explicit type hints. + remainder: VariantArrayBuilder, + }, AutoExpanding { /// Paths declared by type hints and always stored as dedicated Struct fields. explicit_type: JsonNativeType, @@ -69,6 +78,7 @@ impl JsonVectorBuilderState { fn native_type(&self) -> JsonNativeType { match self { Self::Legacy { merged_type, .. } => merged_type.clone(), + Self::ExplicitOnly { explicit_type, .. } => explicit_type.clone(), Self::AutoExpanding { explicit_type, max_auto_expanded_paths, @@ -80,6 +90,7 @@ impl JsonVectorBuilderState { fn len(&self) -> usize { match self { Self::Legacy { values, .. } | Self::AutoExpanding { values, .. } => values.len(), + Self::ExplicitOnly { explicit, .. } => explicit.len(), } } @@ -89,6 +100,14 @@ impl JsonVectorBuilderState { merged_type, values, } => build_legacy(values, merged_type), + Self::ExplicitOnly { + explicit, + remainder, + .. + } => { + let remainder = std::mem::replace(remainder, VariantArrayBuilder::new(0)).build(); + finish_vector(explicit.to_vector(), ArrayRef::from(remainder)) + } Self::AutoExpanding { explicit_type, max_auto_expanded_paths, @@ -101,6 +120,38 @@ impl JsonVectorBuilderState { } } + fn try_build_cloned(&self) -> Result { + let mut state = match self { + Self::Legacy { + merged_type, + values, + } => Self::Legacy { + merged_type: merged_type.clone(), + values: values.clone(), + }, + Self::AutoExpanding { + explicit_type, + max_auto_expanded_paths, + values, + } => Self::AutoExpanding { + explicit_type: explicit_type.clone(), + max_auto_expanded_paths: *max_auto_expanded_paths, + values: values.clone(), + }, + // Only TimeSeriesMemtable requires a non-consuming snapshot, while JSON2 targets + // BulkMemtable. We've tried our best to support it above, but if this match arm does + // not, it's OK. The only reason it doesn't is because of `VariantArrayBuilder`. We'll + // track the upstream and see. + Self::ExplicitOnly { .. } => { + return UnimplementedSnafu { + feat: "no auto expanded JSON2 array builder", + } + .fail(); + } + }; + state.try_build() + } + fn try_push_value_ref(&mut self, value: &ValueRef) -> Result<()> { let ValueRef::Json(value) = value else { return TryFromValueSnafu { @@ -125,6 +176,22 @@ impl JsonVectorBuilderState { } values.push(JsonVariant::from(value.variant())); } + Self::ExplicitOnly { + explicit_type, + struct_type, + explicit, + remainder, + } => { + if value.is_null() { + explicit.push_null(); + remainder.append_null(); + } else { + let (value, rest) = + split_to_explicit_ref(value.variant(), explicit_type, struct_type)?; + explicit.push_struct_value_ref(value)?; + append_json_variant_ref(remainder, &rest).context(ArrowComputeSnafu)?; + } + } Self::AutoExpanding { values, .. } => { values.push(JsonVariant::from(value.variant())); } @@ -137,6 +204,14 @@ impl JsonVectorBuilderState { Self::Legacy { values, .. } | Self::AutoExpanding { values, .. } => { values.push(JsonVariant::Null) } + Self::ExplicitOnly { + explicit, + remainder, + .. + } => { + explicit.push_null(); + remainder.append_null(); + } } } } @@ -163,13 +238,28 @@ impl JsonVectorBuilder { for hint in settings.type_hints() { insert_dynamic_type(&mut explicit_type, &hint.path, (&hint.data_type).into()); } - Self { - state: JsonVectorBuilderState::AutoExpanding { + let state = if settings.max_auto_expanded_paths() == Some(0) { + let DataType::Struct(fields) = explicit_type.as_arrow_type() else { + unreachable!("JSON2 explicit type must map to Arrow Struct") + }; + let struct_type = StructType::from(&fields); + JsonVectorBuilderState::ExplicitOnly { + explicit_type, + explicit: StructVectorBuilder::with_type_and_capacity( + struct_type.clone(), + capacity, + ), + struct_type, + remainder: VariantArrayBuilder::new(capacity), + } + } else { + JsonVectorBuilderState::AutoExpanding { explicit_type, max_auto_expanded_paths: settings.max_auto_expanded_paths().unwrap_or(u32::MAX), values: Vec::with_capacity(capacity), - }, - } + } + }; + Self { state } } fn try_build(&mut self) -> Result { @@ -451,6 +541,132 @@ fn insert_dynamic_type>( ) } +fn split_to_explicit_ref<'a>( + value: &JsonVariantRef<'a>, + explicit_type: &JsonNativeType, + struct_type: &StructType, +) -> Result<(StructValueRef<'a>, JsonVariantRef<'a>)> { + let JsonVariantRef::Object(object) = value else { + return TryFromValueSnafu { + reason: "expected json object value".to_string(), + } + .fail(); + }; + let explicit = json_object_ref_into_struct_value_ref(object, struct_type)?; + let remainder = remainder_ref(object, explicit_type)?; + Ok((explicit, JsonVariantRef::Object(remainder))) +} + +fn json_object_ref_into_struct_value_ref<'a>( + object: &BTreeMap<&'a str, JsonVariantRef<'a>>, + struct_type: &StructType, +) -> Result> { + let mut values = Vec::with_capacity(struct_type.fields().len()); + for field in struct_type.fields().iter() { + let value = match object.get(field.name()) { + Some(value) => json_variant_ref_into_value_ref(value, field.data_type())?, + None => ValueRef::Null, + }; + values.push(value); + } + Ok(StructValueRef::RefList { + val: values, + fields: struct_type.clone(), + }) +} + +fn json_variant_ref_into_value_ref<'a>( + value: &JsonVariantRef<'a>, + expected_type: &ConcreteDataType, +) -> Result> { + let value = match (value, expected_type) { + (JsonVariantRef::Null, _) | (_, ConcreteDataType::Null(_)) => ValueRef::Null, + (JsonVariantRef::Object(object), ConcreteDataType::Struct(struct_type)) => { + ValueRef::Struct(json_object_ref_into_struct_value_ref(object, struct_type)?) + } + (JsonVariantRef::Bool(x), ConcreteDataType::Boolean(_)) => ValueRef::Boolean(*x), + (JsonVariantRef::Number(x), ConcreteDataType::UInt64(_)) => { + let Some(x) = x.as_u64() else { + return TryFromValueSnafu { + reason: format!("unable to convert {x:?} to UInt64"), + } + .fail(); + }; + ValueRef::UInt64(x) + } + (JsonVariantRef::Number(x), ConcreteDataType::Int64(_)) => { + let x = match x { + JsonNumber::PosInt(x) => i64::try_from(*x).ok(), + JsonNumber::NegInt(x) => Some(*x), + JsonNumber::Float(_) => None, + }; + let Some(x) = x else { + return TryFromValueSnafu { + reason: format!("unable to convert {x:?} to Int64"), + } + .fail(); + }; + ValueRef::Int64(x) + } + (JsonVariantRef::Number(JsonNumber::PosInt(x)), ConcreteDataType::Float64(_)) => { + ValueRef::Float64((*x as f64).into()) + } + (JsonVariantRef::Number(JsonNumber::NegInt(x)), ConcreteDataType::Float64(_)) => { + ValueRef::Float64((*x as f64).into()) + } + (JsonVariantRef::Number(JsonNumber::Float(x)), ConcreteDataType::Float64(_)) => { + ValueRef::Float64(*x) + } + (JsonVariantRef::String(x), ConcreteDataType::String(_)) => ValueRef::String(x), + (value, expected_type) => { + return TryFromValueSnafu { + reason: format!("unable to convert json value {value:?} to {expected_type}"), + } + .fail(); + } + }; + Ok(value) +} + +fn remainder_ref<'a>( + object: &BTreeMap<&'a str, JsonVariantRef<'a>>, + explicit_type: &JsonNativeType, +) -> Result>> { + let JsonNativeType::Object(fields) = explicit_type else { + return UnexpectedSnafu { + reason: "JSON2 explicit type must be an object", + } + .fail(); + }; + let mut remainder = BTreeMap::new(); + for (&name, value) in object { + match fields.get(name) { + Some(data_type @ JsonNativeType::Object(_)) => match value { + JsonVariantRef::Null => {} + JsonVariantRef::Object(object) => { + let child = remainder_ref(object, data_type)?; + if !child.is_empty() { + remainder.insert(name, JsonVariantRef::Object(child)); + } + } + _ => { + return TryFromValueSnafu { + reason: "expected json object value".to_string(), + } + .fail(); + } + }, + // A non-object entry in the explicit type tree is an explicit leaf and is + // already written to the Struct builder. So here does nothing. + Some(_) => {} + None => { + remainder.insert(name, value.clone()); + } + } + } + Ok(remainder) +} + fn split_to_explicit( value: JsonVariant, explicit_type: &JsonNativeType, @@ -611,11 +827,19 @@ impl MutableVector for JsonVectorBuilder { } fn to_vector(&mut self) -> VectorRef { - self.try_build().unwrap_or_else(|e| panic!("{:?}", e)) + self.try_build().unwrap_or_else(|e| { + // Just try to avoid panicking here. + common_telemetry::error!(e; "Unable to build JSON2 vector"); + Arc::new(NullVector::new(self.len())) + }) } fn to_vector_cloned(&self) -> VectorRef { - self.clone().to_vector() + self.state.try_build_cloned().unwrap_or_else(|e| { + // Just try to avoid panicking here. + common_telemetry::error!(e; "Unable to build JSON2 vector"); + Arc::new(NullVector::new(self.len())) + }) } fn try_push_value_ref(&mut self, value: &ValueRef) -> Result<()> { @@ -772,7 +996,7 @@ mod tests { } #[test] - fn test_zero_budget_builder_uses_fixed_schema_and_remainder() -> Result<()> { + fn test_zero_budget_builder_uses_explicit_only_schema_and_remainder() -> Result<()> { let settings = JsonSettings::try_new( vec![ JsonTypeHint { @@ -800,6 +1024,10 @@ mod tests { Some(0), )?; let mut builder = JsonVectorBuilder::with_settings(&settings, 2); + assert!(matches!( + &builder.state, + JsonVectorBuilderState::ExplicitOnly { .. } + )); let values = [ json!({ "kind": "record", diff --git a/src/datatypes/src/vectors/json/variant.rs b/src/datatypes/src/vectors/json/variant.rs index 75d87c9986..9f90f0aaac 100644 --- a/src/datatypes/src/vectors/json/variant.rs +++ b/src/datatypes/src/vectors/json/variant.rs @@ -24,7 +24,7 @@ use parquet_variant_json::VariantToJson; use snafu::ResultExt; use crate::error::{ArrowComputeSnafu, Result}; -use crate::json::value::{JsonNumber, JsonVariant, decode_json_variant}; +use crate::json::value::{JsonNumber, JsonVariant, JsonVariantRef, decode_json_variant}; /// Returns the canonical Arrow field for an unshredded Parquet Variant array. pub fn variant_field(name: impl Into, nullable: bool) -> Field { @@ -110,6 +110,52 @@ pub(super) fn append_json_variant( Ok(()) } +pub(super) fn append_json_variant_ref( + builder: &mut impl VariantBuilderExt, + value: &JsonVariantRef<'_>, +) -> std::result::Result<(), ArrowError> { + match value { + JsonVariantRef::Null => builder.append_value(Variant::Null), + JsonVariantRef::Bool(value) => builder.append_value(*value), + JsonVariantRef::Number(JsonNumber::PosInt(value)) => { + if let Ok(value) = i64::try_from(*value) { + builder.append_value(value); + } else { + append_large_u64(builder, *value)?; + } + } + JsonVariantRef::Number(JsonNumber::NegInt(value)) => builder.append_value(*value), + JsonVariantRef::Number(JsonNumber::Float(value)) => { + if value.0.is_finite() { + builder.append_value(value.0) + } else { + builder.append_value("NaN") + } + } + JsonVariantRef::String(value) => builder.append_value(*value), + JsonVariantRef::Array(values) => { + let mut list = builder.try_new_list()?; + for value in values { + append_json_variant_ref(&mut list, value)?; + } + list.finish(); + } + JsonVariantRef::Object(values) => { + let mut object = builder.try_new_object()?; + for (name, value) in values { + append_json_variant_ref(&mut ObjectFieldBuilder::new(name, &mut object), value)?; + } + object.finish(); + } + JsonVariantRef::Variant(value) => { + let value = decode_json_variant(value) + .map_err(|e| ArrowError::JsonError(format!("Failed to decode JSONB: {e}")))?; + append_json_value(builder, &value)?; + } + } + Ok(()) +} + fn append_json_value( builder: &mut impl VariantBuilderExt, value: &serde_json::Value,