laminar_connectors/serde/csv/
mod.rs1use arrow_array::RecordBatch;
7use arrow_schema::SchemaRef;
8
9use super::{Format, RecordDeserializer, RecordSerializer};
10use crate::error::SerdeError;
11use crate::schema::csv::{CsvDecoder, CsvDecoderConfig, CsvEncoder, CsvEncoderConfig};
12use crate::schema::traits::{FormatDecoder, FormatEncoder};
13use crate::schema::types::RawRecord;
14
15#[derive(Debug, Clone)]
17pub struct CsvDeserializer {
18 delimiter: u8,
19}
20
21impl CsvDeserializer {
22 #[must_use]
24 pub fn new() -> Self {
25 Self { delimiter: b',' }
26 }
27
28 #[must_use]
30 pub fn with_delimiter(delimiter: u8) -> Self {
31 Self { delimiter }
32 }
33}
34
35impl Default for CsvDeserializer {
36 fn default() -> Self {
37 Self::new()
38 }
39}
40
41impl RecordDeserializer for CsvDeserializer {
42 fn deserialize(&self, data: &[u8], schema: &SchemaRef) -> Result<RecordBatch, SerdeError> {
43 let config = CsvDecoderConfig {
44 delimiter: self.delimiter,
45 has_header: false,
46 ..CsvDecoderConfig::default()
47 };
48 let decoder = CsvDecoder::with_config(schema.clone(), config);
49 let record = RawRecord::new(data.to_vec());
50 decoder
51 .decode_one(&record)
52 .map_err(|e| SerdeError::Csv(e.to_string()))
53 }
54
55 fn format(&self) -> Format {
56 Format::Csv
57 }
58}
59
60#[derive(Debug, Clone)]
62pub struct CsvSerializer {
63 delimiter: u8,
64}
65
66impl CsvSerializer {
67 #[must_use]
69 pub fn new() -> Self {
70 Self { delimiter: b',' }
71 }
72
73 #[must_use]
75 pub fn with_delimiter(delimiter: u8) -> Self {
76 Self { delimiter }
77 }
78}
79
80impl Default for CsvSerializer {
81 fn default() -> Self {
82 Self::new()
83 }
84}
85
86impl RecordSerializer for CsvSerializer {
87 fn serialize(&self, batch: &RecordBatch) -> Result<Vec<Vec<u8>>, SerdeError> {
88 let config = CsvEncoderConfig {
89 delimiter: self.delimiter,
90 has_header: false,
91 };
92 let encoder = CsvEncoder::with_config(batch.schema(), config);
93 encoder
94 .encode_batch(batch)
95 .map_err(|e| SerdeError::Csv(e.to_string()))
96 }
97
98 fn format(&self) -> Format {
99 Format::Csv
100 }
101}
102
103#[cfg(test)]
104mod tests;