Skip to main content

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