Skip to main content

laminar_connectors/serde/
json.rs

1//! JSON serialization and deserialization.
2//!
3//! Implements [`RecordDeserializer`] / [`RecordSerializer`] by delegating
4//! to [`JsonDecoder`] and [`JsonEncoder`].
5
6use std::sync::Arc;
7
8use arrow_array::RecordBatch;
9use arrow_schema::SchemaRef;
10use parking_lot::Mutex;
11use serde_json::Value;
12
13use super::{Format, RecordDeserializer, RecordSerializer};
14use crate::error::SerdeError;
15use crate::schema::json::decoder::JsonDecoder;
16use crate::schema::json::encoder::JsonEncoder;
17use crate::schema::traits::FormatEncoder;
18
19/// JSON record deserializer. Delegates to a cached [`JsonDecoder`].
20#[derive(Debug, Clone)]
21pub struct JsonDeserializer {
22    #[allow(clippy::type_complexity)]
23    decoder: Arc<Mutex<Option<(SchemaRef, Arc<JsonDecoder>)>>>,
24}
25
26impl JsonDeserializer {
27    /// Creates a new JSON deserializer.
28    #[must_use]
29    pub fn new() -> Self {
30        Self {
31            decoder: Arc::new(Mutex::new(None)),
32        }
33    }
34
35    /// Returns a cached [`JsonDecoder`] for `schema`, rebuilding it only when
36    /// the cached schema differs.
37    fn decoder_for(&self, schema: &SchemaRef) -> Arc<JsonDecoder> {
38        let mut cache = self.decoder.lock();
39        if let Some((cached_schema, cached)) = cache.as_ref() {
40            if Arc::ptr_eq(cached_schema, schema) || cached_schema == schema {
41                return cached.clone();
42            }
43        }
44        let decoder = Arc::new(JsonDecoder::new(schema.clone()));
45        *cache = Some((schema.clone(), decoder.clone()));
46        decoder
47    }
48
49    /// Deserializes a pre-parsed JSON [`Value`] into a [`RecordBatch`].
50    ///
51    /// Used by [`DebeziumDeserializer`](super::debezium::DebeziumDeserializer)
52    /// to avoid double-parsing the envelope.
53    ///
54    /// # Errors
55    ///
56    /// Returns `SerdeError` if the value cannot be decoded.
57    pub fn deserialize_value(
58        &self,
59        value: &Value,
60        schema: &SchemaRef,
61    ) -> Result<RecordBatch, SerdeError> {
62        let bytes = serde_json::to_vec(value).map_err(|e| SerdeError::Json(e.to_string()))?;
63        self.deserialize(&bytes, schema)
64    }
65}
66
67impl Default for JsonDeserializer {
68    fn default() -> Self {
69        Self::new()
70    }
71}
72
73impl RecordDeserializer for JsonDeserializer {
74    fn deserialize(&self, data: &[u8], schema: &SchemaRef) -> Result<RecordBatch, SerdeError> {
75        let d = self.decoder_for(schema);
76        d.decode_slices(&[data])
77            .map_err(|e| SerdeError::Json(e.to_string()))
78    }
79
80    fn deserialize_batch(
81        &self,
82        records: &[&[u8]],
83        schema: &SchemaRef,
84    ) -> Result<RecordBatch, SerdeError> {
85        if records.is_empty() {
86            return Ok(RecordBatch::new_empty(schema.clone()));
87        }
88        let d = self.decoder_for(schema);
89        d.decode_slices(records)
90            .map_err(|e| SerdeError::Json(e.to_string()))
91    }
92
93    fn format(&self) -> Format {
94        Format::Json
95    }
96}
97
98/// JSON record serializer. Delegates to [`JsonEncoder`].
99#[derive(Debug, Clone)]
100pub struct JsonSerializer {
101    _private: (),
102}
103
104impl JsonSerializer {
105    /// Creates a new JSON serializer.
106    #[must_use]
107    pub fn new() -> Self {
108        Self { _private: () }
109    }
110}
111
112impl Default for JsonSerializer {
113    fn default() -> Self {
114        Self::new()
115    }
116}
117
118impl RecordSerializer for JsonSerializer {
119    fn serialize(&self, batch: &RecordBatch) -> Result<Vec<Vec<u8>>, SerdeError> {
120        let encoder = JsonEncoder::new(batch.schema());
121        encoder
122            .encode_batch(batch)
123            .map_err(|e| SerdeError::Json(e.to_string()))
124    }
125
126    fn format(&self) -> Format {
127        Format::Json
128    }
129}
130
131#[cfg(test)]
132mod tests {
133    use super::*;
134    use arrow_schema::{DataType, Field, Schema};
135    use std::sync::Arc;
136
137    fn test_schema() -> SchemaRef {
138        Arc::new(Schema::new(vec![
139            Field::new("id", DataType::Int64, false),
140            Field::new("name", DataType::Utf8, false),
141            Field::new("score", DataType::Float64, true),
142        ]))
143    }
144
145    #[test]
146    fn test_json_deserialize_basic() {
147        let deser = JsonDeserializer::new();
148        let schema = test_schema();
149        let data = br#"{"id": 1, "name": "Alice", "score": 95.5}"#;
150
151        let batch = deser.deserialize(data, &schema).unwrap();
152        assert_eq!(batch.num_rows(), 1);
153        assert_eq!(batch.num_columns(), 3);
154    }
155
156    #[test]
157    fn test_json_serialize_roundtrip() {
158        let deser = JsonDeserializer::new();
159        let ser = JsonSerializer::new();
160        let schema = test_schema();
161
162        let data = br#"{"id": 42, "name": "Charlie", "score": 88.5}"#;
163        let batch = deser.deserialize(data, &schema).unwrap();
164
165        let serialized = ser.serialize(&batch).unwrap();
166        assert_eq!(serialized.len(), 1);
167
168        let roundtrip: Value = serde_json::from_slice(&serialized[0]).unwrap();
169        assert_eq!(roundtrip["id"], 42);
170        assert_eq!(roundtrip["name"], "Charlie");
171    }
172
173    #[test]
174    fn test_json_deserialize_batch() {
175        let deser = JsonDeserializer::new();
176        let schema = test_schema();
177
178        let r1 = br#"{"id": 1, "name": "A", "score": 10.0}"#;
179        let r2 = br#"{"id": 2, "name": "B", "score": 20.0}"#;
180        let records: Vec<&[u8]> = vec![r1, r2];
181
182        let batch = deser.deserialize_batch(&records, &schema).unwrap();
183        assert_eq!(batch.num_rows(), 2);
184    }
185
186    #[test]
187    fn test_json_deserialize_coercion() {
188        let deser = JsonDeserializer::new();
189        let schema = Arc::new(Schema::new(vec![
190            Field::new("qty", DataType::Int64, false),
191            Field::new("price", DataType::Float64, false),
192        ]));
193
194        let data = br#"{"qty": "100", "price": "187.52"}"#;
195        let batch = deser.deserialize(data, &schema).unwrap();
196
197        let qty = batch
198            .column(0)
199            .as_any()
200            .downcast_ref::<arrow_array::Int64Array>()
201            .unwrap();
202        assert_eq!(qty.value(0), 100);
203    }
204}