1use std::cmp::Ordering;
16use std::sync::Arc;
17
18use arrow::compute;
19use arrow::util::display::{ArrayFormatter, FormatOptions};
20use arrow_array::cast::AsArray;
21use arrow_array::types::{Float64Type, Int64Type, UInt64Type};
22use arrow_array::{Array, ArrayRef, GenericListArray, ListArray, StructArray, new_null_array};
23use arrow_schema::{DataType, FieldRef};
24use serde_json::Value;
25use snafu::{OptionExt, ResultExt};
26
27use crate::arrow_array::{
28 MutableBinaryArray, StringViewArray, binary_array_value, string_array_value,
29};
30use crate::data_type::ConcreteDataType;
31use crate::error::{
32 AlignJsonArraySnafu, ArrowComputeSnafu, InvalidJsonSnafu, InvalidJsonbSnafu, Result,
33};
34use crate::json::value::{JsonVariant, decode_json_variant, encode_serde_json_as_jsonb};
35use crate::prelude::DataType as _;
36use crate::vectors::json::builder::{
37 json_variant_into_projected_value, null_on_json_type_mismatch,
38};
39
40pub struct JsonArray<'a> {
41 inner: &'a ArrayRef,
42}
43
44impl JsonArray<'_> {
45 pub fn try_get_value(&self, i: usize) -> Result<Value> {
47 let array = self.inner;
48 if array.is_null(i) {
49 return Ok(Value::Null);
50 }
51
52 let value = match array.data_type() {
53 DataType::Null => Value::Null,
54 DataType::Boolean => Value::Bool(array.as_boolean().value(i)),
55 DataType::Int64 => Value::from(array.as_primitive::<Int64Type>().value(i)),
56 DataType::UInt64 => Value::from(array.as_primitive::<UInt64Type>().value(i)),
57 DataType::Float64 => Value::from(array.as_primitive::<Float64Type>().value(i)),
58 DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View => {
59 Value::String(string_array_value(array, i).to_string())
60 }
61 DataType::Binary | DataType::LargeBinary | DataType::BinaryView => {
62 let bytes = binary_array_value(array, i);
63 decode_json_variant(bytes).map_err(|error| InvalidJsonbSnafu { error }.build())?
64 }
65 DataType::Struct(_) => {
66 let structs = array.as_struct();
67 let object = structs
68 .fields()
69 .iter()
70 .zip(structs.columns())
71 .map(|(field, column)| {
72 JsonArray::from(column)
73 .try_get_value(i)
74 .map(|v| (field.name().clone(), v))
75 })
76 .collect::<Result<_>>()?;
77 Value::Object(object)
78 }
79 DataType::List(_) => {
80 let lists = array.as_list::<i32>();
81 let list = lists.value(i);
82 let list = JsonArray::from(&list);
83 let mut values = Vec::with_capacity(list.inner.len());
84 for i in 0..list.inner.len() {
85 values.push(list.try_get_value(i)?);
86 }
87 Value::Array(values)
88 }
89 t => {
90 return InvalidJsonSnafu {
91 value: format!("unknown JSON type {t}"),
92 }
93 .fail();
94 }
95 };
96 Ok(value)
97 }
98
99 pub fn try_align(&self, expect: &DataType) -> Result<ArrayRef> {
106 if self.inner.data_type() == expect {
107 return Ok(self.inner.clone());
108 }
109
110 common_telemetry::trace!(
111 "Try aligning JSON array {} to data type {}",
112 self.inner.data_type(),
113 expect
114 );
115
116 if self.inner.data_type().is_binary() && matches!(expect, DataType::Struct(_)) {
117 return self.decode_variant(expect);
118 }
119
120 let struct_array = self.inner.as_struct_opt().context(AlignJsonArraySnafu {
121 reason: "expect struct array",
122 })?;
123 let array_fields = struct_array.fields();
124 let array_columns = struct_array.columns();
125 let DataType::Struct(expect_fields) = expect else {
126 return AlignJsonArraySnafu {
127 reason: "expect struct datatype",
128 }
129 .fail();
130 };
131 let mut aligned = Vec::with_capacity(expect_fields.len());
132
133 debug_assert!(expect_fields.iter().map(|f| f.name()).is_sorted());
138 debug_assert!(array_fields.iter().map(|f| f.name()).is_sorted());
139
140 let mut i = 0; let mut j = 0; while i < expect_fields.len() && j < array_fields.len() {
143 let expect_field = &expect_fields[i];
144 let array_field = &array_fields[j];
145 match expect_field.name().cmp(array_field.name()) {
146 Ordering::Equal => {
147 if expect_field.data_type() == array_field.data_type() {
148 aligned.push(array_columns[j].clone());
149 } else {
150 let expect_type = expect_field.data_type();
151 let array_type = array_field.data_type();
152 let array = match (expect_type, array_type) {
153 (DataType::Struct(_), DataType::Struct(_)) => {
154 JsonArray::from(&array_columns[j]).try_align(expect_type)?
155 }
156 (DataType::List(expect_item), DataType::List(array_item)) => {
157 let list_array = array_columns[j].as_list::<i32>();
158 try_align_list(list_array, expect_item, array_item)?
159 }
160 _ => JsonArray::from(&array_columns[j]).try_cast(expect_type)?,
161 };
162 aligned.push(array);
163 }
164 i += 1;
165 j += 1;
166 }
167 Ordering::Less => {
168 aligned.push(new_null_array(expect_field.data_type(), struct_array.len()));
169 i += 1;
170 }
171 Ordering::Greater => {
172 j += 1;
173 }
174 }
175 }
176 if i < expect_fields.len() {
177 for field in &expect_fields[i..] {
178 aligned.push(new_null_array(field.data_type(), struct_array.len()));
179 }
180 }
181
182 let json_array = StructArray::try_new(
183 expect_fields.clone(),
184 aligned,
185 struct_array.nulls().cloned(),
186 )
187 .map_err(|e| {
188 AlignJsonArraySnafu {
189 reason: e.to_string(),
190 }
191 .build()
192 })?;
193 Ok(Arc::new(json_array))
194 }
195
196 fn try_cast(&self, to_type: &DataType) -> Result<ArrayRef> {
197 let from_type = self.inner.data_type();
198 if from_type == to_type {
199 return Ok(self.inner.clone());
200 }
201
202 if to_type == &DataType::Utf8View {
203 let values = (0..self.inner.len())
204 .map(|i| {
205 if self.inner.is_null(i) {
206 return Ok(None);
207 }
208 let value = match self.try_get_value(i)? {
209 Value::Null => return Ok(None),
210 Value::String(value) => value,
211 value => value.to_string(),
212 };
213 Ok(Some(value))
214 })
215 .collect::<Result<Vec<_>>>()?;
216 return Ok(Arc::new(StringViewArray::from(values)) as ArrayRef);
217 }
218
219 if from_type.is_binary() && !to_type.is_binary() {
220 return self.decode_variant(to_type);
221 }
222
223 if !from_type.is_binary() && to_type.is_binary() {
224 return self.encode_variant();
225 }
226
227 if compute::can_cast_types(from_type, to_type) {
228 return compute::cast(&self.inner, to_type).context(ArrowComputeSnafu);
229 }
230
231 let formatter = ArrayFormatter::try_new(&self.inner, &FormatOptions::default())
232 .context(ArrowComputeSnafu)?;
233
234 let values = (0..self.inner.len())
235 .map(|i| {
236 self.inner
237 .is_valid(i)
238 .then(|| formatter.value(i).to_string())
239 })
240 .collect::<Vec<_>>();
241 Ok(Arc::new(StringViewArray::from(values)))
242 }
243
244 fn encode_variant(&self) -> Result<ArrayRef> {
245 let len = self.inner.len();
246 let mut encoded = Vec::with_capacity(len);
247 let mut total_bytes = 0;
248
249 for i in 0..len {
250 let value = self.try_get_value(i)?;
251 if value.is_null() {
252 encoded.push(None);
253 } else {
254 let bytes = encode_serde_json_as_jsonb(value);
255 total_bytes += bytes.len();
256 encoded.push(Some(bytes));
257 }
258 }
259
260 let mut builder = MutableBinaryArray::with_capacity(len, total_bytes);
261 for value in encoded {
262 builder.append_option(value);
263 }
264 Ok(Arc::new(builder.finish()))
265 }
266
267 fn decode_variant(&self, to_type: &DataType) -> Result<ArrayRef> {
268 let values = (0..self.inner.len())
269 .map(|i| self.try_get_value(i))
270 .collect::<Result<Vec<_>>>()?;
271 decode_json_values(values, to_type)
272 }
273}
274
275fn decode_json_values(values: Vec<Value>, to_type: &DataType) -> Result<ArrayRef> {
276 let concrete_type = ConcreteDataType::from_arrow_type(to_type);
277 let mut builder = concrete_type.create_mutable_vector(values.len());
278 for value in values {
279 let value = null_on_json_type_mismatch(json_variant_into_projected_value(
280 JsonVariant::from(value),
281 &concrete_type,
282 ))?;
283 builder.try_push_value_ref(&value.as_value_ref())?;
284 }
285 Ok(builder.to_vector().to_arrow_array())
286}
287
288fn try_align_list(
289 list_array: &ListArray,
290 expect_item: &FieldRef,
291 array_item: &FieldRef,
292) -> Result<ArrayRef> {
293 let item_aligned = match (expect_item.data_type(), array_item.data_type()) {
294 (DataType::Struct(_), DataType::Struct(_)) => {
295 JsonArray::from(list_array.values()).try_align(expect_item.data_type())?
296 }
297 (DataType::List(expect_item), DataType::List(array_item)) => {
298 let list_array = list_array.values().as_list::<i32>();
299 try_align_list(list_array, expect_item, array_item)?
300 }
301 _ => JsonArray::from(list_array.values()).try_cast(expect_item.data_type())?,
302 };
303 Ok(Arc::new(
304 GenericListArray::<i32>::try_new(
305 expect_item.clone(),
306 list_array.offsets().clone(),
307 item_aligned,
308 list_array.nulls().cloned(),
309 )
310 .context(ArrowComputeSnafu)?,
311 ))
312}
313
314impl<'a> From<&'a ArrayRef> for JsonArray<'a> {
315 fn from(inner: &'a ArrayRef) -> Self {
316 Self { inner }
317 }
318}
319
320#[cfg(test)]
321mod test {
322 use arrow_array::types::Int64Type;
323 use arrow_array::{
324 BinaryArray, BooleanArray, Float64Array, Int32Array, Int64Array, ListArray, StringArray,
325 };
326 use arrow_schema::{Field, Fields};
327 use serde_json::json;
328
329 use super::*;
330
331 #[test]
332 fn test_try_get_value() -> Result<()> {
333 let nulls = new_null_array(&DataType::Null, 2);
334 assert_eq!(JsonArray::from(&nulls).try_get_value(0)?, Value::Null);
335
336 let bools: ArrayRef = Arc::new(BooleanArray::from(vec![Some(true), None]));
337 assert_eq!(JsonArray::from(&bools).try_get_value(0)?, json!(true));
338 assert_eq!(JsonArray::from(&bools).try_get_value(1)?, Value::Null);
339
340 let ints: ArrayRef = Arc::new(Int64Array::from(vec![Some(-7), None]));
341 assert_eq!(JsonArray::from(&ints).try_get_value(0)?, json!(-7));
342 assert_eq!(JsonArray::from(&ints).try_get_value(1)?, Value::Null);
343
344 let floats: ArrayRef = Arc::new(Float64Array::from(vec![Some(1.5)]));
345 assert_eq!(JsonArray::from(&floats).try_get_value(0)?, json!(1.5));
346
347 let strings: ArrayRef = Arc::new(StringArray::from(vec![Some("hello"), None]));
348 assert_eq!(JsonArray::from(&strings).try_get_value(0)?, json!("hello"));
349 assert_eq!(JsonArray::from(&strings).try_get_value(1)?, Value::Null);
350
351 let nested = jsonb::parse_value(br#"{"nested":[1,null,"x"]}"#)
352 .unwrap()
353 .to_vec();
354 let null = jsonb::parse_value(b"null").unwrap().to_vec();
355 let binaries: ArrayRef =
356 Arc::new(BinaryArray::from(vec![nested.as_slice(), null.as_slice()]));
357 assert_eq!(
358 JsonArray::from(&binaries).try_get_value(0)?,
359 json!({"nested": [1, null, "x"]})
360 );
361 assert_eq!(JsonArray::from(&binaries).try_get_value(1)?, Value::Null);
362
363 let lists: ArrayRef = Arc::new(ListArray::from_iter_primitive::<Int64Type, _, _>(vec![
364 Some(vec![Some(1), None, Some(3)]),
365 None,
366 ]));
367 assert_eq!(
368 JsonArray::from(&lists).try_get_value(0)?,
369 json!([1, null, 3])
370 );
371 assert_eq!(JsonArray::from(&lists).try_get_value(1)?, Value::Null);
372
373 let structs: ArrayRef = Arc::new(StructArray::from(vec![
374 (
375 Arc::new(Field::new("flag", DataType::Boolean, true)),
376 Arc::new(BooleanArray::from(vec![Some(true), None])) as ArrayRef,
377 ),
378 (
379 Arc::new(Field::new_list(
380 "items",
381 Field::new_list_field(DataType::Int64, true),
382 true,
383 )),
384 Arc::new(ListArray::from_iter_primitive::<Int64Type, _, _>(vec![
385 Some(vec![Some(1), None]),
386 Some(vec![Some(2)]),
387 ])) as ArrayRef,
388 ),
389 ]));
390 assert_eq!(
391 JsonArray::from(&structs).try_get_value(0)?,
392 json!({"flag": true, "items": [1, null]})
393 );
394 assert_eq!(
395 JsonArray::from(&structs).try_get_value(1)?,
396 json!({"flag": null, "items": [2]})
397 );
398
399 let unsupported: ArrayRef = Arc::new(Int32Array::from(vec![1]));
400 assert_eq!(
401 JsonArray::from(&unsupported)
402 .try_get_value(0)
403 .unwrap_err()
404 .to_string(),
405 "Invalid JSON: unknown JSON type Int32"
406 );
407
408 Ok(())
409 }
410
411 #[test]
412 fn test_cast_variant_to_utf8_view_preserves_json_null() -> Result<()> {
413 let encode = |json: &[u8]| jsonb::parse_value(json).unwrap().to_vec();
414 let json_null = encode(b"null");
415 let object = encode(br#"{"value":1}"#);
416 let string = encode(br#""text""#);
417 let variants: ArrayRef = Arc::new(BinaryArray::from(vec![
418 Some(json_null.as_slice()),
419 Some(object.as_slice()),
420 Some(string.as_slice()),
421 None,
422 ]));
423
424 let casted = JsonArray::from(&variants).try_cast(&DataType::Utf8View)?;
425 let casted = casted.as_string_view();
426 assert!(casted.is_null(0));
427 assert_eq!(casted.value(1), r#"{"value":1}"#);
428 assert_eq!(casted.value(2), "text");
429 assert!(casted.is_null(3));
430
431 Ok(())
432 }
433
434 #[test]
435 fn test_align_json_array() -> Result<()> {
436 struct TestCase {
437 json_array: ArrayRef,
438 schema_type: DataType,
439 expected: std::result::Result<ArrayRef, String>,
440 }
441
442 impl TestCase {
443 fn new(
444 json_array: StructArray,
445 schema_type: Fields,
446 expected: std::result::Result<Vec<ArrayRef>, String>,
447 ) -> Self {
448 Self {
449 json_array: Arc::new(json_array),
450 schema_type: DataType::Struct(schema_type.clone()),
451 expected: expected
452 .map(|x| Arc::new(StructArray::new(schema_type, x, None)) as ArrayRef),
453 }
454 }
455
456 fn test(self) -> Result<()> {
457 let result = JsonArray::from(&self.json_array).try_align(&self.schema_type);
458 match (result, self.expected) {
459 (Ok(json_array), Ok(expected)) => assert_eq!(&json_array, &expected),
460 (Ok(json_array), Err(e)) => {
461 panic!("expecting error {e} but actually get: {json_array:?}")
462 }
463 (Err(e), Err(expected)) => assert_eq!(e.to_string(), expected),
464 (Err(e), Ok(_)) => return Err(e),
465 }
466 Ok(())
467 }
468 }
469
470 TestCase::new(
472 StructArray::new_empty_fields(2, None),
473 Fields::from(vec![
474 Field::new("int", DataType::Int64, true),
475 Field::new_struct(
476 "nested",
477 vec![Field::new("bool", DataType::Boolean, true)],
478 true,
479 ),
480 Field::new("string", DataType::Utf8, true),
481 ]),
482 Ok(vec![
483 Arc::new(Int64Array::new_null(2)) as ArrayRef,
484 Arc::new(StructArray::new_null(
485 Fields::from(vec![Arc::new(Field::new("bool", DataType::Boolean, true))]),
486 2,
487 )),
488 Arc::new(StringArray::new_null(2)),
489 ]),
490 )
491 .test()?;
492
493 TestCase::new(
495 StructArray::from(vec![(
496 Arc::new(Field::new("float", DataType::Float64, true)),
497 Arc::new(Float64Array::from(vec![1.0, 2.0, 3.0])) as ArrayRef,
498 )]),
499 Fields::from(vec![
500 Field::new("float", DataType::Float64, true),
501 Field::new("string", DataType::Utf8, true),
502 ]),
503 Ok(vec![
504 Arc::new(Float64Array::from(vec![1.0, 2.0, 3.0])) as ArrayRef,
505 Arc::new(StringArray::new_null(3)),
506 ]),
507 )
508 .test()?;
509
510 TestCase::new(
512 StructArray::from(vec![
513 (
514 Arc::new(Field::new_list(
515 "list",
516 Field::new_list_field(DataType::Int64, true),
517 true,
518 )),
519 Arc::new(ListArray::from_iter_primitive::<Int64Type, _, _>(vec![
520 Some(vec![Some(1)]),
521 None,
522 Some(vec![Some(2), Some(3)]),
523 ])) as ArrayRef,
524 ),
525 (
526 Arc::new(Field::new_struct(
527 "nested",
528 vec![Field::new("int", DataType::Int64, true)],
529 true,
530 )),
531 Arc::new(StructArray::from(vec![(
532 Arc::new(Field::new("int", DataType::Int64, true)),
533 Arc::new(Int64Array::from(vec![-1, -2, -3])) as ArrayRef,
534 )])),
535 ),
536 (
537 Arc::new(Field::new("string", DataType::Utf8, true)),
538 Arc::new(StringArray::from(vec!["a", "b", "c"])),
539 ),
540 ]),
541 Fields::from(vec![
542 Field::new("bool", DataType::Boolean, true),
543 Field::new_list("list", Field::new_list_field(DataType::Int64, true), true),
544 Field::new_struct(
545 "nested",
546 vec![
547 Field::new("float", DataType::Float64, true),
548 Field::new("int", DataType::Int64, true),
549 ],
550 true,
551 ),
552 Field::new("string", DataType::Utf8, true),
553 ]),
554 Ok(vec![
555 Arc::new(BooleanArray::new_null(3)) as ArrayRef,
556 Arc::new(ListArray::from_iter_primitive::<Int64Type, _, _>(vec![
557 Some(vec![Some(1)]),
558 None,
559 Some(vec![Some(2), Some(3)]),
560 ])),
561 Arc::new(StructArray::from(vec![
562 (
563 Arc::new(Field::new("float", DataType::Float64, true)),
564 Arc::new(Float64Array::new_null(3)) as ArrayRef,
565 ),
566 (
567 Arc::new(Field::new("int", DataType::Int64, true)),
568 Arc::new(Int64Array::from(vec![-1, -2, -3])),
569 ),
570 ])),
571 Arc::new(StringArray::from(vec!["a", "b", "c"])),
572 ]),
573 )
574 .test()?;
575
576 Ok(())
577 }
578
579 #[test]
580 fn test_align_variant_to_struct() -> Result<()> {
581 let encode = |json: &[u8]| jsonb::parse_value(json).unwrap().to_vec();
582 let object =
583 encode(br#"{"nested":{"flag":true,"items":[1,2],"raw":{"x":1},"text":42,"value":42}}"#);
584 let scalar = encode(b"1");
585 let variants: ArrayRef = Arc::new(BinaryArray::from(vec![
586 Some(object.as_slice()),
587 None,
588 Some(scalar.as_slice()),
589 ]));
590 let expected_type = DataType::Struct(Fields::from(vec![Field::new_struct(
591 "nested",
592 vec![
593 Field::new("flag", DataType::Boolean, true),
594 Field::new_list("items", Field::new_list_field(DataType::UInt64, true), true),
595 Field::new("raw", DataType::Binary, true),
596 Field::new("text", DataType::Utf8View, true),
597 Field::new("value", DataType::UInt64, true),
598 ],
599 true,
600 )]));
601
602 let aligned = JsonArray::from(&variants).try_align(&expected_type)?;
603 assert_eq!(&expected_type, aligned.data_type());
604 assert_eq!(
605 json!({
606 "nested": {
607 "flag": true,
608 "items": [1, 2],
609 "raw": {"x": 1},
610 "text": "42",
611 "value": 42
612 }
613 }),
614 JsonArray::from(&aligned).try_get_value(0)?
615 );
616 assert!(aligned.is_null(1));
617 assert!(aligned.is_null(2));
618
619 Ok(())
620 }
621
622 #[test]
623 fn test_align_nested_variant_to_struct() -> Result<()> {
624 let object = jsonb::parse_value(br#"{"flag":true,"value":42}"#)
625 .unwrap()
626 .to_vec();
627 let variants: ArrayRef = Arc::new(BinaryArray::from(vec![Some(object.as_slice()), None]));
628 let input: ArrayRef = Arc::new(StructArray::from(vec![(
629 Arc::new(Field::new("nested", DataType::Binary, true)),
630 variants,
631 )]));
632 let expected_type = DataType::Struct(Fields::from(vec![Field::new_struct(
633 "nested",
634 vec![
635 Field::new("flag", DataType::Boolean, true),
636 Field::new("value", DataType::UInt64, true),
637 ],
638 true,
639 )]));
640
641 let aligned = JsonArray::from(&input).try_align(&expected_type)?;
642 assert_eq!(&expected_type, aligned.data_type());
643 assert_eq!(
644 json!({"nested": {"flag": true, "value": 42}}),
645 JsonArray::from(&aligned).try_get_value(0)?
646 );
647 assert_eq!(
648 json!({"nested": null}),
649 JsonArray::from(&aligned).try_get_value(1)?
650 );
651
652 Ok(())
653 }
654}