1use std::any::Any;
16use std::sync::Arc;
17
18use arrow_schema::DataType;
19
20use crate::data_type::ConcreteDataType;
21use crate::error::{Error, Result, TryFromValueSnafu, UnexpectedSnafu, UnsupportedOperationSnafu};
22use crate::json::value::{JsonNumber, JsonVariant, encode_json_variant};
23use crate::prelude::{ValueRef, Vector, VectorRef};
24use crate::types::StructType;
25use crate::types::json_type::{JsonNativeType, is_include};
26use crate::value::{ListValue, StructValue, StructValueRef, Value};
27use crate::vectors::{MutableVector, StructVectorBuilder};
28
29#[derive(Clone)]
30pub(crate) struct JsonVectorBuilder {
31 merged_type: JsonNativeType,
32 values: Vec<JsonVariant>,
33}
34
35impl JsonVectorBuilder {
36 pub(crate) fn new(initial_native_type: JsonNativeType, capacity: usize) -> Self {
37 debug_assert!(matches!(
38 initial_native_type,
39 JsonNativeType::Object(_) | JsonNativeType::Null
40 ));
41 Self {
42 merged_type: initial_native_type,
43 values: Vec::with_capacity(capacity),
44 }
45 }
46
47 fn try_build(&mut self) -> Result<VectorRef> {
48 let DataType::Struct(fields) = self.merged_type.as_arrow_type() else {
49 return UnexpectedSnafu {
50 reason: "merged JSON2 type must map to Arrow Struct in JsonVectorBuilder",
51 }
52 .fail();
53 };
54 let struct_type = StructType::from(&fields);
56
57 let mut builder =
58 StructVectorBuilder::with_type_and_capacity(struct_type.clone(), self.values.len());
59 for value in std::mem::take(&mut self.values) {
60 if matches!(&value, JsonVariant::Null) {
61 builder.push_null();
62 continue;
63 }
64 let value = json_variant_into_struct_value(value, struct_type.clone())?;
65 builder.push_struct_value_ref(StructValueRef::Ref(&value))?;
66 }
67 Ok(builder.to_vector())
68 }
69}
70
71fn json_variant_into_struct_value(
72 value: JsonVariant,
73 struct_type: StructType,
74) -> Result<StructValue> {
75 let JsonVariant::Object(object) = value else {
76 return TryFromValueSnafu {
77 reason: format!("expected json object value, got {value:?}"),
78 }
79 .fail();
80 };
81
82 let mut entries = object.into_iter();
83 let mut entry = entries.next();
84 let mut values = Vec::with_capacity(struct_type.fields().len());
85 for field in struct_type.fields().iter() {
86 let value = match entry.take() {
87 Some((name, value)) if name == field.name() => {
88 entry = entries.next();
89 json_variant_into_value(value, field.data_type())?
90 }
91 Some((name, _)) if name.as_str() < field.name() => {
92 return TryFromValueSnafu {
93 reason: format!("field {name} is missing from merged JSON type"),
94 }
95 .fail();
96 }
97 next => {
98 entry = next;
99 Value::Null
100 }
101 };
102 values.push(value);
103 }
104 if let Some((name, _)) = entry {
105 return TryFromValueSnafu {
106 reason: format!("field {name} is missing from merged JSON type"),
107 }
108 .fail();
109 }
110
111 Ok(StructValue::new(values, struct_type))
112}
113
114fn json_variant_into_value(value: JsonVariant, expected_type: &ConcreteDataType) -> Result<Value> {
115 let value = match (value, expected_type) {
116 (JsonVariant::Null, _) | (_, ConcreteDataType::Null(_)) => Value::Null,
117 (JsonVariant::Object(object), _) if object.is_empty() => Value::Null,
118 (JsonVariant::Bool(x), ConcreteDataType::Boolean(_)) => Value::Boolean(x),
119 (JsonVariant::Number(x), ConcreteDataType::UInt64(_)) => {
120 let Some(x) = x.as_u64() else {
121 return TryFromValueSnafu {
122 reason: format!("unable to convert {x:?} to UInt64"),
123 }
124 .fail();
125 };
126 Value::UInt64(x)
127 }
128 (JsonVariant::Number(x), ConcreteDataType::Int64(_)) => {
129 let x = match x {
130 JsonNumber::PosInt(x) => i64::try_from(x).ok(),
131 JsonNumber::NegInt(x) => Some(x),
132 JsonNumber::Float(_) => None,
133 };
134 let Some(x) = x else {
135 return TryFromValueSnafu {
136 reason: format!("unable to convert {x:?} to Int64"),
137 }
138 .fail();
139 };
140 Value::Int64(x)
141 }
142 (JsonVariant::Number(JsonNumber::PosInt(x)), ConcreteDataType::Float64(_)) => {
143 Value::Float64((x as f64).into())
144 }
145 (JsonVariant::Number(JsonNumber::NegInt(x)), ConcreteDataType::Float64(_)) => {
146 Value::Float64((x as f64).into())
147 }
148 (JsonVariant::Number(JsonNumber::Float(x)), ConcreteDataType::Float64(_)) => {
149 Value::Float64(x)
150 }
151 (JsonVariant::String(x), ConcreteDataType::String(_)) => Value::String(x.into()),
152 (JsonVariant::Array(array), ConcreteDataType::List(list_type)) => {
153 let item_type = list_type.item_type().clone();
154 let values = array
155 .into_iter()
156 .map(|v| json_variant_into_value(v, &item_type))
157 .collect::<Result<Vec<_>>>()?;
158 Value::List(ListValue::new(values, Arc::new(item_type)))
159 }
160 (value @ JsonVariant::Object(_), ConcreteDataType::Struct(struct_type)) => {
161 Value::Struct(json_variant_into_struct_value(value, struct_type.clone())?)
162 }
163 (value, ConcreteDataType::Binary(_)) => Value::from(encode_json_variant(value)?),
164 (value, expected_type) => {
165 return TryFromValueSnafu {
166 reason: format!("unable to convert json value {value:?} to {expected_type}"),
167 }
168 .fail();
169 }
170 };
171 Ok(value)
172}
173
174pub(crate) fn json_variant_into_projected_value(
178 value: JsonVariant,
179 expected_type: &ConcreteDataType,
180) -> Result<Value> {
181 match (value, expected_type) {
182 (JsonVariant::Null, _) => Ok(Value::Null),
183 (JsonVariant::String(value), ConcreteDataType::String(_)) => {
184 Ok(Value::String(value.into()))
185 }
186 (value, ConcreteDataType::String(_)) => Ok(Value::String(value.to_string().into())),
187 (JsonVariant::Object(mut object), ConcreteDataType::Struct(struct_type)) => {
188 let values = struct_type
189 .fields()
190 .iter()
191 .map(|field| {
192 object
193 .remove(field.name())
194 .map(|value| {
195 null_on_json_type_mismatch(json_variant_into_projected_value(
200 value,
201 field.data_type(),
202 ))
203 })
204 .transpose()
205 .map(|value| value.unwrap_or(Value::Null))
206 })
207 .collect::<Result<Vec<_>>>()?;
208 Ok(Value::Struct(StructValue::new(values, struct_type.clone())))
209 }
210 (JsonVariant::Array(array), ConcreteDataType::List(list_type)) => {
211 let item_type = list_type.item_type().clone();
212 let values = array
213 .into_iter()
214 .map(|value| {
215 null_on_json_type_mismatch(json_variant_into_projected_value(value, &item_type))
217 })
218 .collect::<Result<Vec<_>>>()?;
219 Ok(Value::List(ListValue::new(values, Arc::new(item_type))))
220 }
221 (value, expected_type) => json_variant_into_value(value, expected_type),
222 }
223}
224
225pub(crate) fn null_on_json_type_mismatch(result: Result<Value>) -> Result<Value> {
226 match result {
227 Ok(value) => Ok(value),
228 Err(Error::TryFromValue { .. }) => Ok(Value::Null),
229 Err(error) => Err(error),
230 }
231}
232
233impl MutableVector for JsonVectorBuilder {
234 fn data_type(&self) -> ConcreteDataType {
235 ConcreteDataType::json2(self.merged_type.clone())
236 }
237
238 fn len(&self) -> usize {
239 self.values.len()
240 }
241
242 fn as_any(&self) -> &dyn Any {
243 self
244 }
245
246 fn as_mut_any(&mut self) -> &mut dyn Any {
247 self
248 }
249
250 fn to_vector(&mut self) -> VectorRef {
251 self.try_build().unwrap_or_else(|e| panic!("{:?}", e))
252 }
253
254 fn to_vector_cloned(&self) -> VectorRef {
255 self.clone().to_vector()
256 }
257
258 fn try_push_value_ref(&mut self, value: &ValueRef) -> Result<()> {
259 let ValueRef::Json(value) = value else {
260 return TryFromValueSnafu {
261 reason: format!("expected json value, got {value:?}"),
262 }
263 .fail();
264 };
265 let json_type = value.json_type();
266 let json_type = json_type.as_ref();
267 if !matches!(json_type, JsonNativeType::Object(_) | JsonNativeType::Null) {
268 return TryFromValueSnafu {
269 reason: format!("expected json object value, got {value:?}"),
270 }
271 .fail();
272 }
273 if !is_include(&self.merged_type, json_type) {
274 self.merged_type.merge(json_type);
275 }
276
277 self.values.push(JsonVariant::from(value.variant()));
278 Ok(())
279 }
280
281 fn push_null(&mut self) {
282 self.values.push(JsonVariant::Null)
283 }
284
285 fn extend_slice_of(&mut self, _: &dyn Vector, _: usize, _: usize) -> Result<()> {
286 UnsupportedOperationSnafu {
287 op: "extend_slice_of",
288 vector_type: "JsonVector",
289 }
290 .fail()
291 }
292}
293
294#[cfg(test)]
295mod tests {
296 use std::sync::Arc;
297
298 use common_base::bytes::Bytes;
299
300 use super::*;
301 use crate::data_type::ConcreteDataType;
302 use crate::types::StructField;
303 use crate::types::json_type::JsonObjectType;
304 use crate::value::{ListValue, StructValue, Value, ValueRef};
305
306 #[test]
307 fn test_json_vector_builder() -> Result<()> {
308 fn parse_json_value(json: &str) -> Value {
309 let value: serde_json::Value = serde_json::from_str(json).unwrap();
310 Value::Json(Box::new(value.into()))
311 }
312
313 fn jsonb_bytes(json: &str) -> Bytes {
314 Bytes::from(jsonb::parse_value(json.as_bytes()).unwrap().to_vec())
315 }
316
317 let mut builder = JsonVectorBuilder::new(JsonNativeType::Object(Default::default()), 3);
320 let first = parse_json_value(r#"{"id":1,"payload":{"name":"foo"}}"#);
321 let second = parse_json_value(r#"{"id":2,"extra":true,"payload":"raw"}"#);
322 builder.try_push_value_ref(&first.as_value_ref())?;
323 builder.push_null();
324 builder.try_push_value_ref(&second.as_value_ref())?;
325
326 let merged_type = JsonNativeType::Object(JsonObjectType::from([
327 ("extra".to_string(), JsonNativeType::Bool),
328 ("id".to_string(), JsonNativeType::i64()),
329 ("payload".to_string(), JsonNativeType::Variant),
330 ]));
331 assert_eq!(
332 builder.data_type(),
333 ConcreteDataType::json2(merged_type.clone())
334 );
335
336 let DataType::Struct(fields) = merged_type.as_arrow_type() else {
337 unreachable!()
338 };
339 let merged_struct_type = StructType::from(&fields);
340 let vector = builder.to_vector();
341 assert_eq!(vector.len(), 3);
342 assert_eq!(
343 vector.get(0),
344 Value::Struct(StructValue::new(
345 vec![
346 Value::Null,
347 Value::Int64(1),
348 Value::Binary(jsonb_bytes(r#"{"name":"foo"}"#)),
349 ],
350 merged_struct_type.clone(),
351 ))
352 );
353 assert_eq!(vector.get(1), Value::Null);
354 assert_eq!(
355 vector.get(2),
356 Value::Struct(StructValue::new(
357 vec![
358 Value::Boolean(true),
359 Value::Int64(2),
360 Value::Binary(jsonb_bytes(r#""raw""#)),
361 ],
362 merged_struct_type,
363 ))
364 );
365
366 let mut inferred_builder = JsonVectorBuilder::new(JsonNativeType::Null, 2);
369 let inferred_value = parse_json_value(r#"{"id":3}"#);
370 inferred_builder.push_null();
371 inferred_builder.try_push_value_ref(&inferred_value.as_value_ref())?;
372
373 let inferred_type = JsonNativeType::Object(JsonObjectType::from([(
374 "id".to_string(),
375 JsonNativeType::i64(),
376 )]));
377 assert_eq!(
378 inferred_builder.data_type(),
379 ConcreteDataType::json2(inferred_type.clone())
380 );
381
382 let DataType::Struct(fields) = inferred_type.as_arrow_type() else {
383 unreachable!()
384 };
385 let inferred_struct_type = StructType::from(&fields);
386 let vector = inferred_builder.to_vector();
387 assert_eq!(vector.get(0), Value::Null);
388 assert_eq!(
389 vector.get(1),
390 Value::Struct(StructValue::new(
391 vec![Value::Int64(3)],
392 inferred_struct_type,
393 ))
394 );
395
396 let result = std::panic::catch_unwind(|| JsonVectorBuilder::new(JsonNativeType::Bool, 2));
398 assert!(result.is_err());
399
400 let mut object_builder =
402 JsonVectorBuilder::new(JsonNativeType::Object(Default::default()), 2);
403 let object = parse_json_value(r#"{"k":1}"#);
404 let boolean = parse_json_value("true");
405 let err = object_builder
406 .try_push_value_ref(&boolean.as_value_ref())
407 .unwrap_err();
408 assert!(err.to_string().contains("expected json object value"));
409 object_builder.try_push_value_ref(&object.as_value_ref())?;
410
411 let mut invalid_builder =
413 JsonVectorBuilder::new(JsonNativeType::Object(Default::default()), 1);
414 let err = invalid_builder
415 .try_push_value_ref(&ValueRef::Boolean(true))
416 .unwrap_err();
417 assert!(err.to_string().contains("expected json value"));
418
419 Ok(())
420 }
421
422 #[test]
423 fn test_json_variant_into_struct_value() -> Result<()> {
424 assert_eq!(
425 json_variant_into_value(
426 JsonVariant::Object(Default::default()),
427 &ConcreteDataType::string_datatype(),
428 )?,
429 Value::Null
430 );
431
432 let item_type =
433 ConcreteDataType::struct_datatype(StructType::new(Arc::new(vec![StructField::new(
434 "id".to_string(),
435 ConcreteDataType::int64_datatype(),
436 true,
437 )])));
438 let struct_type = StructType::new(Arc::new(vec![
439 StructField::new(
440 "items".to_string(),
441 ConcreteDataType::list_datatype(Arc::new(item_type.clone())),
442 true,
443 ),
444 StructField::new(
445 "meta".to_string(),
446 ConcreteDataType::struct_datatype(StructType::new(Arc::new(vec![
447 StructField::new(
448 "name".to_string(),
449 ConcreteDataType::string_datatype(),
450 true,
451 ),
452 ]))),
453 true,
454 ),
455 ]));
456 let variant = JsonVariant::from([
457 (
458 "items",
459 JsonVariant::Array(vec![
460 JsonVariant::from([("id", JsonVariant::from(1i64))]),
461 JsonVariant::from([("id", JsonVariant::from(2i64))]),
462 ]),
463 ),
464 (
465 "meta",
466 JsonVariant::from([("name", JsonVariant::from("foo"))]),
467 ),
468 ]);
469 let value = Value::Struct(json_variant_into_struct_value(
470 variant,
471 struct_type.clone(),
472 )?);
473
474 assert_eq!(
475 value,
476 Value::Struct(StructValue::new(
477 vec![
478 Value::List(ListValue::new(
479 vec![
480 Value::Struct(StructValue::new(
481 vec![Value::Int64(1)],
482 StructType::new(Arc::new(vec![StructField::new(
483 "id".to_string(),
484 ConcreteDataType::int64_datatype(),
485 true,
486 )]))
487 )),
488 Value::Struct(StructValue::new(
489 vec![Value::Int64(2)],
490 StructType::new(Arc::new(vec![StructField::new(
491 "id".to_string(),
492 ConcreteDataType::int64_datatype(),
493 true,
494 )]))
495 )),
496 ],
497 Arc::new(item_type),
498 )),
499 Value::Struct(StructValue::new(
500 vec![Value::String("foo".into())],
501 StructType::new(Arc::new(vec![StructField::new(
502 "name".to_string(),
503 ConcreteDataType::string_datatype(),
504 true,
505 )])),
506 )),
507 ],
508 struct_type,
509 ))
510 );
511 Ok(())
512 }
513
514 #[test]
515 fn test_projected_value_nulls_nested_type_mismatches() {
516 let struct_type = ConcreteDataType::struct_datatype(StructType::new(Arc::new(vec![
517 StructField::new(
518 "value".to_string(),
519 ConcreteDataType::uint64_datatype(),
520 true,
521 ),
522 StructField::new(
523 "items".to_string(),
524 ConcreteDataType::list_datatype(Arc::new(ConcreteDataType::uint64_datatype())),
525 true,
526 ),
527 StructField::new(
528 "payload".to_string(),
529 ConcreteDataType::binary_datatype(),
530 true,
531 ),
532 ])));
533
534 let invalid_field = JsonVariant::from([("value", JsonVariant::from("invalid"))]);
535 let Value::Struct(value) = json_variant_into_projected_value(invalid_field, &struct_type)
536 .expect("type mismatch should become null")
537 else {
538 unreachable!()
539 };
540 assert_eq!(value.items()[0], Value::Null);
541
542 let invalid_item = JsonVariant::from([(
543 "items",
544 JsonVariant::Array(vec![JsonVariant::from("invalid")]),
545 )]);
546 let Value::Struct(value) = json_variant_into_projected_value(invalid_item, &struct_type)
547 .expect("type mismatch should become null")
548 else {
549 unreachable!()
550 };
551 let Value::List(items) = &value.items()[1] else {
552 unreachable!()
553 };
554 assert_eq!(items.items()[0], Value::Null);
555
556 let invalid_json = JsonVariant::from([("payload", JsonVariant::from(f64::NAN))]);
557 let error = json_variant_into_projected_value(invalid_json, &struct_type).unwrap_err();
558 assert!(matches!(error, Error::InvalidJson { .. }));
559 }
560}