Skip to main content

micromegas_datafusion_extensions/jsonb/
parse.rs

1use 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/// A scalar UDF that parses a JSON string into a JSONB value.
13///
14/// Accepts both Utf8 and Dictionary<Int32, Utf8> inputs.
15/// Returns Dictionary<Int32, Binary> for memory efficiency.
16#[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
143/// Creates a user-defined function to parse a JSON string into a JSONB value.
144pub fn make_jsonb_parse_udf() -> ScalarUDF {
145    ScalarUDF::new_from_impl(JsonbParse::new())
146}