Skip to main content

laminar_connectors/schema/json/decoder/
batch.rs

1//! Per-batch JSON parsing, expansion, and Arrow column assembly.
2
3use std::sync::Arc;
4
5use arrow_array::builder::LargeBinaryBuilder;
6use arrow_array::{ArrayRef, RecordBatch};
7use arrow_schema::DataType;
8
9use super::value::{append_null, append_value, create_builders, ColumnBuilder};
10use super::{ColumnExtraction, JsonDecoder, UnknownFieldStrategy};
11use crate::schema::error::{SchemaError, SchemaResult};
12use crate::schema::json::jsonb::JsonbEncoder;
13
14impl JsonDecoder {
15    /// Decode borrowed JSON payload slices into a `RecordBatch` without copying
16    /// the input payloads.
17    ///
18    /// # Errors
19    /// Returns a decode error if a payload is invalid or cannot be represented
20    /// by the configured Arrow schema.
21    pub fn decode_slices(&self, values: &[&[u8]]) -> SchemaResult<RecordBatch> {
22        self.decode_slices_bounded(values, usize::MAX)
23    }
24
25    /// Decode borrowed JSON payloads while rejecting output expansion beyond
26    /// `max_rows` before values are appended to Arrow builders.
27    ///
28    /// # Errors
29    /// Returns a decode error when parsing, coercion, or `json.explode` exceeds
30    /// the row bound.
31    pub fn decode_slices_bounded(
32        &self,
33        values: &[&[u8]],
34        max_rows: usize,
35    ) -> SchemaResult<RecordBatch> {
36        if values.is_empty() {
37            return Ok(RecordBatch::new_empty(self.schema.clone()));
38        }
39
40        let mut batch = BatchDecodeState::new(self, values.len(), max_rows);
41        for &bytes in values {
42            batch.decode_record(bytes)?;
43        }
44        batch.finish()
45    }
46
47    /// O(n) field lookup. Linear scan is faster for typical schemas with fewer
48    /// than 50 fields.
49    fn field_index(&self, name: &str) -> Option<usize> {
50        self.field_indices
51            .iter()
52            .find(|(field, _)| field == name)
53            .map(|(_, index)| *index)
54    }
55}
56
57struct BatchDecodeState<'a> {
58    decoder: &'a JsonDecoder,
59    builders: Vec<Box<dyn ColumnBuilder>>,
60    extra_builder: Option<LargeBinaryBuilder>,
61    jsonb_encoder: Option<JsonbEncoder>,
62    populated: Vec<bool>,
63    output_rows: usize,
64    max_rows: usize,
65}
66
67impl<'a> BatchDecodeState<'a> {
68    fn new(decoder: &'a JsonDecoder, capacity: usize, max_rows: usize) -> Self {
69        let extra_builder = matches!(
70            decoder.config.unknown_fields,
71            UnknownFieldStrategy::CollectExtra
72        )
73        .then(|| LargeBinaryBuilder::with_capacity(capacity, capacity * 64));
74        let jsonb_encoder = decoder.config.nested_as_jsonb.then(JsonbEncoder::new);
75        Self {
76            decoder,
77            builders: create_builders(&decoder.schema, capacity),
78            extra_builder,
79            jsonb_encoder,
80            populated: vec![false; decoder.schema.fields().len()],
81            output_rows: 0,
82            max_rows,
83        }
84    }
85
86    #[inline]
87    fn decode_record(&mut self, bytes: &[u8]) -> SchemaResult<()> {
88        let value: serde_json::Value = serde_json::from_slice(bytes)
89            .map_err(|error| SchemaError::DecodeError(format!("JSON parse error: {error}")))?;
90        let Some(default_target) =
91            navigate_path_opt(&value, self.decoder.config.json_path.as_deref())
92        else {
93            // Some sources interleave non-data frames that lack the configured path.
94            return Ok(());
95        };
96
97        if self.decoder.explode_col_indices.is_some() {
98            self.decode_exploded(default_target)
99        } else {
100            self.decode_object(&value, default_target)
101        }
102    }
103
104    fn decode_exploded(&mut self, target: &serde_json::Value) -> SchemaResult<()> {
105        let elements = target.as_array().ok_or_else(|| {
106            SchemaError::DecodeError("json.explode target must be an array".into())
107        })?;
108        self.reserve_rows(elements.len())?;
109        for element in elements {
110            self.populated.fill(false);
111            self.append_exploded_element(element)?;
112            self.append_missing_fields();
113            if let Some(extra) = &mut self.extra_builder {
114                extra.append_null();
115            }
116        }
117        Ok(())
118    }
119
120    fn append_exploded_element(&mut self, element: &serde_json::Value) -> SchemaResult<()> {
121        match element {
122            serde_json::Value::Array(items) => {
123                let positions = self
124                    .decoder
125                    .explode_col_indices
126                    .as_ref()
127                    .expect("explode positions exist in explode mode")
128                    .len();
129                for position in 0..positions {
130                    let column = self
131                        .decoder
132                        .explode_col_indices
133                        .as_ref()
134                        .and_then(|indices| indices.get(position).copied().flatten());
135                    if let Some(column) = column {
136                        let value = items.get(position).unwrap_or(&serde_json::Value::Null);
137                        self.append_column(column, value)?;
138                        self.populated[column] = true;
139                    }
140                }
141            }
142            serde_json::Value::Object(object) => {
143                for column in 0..self.decoder.field_indices.len() {
144                    let name = self.decoder.field_indices[column].0.as_str();
145                    if let Some(value) = object.get(name) {
146                        self.append_column(column, value)?;
147                        self.populated[column] = true;
148                    }
149                }
150            }
151            _ => {
152                return Err(SchemaError::DecodeError(
153                    "json.explode array elements must be arrays or objects".into(),
154                ));
155            }
156        }
157        Ok(())
158    }
159
160    fn decode_object(
161        &mut self,
162        root: &serde_json::Value,
163        target: &serde_json::Value,
164    ) -> SchemaResult<()> {
165        self.reserve_rows(1)?;
166        let object = target
167            .as_object()
168            .ok_or_else(|| SchemaError::DecodeError("JSON value must be an object".into()))?;
169        self.populated.fill(false);
170        self.append_object_fields(root, object)?;
171        let extra_fields = self.collect_or_reject_unknown_fields(object)?;
172        self.append_missing_fields();
173        self.append_extra_fields(extra_fields.as_ref());
174        Ok(())
175    }
176
177    fn append_object_fields(
178        &mut self,
179        root: &serde_json::Value,
180        object: &serde_json::Map<String, serde_json::Value>,
181    ) -> SchemaResult<()> {
182        for column in 0..self.decoder.field_indices.len() {
183            let name = self.decoder.field_indices[column].0.as_str();
184            let value = match &self.decoder.column_extractions[column] {
185                ColumnExtraction::DefaultPath => object.get(name),
186                ColumnExtraction::CustomPath { segments } => navigate_path(root, segments),
187            };
188            if let Some(value) = value {
189                self.append_column(column, value)?;
190                self.populated[column] = true;
191            }
192        }
193        Ok(())
194    }
195
196    fn collect_or_reject_unknown_fields(
197        &self,
198        object: &serde_json::Map<String, serde_json::Value>,
199    ) -> SchemaResult<Option<serde_json::Map<String, serde_json::Value>>> {
200        match self.decoder.config.unknown_fields {
201            UnknownFieldStrategy::CollectExtra => {
202                let mut extra = serde_json::Map::new();
203                for (key, value) in object {
204                    if self.decoder.field_index(key).is_none() {
205                        extra.insert(key.clone(), value.clone());
206                    }
207                }
208                Ok(Some(extra))
209            }
210            UnknownFieldStrategy::Reject => {
211                for key in object.keys() {
212                    if self.decoder.field_index(key).is_none() {
213                        return Err(SchemaError::DecodeError(format!(
214                            "unknown field '{key}' not in schema"
215                        )));
216                    }
217                }
218                Ok(None)
219            }
220            UnknownFieldStrategy::Ignore => Ok(None),
221        }
222    }
223
224    #[inline]
225    fn append_column(&mut self, column: usize, value: &serde_json::Value) -> SchemaResult<()> {
226        let field = &self.decoder.schema.fields()[column];
227        append_value(
228            &mut self.builders[column],
229            field.data_type(),
230            value,
231            &self.decoder.config,
232            &self.decoder.mismatch_count,
233            self.jsonb_encoder.as_mut(),
234            self.decoder.column_epoch_units[column],
235        )
236    }
237
238    #[inline]
239    fn append_missing_fields(&mut self) {
240        for (column, populated) in self.populated.iter().enumerate() {
241            if !populated {
242                append_null(&mut self.builders[column]);
243            }
244        }
245    }
246
247    fn append_extra_fields(
248        &mut self,
249        extra_fields: Option<&serde_json::Map<String, serde_json::Value>>,
250    ) {
251        let Some(extra_builder) = &mut self.extra_builder else {
252            return;
253        };
254        let Some(extra_fields) = extra_fields.filter(|extra| !extra.is_empty()) else {
255            extra_builder.append_null();
256            return;
257        };
258
259        let mut encoder = self
260            .jsonb_encoder
261            .as_mut()
262            .map_or_else(JsonbEncoder::new, |_| JsonbEncoder::new());
263        let bytes = encoder.encode(&serde_json::Value::Object(extra_fields.clone()));
264        extra_builder.append_value(&bytes);
265    }
266
267    fn reserve_rows(&mut self, additional: usize) -> SchemaResult<()> {
268        self.output_rows = self
269            .output_rows
270            .checked_add(additional)
271            .filter(|rows| *rows <= self.max_rows)
272            .ok_or_else(|| {
273                SchemaError::DecodeError(format!(
274                    "JSON output exceeds the {}-row batch limit",
275                    self.max_rows
276                ))
277            })?;
278        Ok(())
279    }
280
281    fn finish(self) -> SchemaResult<RecordBatch> {
282        let mut columns: Vec<ArrayRef> = self
283            .builders
284            .into_iter()
285            .map(|mut builder| builder.finish())
286            .collect();
287        let final_schema = if let Some(mut extra_builder) = self.extra_builder {
288            columns.push(Arc::new(extra_builder.finish()));
289            let mut fields = self.decoder.schema.fields().to_vec();
290            fields.push(Arc::new(arrow_schema::Field::new(
291                "_extra",
292                DataType::LargeBinary,
293                true,
294            )));
295            Arc::new(arrow_schema::Schema::new(fields))
296        } else {
297            self.decoder.schema.clone()
298        };
299        RecordBatch::try_new(final_schema, columns)
300            .map_err(|error| SchemaError::DecodeError(format!("RecordBatch construction: {error}")))
301    }
302}
303
304fn navigate_path<'a>(
305    root: &'a serde_json::Value,
306    segments: &[String],
307) -> Option<&'a serde_json::Value> {
308    let mut current = root;
309    for segment in segments {
310        current = current.get(segment.as_str())?;
311    }
312    Some(current)
313}
314
315fn navigate_path_opt<'a>(
316    root: &'a serde_json::Value,
317    segments: Option<&[String]>,
318) -> Option<&'a serde_json::Value> {
319    match segments {
320        Some(segments) => navigate_path(root, segments),
321        None => Some(root),
322    }
323}