laminar_connectors/serde/
json.rs1use 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#[derive(Debug, Clone)]
21pub struct JsonDeserializer {
22 #[allow(clippy::type_complexity)]
23 decoder: Arc<Mutex<Option<(SchemaRef, Arc<JsonDecoder>)>>>,
24}
25
26impl JsonDeserializer {
27 #[must_use]
29 pub fn new() -> Self {
30 Self {
31 decoder: Arc::new(Mutex::new(None)),
32 }
33 }
34
35 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 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#[derive(Debug, Clone)]
100pub struct JsonSerializer {
101 _private: (),
102}
103
104impl JsonSerializer {
105 #[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}