micromegas_datafusion_extensions/jsonb/
parse.rs1use datafusion::arrow::array::{Array, BinaryDictionaryBuilder, DictionaryArray, StringArray};
2use datafusion::arrow::datatypes::{DataType, Int32Type};
3use datafusion::common::{Result, internal_err};
4use datafusion::error::DataFusionError;
5use datafusion::logical_expr::{
6 ColumnarValue, ScalarFunctionArgs, ScalarUDF, ScalarUDFImpl, Signature, Volatility,
7};
8use jsonb::parse_value;
9use micromegas_tracing::warn;
10use std::sync::Arc;
11
12#[derive(Debug, PartialEq, Eq, Hash)]
17pub struct JsonbParse {
18 signature: Signature,
19}
20
21impl JsonbParse {
22 pub fn new() -> Self {
23 Self {
24 signature: Signature::any(1, Volatility::Immutable),
25 }
26 }
27}
28
29impl Default for JsonbParse {
30 fn default() -> Self {
31 Self::new()
32 }
33}
34
35fn parse_json_to_jsonb(json_str: &str) -> Option<Vec<u8>> {
36 match parse_value(json_str.as_bytes()) {
37 Ok(parsed) => {
38 let mut buffer = vec![];
39 parsed.write_to_vec(&mut buffer);
40 Some(buffer)
41 }
42 Err(e) => {
43 warn!("error parsing json={json_str} error={e:?}");
44 None
45 }
46 }
47}
48
49impl ScalarUDFImpl for JsonbParse {
50 fn name(&self) -> &str {
51 "jsonb_parse"
52 }
53
54 fn signature(&self) -> &Signature {
55 &self.signature
56 }
57
58 fn return_type(&self, _args: &[DataType]) -> Result<DataType> {
59 Ok(DataType::Dictionary(
60 Box::new(DataType::Int32),
61 Box::new(DataType::Binary),
62 ))
63 }
64
65 fn invoke_with_args(&self, args: ScalarFunctionArgs) -> Result<ColumnarValue> {
66 let args = ColumnarValue::values_to_arrays(&args.args)?;
67 if args.len() != 1 {
68 return internal_err!("wrong number of arguments to jsonb_parse()");
69 }
70
71 match args[0].data_type() {
72 DataType::Utf8 => {
73 let string_array =
74 args[0]
75 .as_any()
76 .downcast_ref::<StringArray>()
77 .ok_or_else(|| {
78 DataFusionError::Internal("error casting to string array".into())
79 })?;
80
81 let mut dict_builder = BinaryDictionaryBuilder::<Int32Type>::new();
82 for i in 0..string_array.len() {
83 if string_array.is_null(i) {
84 dict_builder.append_null();
85 } else {
86 let json_str = string_array.value(i);
87 if let Some(jsonb_bytes) = parse_json_to_jsonb(json_str) {
88 dict_builder.append_value(&jsonb_bytes);
89 } else {
90 dict_builder.append_null();
91 }
92 }
93 }
94 Ok(ColumnarValue::Array(Arc::new(dict_builder.finish())))
95 }
96 DataType::Dictionary(_, value_type)
97 if matches!(value_type.as_ref(), DataType::Utf8) =>
98 {
99 let dict_array = args[0]
100 .as_any()
101 .downcast_ref::<DictionaryArray<Int32Type>>()
102 .ok_or_else(|| {
103 DataFusionError::Internal("error casting dictionary array".into())
104 })?;
105
106 let string_values = dict_array
107 .values()
108 .as_any()
109 .downcast_ref::<StringArray>()
110 .ok_or_else(|| {
111 DataFusionError::Internal("dictionary values are not a string array".into())
112 })?;
113
114 let mut dict_builder = BinaryDictionaryBuilder::<Int32Type>::new();
115 for i in 0..dict_array.len() {
116 if dict_array.is_null(i) {
117 dict_builder.append_null();
118 } else {
119 let key_index = dict_array.keys().value(i) as usize;
120 if key_index < string_values.len() {
121 let json_str = string_values.value(key_index);
122 if let Some(jsonb_bytes) = parse_json_to_jsonb(json_str) {
123 dict_builder.append_value(&jsonb_bytes);
124 } else {
125 dict_builder.append_null();
126 }
127 } else {
128 return internal_err!(
129 "Dictionary key index out of bounds in jsonb_parse"
130 );
131 }
132 }
133 }
134 Ok(ColumnarValue::Array(Arc::new(dict_builder.finish())))
135 }
136 _ => internal_err!(
137 "jsonb_parse: unsupported input type, expected Utf8 or Dictionary<Int32, Utf8>"
138 ),
139 }
140 }
141}
142
143pub fn make_jsonb_parse_udf() -> ScalarUDF {
145 ScalarUDF::new_from_impl(JsonbParse::new())
146}