laminar_connectors/schema/json/decoder/
batch.rs1use 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 pub fn decode_slices(&self, values: &[&[u8]]) -> SchemaResult<RecordBatch> {
22 self.decode_slices_bounded(values, usize::MAX)
23 }
24
25 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 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 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}