laminar_connectors/serde/json/
mod.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;