diff --git a/crates/adapterlib/src/catalog.rs b/crates/adapterlib/src/catalog.rs index 6260edc6c33..29f2c6d8954 100644 --- a/crates/adapterlib/src/catalog.rs +++ b/crates/adapterlib/src/catalog.rs @@ -21,7 +21,7 @@ use dbsp::dynamic::{ClonableTrait, DynData, DynVec, Erase, Factory}; use dbsp::operator::StagedBuffers; use dyn_clone::DynClone; use feldera_sqllib::Variant; -use feldera_types::format::csv::CsvParserConfig; +use feldera_types::format::csv::CsvFormatConfig; use feldera_types::format::json::JsonFlavor; use feldera_types::program_schema::{Relation, SqlIdentifier}; use feldera_types::serde_with_context::SqlSerdeConfig; @@ -43,7 +43,7 @@ pub enum RecordFormat { // raw encoding of this column only. This is particularly useful for // tables that store raw JSON or binary data to be parsed using SQL. Json(JsonFlavor), - Csv(CsvParserConfig), + Csv(CsvFormatConfig), Parquet(SqlSerdeConfig), #[cfg(feature = "with-avro")] Avro, diff --git a/crates/adapters/src/format/csv.rs b/crates/adapters/src/format/csv.rs index 03d6fd01f54..05c4b435648 100644 --- a/crates/adapters/src/format/csv.rs +++ b/crates/adapters/src/format/csv.rs @@ -12,7 +12,7 @@ use feldera_adapterlib::{ConnectorMetadata, catalog::SerCursorFlattened}; use feldera_sqllib::Variant; use feldera_types::{ config::ConnectorConfig, - format::csv::{CsvEncoderConfig, CsvParserConfig}, + format::csv::{CsvEncoderConfig, CsvFormatConfig}, }; use serde::Deserialize; use serde_json::json; @@ -45,7 +45,7 @@ fn csv_trim(t: &feldera_types::format::csv::CsvTrim) -> csv::Trim { /// /// `has_headers` is always `false` — the adapter layer skips the header row /// itself and must not pass it to the underlying CSV reader. -pub(crate) fn csv_reader_builder(config: &CsvParserConfig) -> csv::ReaderBuilder { +pub(crate) fn csv_reader_builder(config: &CsvFormatConfig) -> csv::ReaderBuilder { let mut b = csv::ReaderBuilder::new(); b.has_headers(false) .delimiter(config.delimiter().0) @@ -75,7 +75,7 @@ impl InputFormat for CsvInputFormat { _endpoint_name: &str, _request: &HttpRequest, ) -> Result, ControllerError> { - Ok(Box::new(CsvParserConfig::default())) + Ok(Box::new(CsvFormatConfig::default())) } fn new_parser( @@ -85,7 +85,7 @@ impl InputFormat for CsvInputFormat { config: &serde_json::Value, ) -> Result, ControllerError> { let config = if config.is_null() { &json!({}) } else { config }; - let config = CsvParserConfig::deserialize(config).map_err(|e| { + let config = CsvFormatConfig::deserialize(config).map_err(|e| { ControllerError::parser_config_parse_error( endpoint_name, &e, @@ -121,7 +121,7 @@ struct CsvParser { } impl CsvParser { - fn new(input_stream: Box, config: &CsvParserConfig) -> Self { + fn new(input_stream: Box, config: &CsvFormatConfig) -> Self { Self { input_stream, last_event_number: 0, @@ -404,13 +404,13 @@ impl Encoder for CsvEncoder { let mut cursor = if self.input_is_indexed { CursorWithPolarity::new(Box::new(SerCursorFlattened::new(batch.cursor( - RecordFormat::Csv(CsvParserConfig { + RecordFormat::Csv(CsvFormatConfig { delimiter: self.config.delimiter, ..Default::default() }), )?))) } else { - CursorWithPolarity::new(batch.cursor(RecordFormat::Csv(CsvParserConfig { + CursorWithPolarity::new(batch.cursor(RecordFormat::Csv(CsvFormatConfig { delimiter: self.config.delimiter, ..Default::default() }))?) @@ -482,7 +482,7 @@ mod test { use super::CsvSplitter; use crate::format::Splitter; - use feldera_types::format::csv::{CsvParserConfig, CsvTrim}; + use feldera_types::format::csv::{CsvFormatConfig, CsvTrim}; #[derive(Debug, Eq, PartialEq)] #[allow(non_snake_case)] @@ -630,7 +630,7 @@ true,bar,buzz"#; /// Push `data` through a fully-configured `CsvParser` and return the /// inserted `TwoStrings` records. - fn e2e_parse_csv(config: CsvParserConfig, data: &[u8]) -> Vec { + fn e2e_parse_csv(config: CsvFormatConfig, data: &[u8]) -> Vec { let format_config = FormatConfig { name: Cow::from("csv"), config: serde_json::to_value(config).unwrap(), @@ -655,7 +655,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_quote() { let rows = e2e_parse_csv( - CsvParserConfig { + CsvFormatConfig { quote: '\'', ..Default::default() }, @@ -682,7 +682,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_escape() { let rows = e2e_parse_csv( - CsvParserConfig { + CsvFormatConfig { escape: Some('\\'), double_quote: false, ..Default::default() @@ -699,7 +699,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_double_quote() { let rows = e2e_parse_csv( - CsvParserConfig::default(), + CsvFormatConfig::default(), // "foo""bar" → foo"bar b"\"foo\"\"bar\",baz\n", ); @@ -714,7 +714,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_double_quote_disabled() { let rows = e2e_parse_csv( - CsvParserConfig { + CsvFormatConfig { double_quote: false, ..Default::default() }, @@ -728,7 +728,7 @@ true,bar,buzz"#; // `trim = None` (default): leading and trailing whitespace is preserved. #[test] fn csv_e2e_trim_none() { - let rows = e2e_parse_csv(CsvParserConfig::default(), b" foo , bar \n"); + let rows = e2e_parse_csv(CsvFormatConfig::default(), b" foo , bar \n"); assert_eq!(rows.len(), 1); assert_eq!(rows[0].first, " foo "); assert_eq!(rows[0].second, " bar "); @@ -738,7 +738,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_quoting_disabled() { let rows = e2e_parse_csv( - CsvParserConfig { + CsvFormatConfig { quoting: false, ..Default::default() }, @@ -754,7 +754,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_comment() { let rows = e2e_parse_csv( - CsvParserConfig { + CsvFormatConfig { comment: Some('#'), ..Default::default() }, @@ -770,7 +770,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_flexible_true() { let rows = e2e_parse_csv( - CsvParserConfig::default(), // flexible=true + CsvFormatConfig::default(), // flexible=true b"foo,bar\nqux,quux,extra\n", ); assert_eq!(rows.len(), 2); @@ -785,7 +785,7 @@ true,bar,buzz"#; fn csv_e2e_flexible_false() { let format_config = FormatConfig { name: Cow::from("csv"), - config: serde_json::to_value(CsvParserConfig { + config: serde_json::to_value(CsvFormatConfig { flexible: false, ..Default::default() }) @@ -811,7 +811,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_trim_fields() { let rows = e2e_parse_csv( - CsvParserConfig { + CsvFormatConfig { trim: CsvTrim::Fields, ..Default::default() }, @@ -826,7 +826,7 @@ true,bar,buzz"#; #[test] fn csv_e2e_headers() { let rows = e2e_parse_csv( - CsvParserConfig { + CsvFormatConfig { headers: true, ..Default::default() }, diff --git a/crates/adapters/src/static_compile/deinput.rs b/crates/adapters/src/static_compile/deinput.rs index 671ab9dc701..03f6f7f8553 100644 --- a/crates/adapters/src/static_compile/deinput.rs +++ b/crates/adapters/src/static_compile/deinput.rs @@ -25,7 +25,7 @@ use erased_serde::Deserializer as ErasedDeserializer; use feldera_adapterlib::catalog::AvroSchemaRefs; use feldera_adapterlib::format::{BufferSize, flatten_nested}; use feldera_sqllib::Variant; -use feldera_types::format::csv::CsvParserConfig; +use feldera_types::format::csv::CsvFormatConfig; use feldera_types::serde_with_context::{DeserializeWithContext, SqlSerdeConfig}; use serde_arrow::Deserializer as ArrowDeserializer; use serde_json::de::SliceRead; @@ -94,8 +94,8 @@ pub struct CsvDeserializerFromBytes { config: SqlSerdeConfig, } -impl DeserializerFromBytes<(SqlSerdeConfig, CsvParserConfig)> for CsvDeserializerFromBytes { - fn create((serde_config, csv_config): (SqlSerdeConfig, CsvParserConfig)) -> Self { +impl DeserializerFromBytes<(SqlSerdeConfig, CsvFormatConfig)> for CsvDeserializerFromBytes { + fn create((serde_config, csv_config): (SqlSerdeConfig, CsvFormatConfig)) -> Self { CsvDeserializerFromBytes { // `csv_reader_builder` sets `has_headers(false)` — the adapter // layer handles header-row skipping itself. diff --git a/crates/adapters/src/static_compile/seroutput.rs b/crates/adapters/src/static_compile/seroutput.rs index 574093ba546..e9afa7de941 100644 --- a/crates/adapters/src/static_compile/seroutput.rs +++ b/crates/adapters/src/static_compile/seroutput.rs @@ -41,7 +41,7 @@ use feldera_types::serde_with_context::serialize::{ SerializeFieldsWithContextWrapper, SerializeWithContextWrapper, }; use feldera_types::{ - format::csv::CsvParserConfig, + format::csv::CsvFormatConfig, serde_with_context::{SerializeWithContext, SqlSerdeConfig}, }; use rand::thread_rng; @@ -197,7 +197,7 @@ struct CsvSerializer { } impl CsvSerializer { - fn create(config: CsvParserConfig) -> Self { + fn create(config: CsvFormatConfig) -> Self { Self { writer: CsvWriterBuilder::new() .has_headers(false) diff --git a/crates/feldera-types/src/format/csv.rs b/crates/feldera-types/src/format/csv.rs index 4f08ac2e625..2e588aaa5d7 100644 --- a/crates/feldera-types/src/format/csv.rs +++ b/crates/feldera-types/src/format/csv.rs @@ -20,19 +20,25 @@ pub enum CsvTrim { #[derive(Clone, Debug, Deserialize, Serialize, ToSchema, PartialEq)] #[serde(default)] -pub struct CsvParserConfig { +pub struct CsvFormatConfig { /// Field delimiter (default `','`). /// /// Must be an ASCII character. + /// + /// Used by: input and output. pub delimiter: char, /// Whether the input begins with a header line (which is skipped). + /// + /// Used by: input only. pub headers: bool, /// The quote character (default `'"'`). /// /// Must be an ASCII character. Set `quoting` to `false` to disable /// quoting entirely. + /// + /// Used by: input only. pub quote: char, /// The escape character for quoted fields (default: `None`). @@ -44,6 +50,8 @@ pub struct CsvParserConfig { /// escaping). /// /// Must be an ASCII character. + /// + /// Used by: input only. pub escape: Option, /// Enable double-quote escaping (default `true`). @@ -51,37 +59,47 @@ pub struct CsvParserConfig { /// When `true`, a quote character inside a quoted field may be escaped by /// doubling it. Setting this to `false` disables double-quote escaping /// (an explicit `escape` character can still be used). + /// + /// Used by: input only. pub double_quote: bool, /// Enable quoting (default `true`). /// /// When `false`, the `quote` and `escape` characters have no special /// meaning and every newline terminates a record regardless of context. + /// + /// Used by: input only. pub quoting: bool, /// Comment character (default: `None`). /// /// When set, lines whose first byte matches this character are treated as /// comments and skipped entirely. Must be an ASCII character. + /// + /// Used by: input only. pub comment: Option, /// Allow records with a variable number of fields (default `true`). /// /// When `true`, records that have fewer or more fields than expected are /// accepted rather than treated as errors. + /// + /// Used by: input only. pub flexible: bool, /// Whitespace trimming policy (default [`CsvTrim::None`]). + /// + /// Used by: input only. pub trim: CsvTrim, } -impl CsvParserConfig { +impl CsvFormatConfig { pub fn delimiter(&self) -> CsvDelimiter { self.delimiter.into() } } -impl Default for CsvParserConfig { +impl Default for CsvFormatConfig { fn default() -> Self { Self { delimiter: CsvDelimiter::default().0.into(), diff --git a/sql-to-dbsp-compiler/SQL-compiler/src/test/java/org/dbsp/sqlCompiler/compiler/sql/OtherTests.java b/sql-to-dbsp-compiler/SQL-compiler/src/test/java/org/dbsp/sqlCompiler/compiler/sql/OtherTests.java index 989ee2409b5..75b8e801228 100644 --- a/sql-to-dbsp-compiler/SQL-compiler/src/test/java/org/dbsp/sqlCompiler/compiler/sql/OtherTests.java +++ b/sql-to-dbsp-compiler/SQL-compiler/src/test/java/org/dbsp/sqlCompiler/compiler/sql/OtherTests.java @@ -436,7 +436,7 @@ pub fn test() { #[test] pub fn test() { use dbsp_adapters::{CircuitCatalog, RecordFormat}; - use feldera_types::format::csv::CsvParserConfig; + use feldera_types::format::csv::CsvFormatConfig; let (mut circuit, catalog) = circuit(CircuitConfig::with_workers(2)) .expect("Failed to build circuit"); @@ -470,7 +470,7 @@ pub fn test() { // Read the produced output let reader = adult.concat(); let mut cursor = reader - .cursor(RecordFormat::Csv(CsvParserConfig::default())) + .cursor(RecordFormat::Csv(CsvFormatConfig::default())) .unwrap(); while cursor.key_valid() { let mut w = cursor.weight(); diff --git a/sql-to-dbsp-compiler/using.md b/sql-to-dbsp-compiler/using.md index eff8552be32..3f2b2ea4014 100644 --- a/sql-to-dbsp-compiler/using.md +++ b/sql-to-dbsp-compiler/using.md @@ -396,7 +396,7 @@ We exercise this circuit by inserting data using a CSV format: #[test] pub fn test() { use dbsp_adapters::{CircuitCatalog, RecordFormat}; - use use feldera_types::format::csv::CsvParserConfig; + use feldera_types::format::csv::CsvFormatConfig; let (mut circuit, catalog) = circuit(2) .expect("Failed to build circuit"); @@ -430,7 +430,7 @@ pub fn test() { // Read the produced output let reader = adult.concat().consolidate(); let mut cursor = reader - .cursor(RecordFormat::Csv(CsvParserConfig::default())) + .cursor(RecordFormat::Csv(CsvFormatConfig::default())) .unwrap(); while cursor.key_valid() { let mut w = cursor.weight();