-
Notifications
You must be signed in to change notification settings - Fork 203
fix: validate CSV expression function options like Spark #2255
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
c98ae60
2999fcc
0c4740b
4954070
094b912
ec9166e
ee671d3
ff14633
f28d04e
b2dcd21
6e3349a
ed78a05
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,3 +1,4 @@ | ||
| mod options; | ||
| mod schema_of_csv; | ||
| pub mod spark_from_csv; | ||
| pub mod spark_to_csv; | ||
|
|
||
| Original file line number | Diff line number | Diff line change | |||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| @@ -0,0 +1,316 @@ | |||||||||||||||||
| use std::ops::Range; | |||||||||||||||||
|
|
|||||||||||||||||
| use datafusion::arrow::array::{Array, MapArray, StringArray}; | |||||||||||||||||
| use datafusion::error::Result; | |||||||||||||||||
| use datafusion_common::exec_err; | |||||||||||||||||
| use sail_common::spec::{SAIL_MAP_KEY_FIELD_NAME, SAIL_MAP_VALUE_FIELD_NAME}; | |||||||||||||||||
|
|
|||||||||||||||||
| /// Which CSV expression function is reading the options. | |||||||||||||||||
| /// | |||||||||||||||||
| /// Spark validates almost every option identically in all three, but not quite: only the writer | |||||||||||||||||
| /// rejects an empty `lineSep`, and only `from_csv` rejects the `DROPMALFORMED` parse mode. | |||||||||||||||||
| #[derive(Clone, Copy, PartialEq, Eq)] | |||||||||||||||||
| pub(super) enum CsvFunction { | |||||||||||||||||
| From, | |||||||||||||||||
| To, | |||||||||||||||||
| SchemaOf, | |||||||||||||||||
| } | |||||||||||||||||
|
|
|||||||||||||||||
| impl CsvFunction { | |||||||||||||||||
| fn name(self) -> &'static str { | |||||||||||||||||
| match self { | |||||||||||||||||
| Self::From => "from_csv", | |||||||||||||||||
| Self::To => "to_csv", | |||||||||||||||||
| Self::SchemaOf => "schema_of_csv", | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
|
|
|||||||||||||||||
| /// The boolean options that Spark's `CSVOptions` parses with `getBool`, which reports | |||||||||||||||||
| /// `<option> flag can be true or false.` on a non-boolean value. | |||||||||||||||||
| const BOOL_OPTIONS: [&str; 9] = [ | |||||||||||||||||
| "header", | |||||||||||||||||
| "inferSchema", | |||||||||||||||||
| "ignoreLeadingWhiteSpace", | |||||||||||||||||
| "ignoreTrailingWhiteSpace", | |||||||||||||||||
| "escapeQuotes", | |||||||||||||||||
| "quoteAll", | |||||||||||||||||
| "enforceSchema", | |||||||||||||||||
| "preferDate", | |||||||||||||||||
| "columnPruning", | |||||||||||||||||
| ]; | |||||||||||||||||
|
|
|||||||||||||||||
| /// Boolean options Spark reads with Scala's `String.toBoolean` instead of `getBool`, so a bad | |||||||||||||||||
| /// value is reported as `For input string: "<value>"`. `multiLine` is one of these; the rest are | |||||||||||||||||
| /// inherited from the file-source options and validated just as eagerly. | |||||||||||||||||
| const TO_BOOLEAN_OPTIONS: [&str; 4] = [ | |||||||||||||||||
| "multiLine", | |||||||||||||||||
| "enableDateTimeParsingFallback", | |||||||||||||||||
| "ignoreCorruptFiles", | |||||||||||||||||
| "ignoreMissingFiles", | |||||||||||||||||
| ]; | |||||||||||||||||
|
|
|||||||||||||||||
| /// Options Spark reads with `CSVOptions.getChar`. An empty value means "no character" and is | |||||||||||||||||
| /// accepted; anything longer than one character is rejected. | |||||||||||||||||
| const SINGLE_CHAR_OPTIONS: [&str; 4] = ["quote", "escape", "comment", "charToEscapeQuoteEscaping"]; | |||||||||||||||||
|
|
|||||||||||||||||
| /// Integer options Spark reads with `getInt`, reporting `<option> should be an integer. Found <v>.`. | |||||||||||||||||
| const INT_OPTIONS: [&str; 2] = ["maxColumns", "maxCharsPerColumn"]; | |||||||||||||||||
| /// Integer options Spark reads with a bare `.toInt`, reporting `For input string: "<value>"`. | |||||||||||||||||
| const TO_INT_OPTIONS: [&str; 1] = ["inputBufferSize"]; | |||||||||||||||||
| const SAMPLING_RATIO_OPTION: &str = "samplingRatio"; | |||||||||||||||||
| const UNESCAPED_QUOTE_HANDLING_OPTION: &str = "unescapedQuoteHandling"; | |||||||||||||||||
| const UNESCAPED_QUOTE_HANDLING_VALUES: [&str; 5] = [ | |||||||||||||||||
| "STOP_AT_CLOSING_QUOTE", | |||||||||||||||||
| "BACK_TO_DELIMITER", | |||||||||||||||||
| "STOP_AT_DELIMITER", | |||||||||||||||||
| "SKIP_VALUE", | |||||||||||||||||
| "RAISE_ERROR", | |||||||||||||||||
| ]; | |||||||||||||||||
| const LINE_SEP_OPTION: &str = "lineSep"; | |||||||||||||||||
| const MODE_OPTION: &str = "mode"; | |||||||||||||||||
| const DROP_MALFORMED_MODE: &str = "DROPMALFORMED"; | |||||||||||||||||
|
|
|||||||||||||||||
| /// Spark validates the charset (`encoding`, and its alias `charset`) against a fixed set in | |||||||||||||||||
| /// `CSVOptions`, reporting `[INVALID_PARAMETER_VALUE.CHARSET]` for anything else. The set is | |||||||||||||||||
| /// narrower than Java's — `Shift_JIS`, for instance, is rejected. | |||||||||||||||||
| const CHARSET_OPTIONS: [&str; 2] = ["encoding", "charset"]; | |||||||||||||||||
| const VALID_CHARSETS: [&str; 7] = [ | |||||||||||||||||
| "iso-8859-1", | |||||||||||||||||
| "us-ascii", | |||||||||||||||||
| "utf-16", | |||||||||||||||||
| "utf-16be", | |||||||||||||||||
| "utf-16le", | |||||||||||||||||
| "utf-32", | |||||||||||||||||
| "utf-8", | |||||||||||||||||
| ]; | |||||||||||||||||
|
|
|||||||||||||||||
| /// Spark validates `compression` against its registered codecs, reporting `[CODEC_NOT_AVAILABLE]` | |||||||||||||||||
| /// for anything else. There is no `gz` alias and `zstd` is not among them. | |||||||||||||||||
| const COMPRESSION_OPTION: &str = "compression"; | |||||||||||||||||
| const VALID_CODECS: [&str; 7] = [ | |||||||||||||||||
| "bzip2", | |||||||||||||||||
| "deflate", | |||||||||||||||||
| "uncompressed", | |||||||||||||||||
| "snappy", | |||||||||||||||||
| "none", | |||||||||||||||||
| "lz4", | |||||||||||||||||
| "gzip", | |||||||||||||||||
| ]; | |||||||||||||||||
|
|
|||||||||||||||||
| /// The key and value columns of the options map, together with the entry range of its first row. | |||||||||||||||||
| /// | |||||||||||||||||
| /// The options argument is one map for the whole batch, but `make_scalar_function` pads a scalar | |||||||||||||||||
| /// argument to one value per row, so the flattened entries repeat that map once per row. Only the | |||||||||||||||||
| /// first row is read: the rest are copies, and walking them would validate the same option once | |||||||||||||||||
| /// per row of the batch. | |||||||||||||||||
| fn first_row_entries(map: &MapArray) -> Option<(&StringArray, &StringArray, Range<usize>)> { | |||||||||||||||||
| if map.is_empty() || map.is_null(0) { | |||||||||||||||||
| return None; | |||||||||||||||||
| } | |||||||||||||||||
| let keys = map | |||||||||||||||||
| .entries() | |||||||||||||||||
| .column_by_name(SAIL_MAP_KEY_FIELD_NAME)? | |||||||||||||||||
| .as_any() | |||||||||||||||||
| .downcast_ref::<StringArray>()?; | |||||||||||||||||
| let values = map | |||||||||||||||||
| .entries() | |||||||||||||||||
| .column_by_name(SAIL_MAP_VALUE_FIELD_NAME)? | |||||||||||||||||
| .as_any() | |||||||||||||||||
| .downcast_ref::<StringArray>()?; | |||||||||||||||||
| let offsets = map.value_offsets(); | |||||||||||||||||
| let start = *offsets.first()? as usize; | |||||||||||||||||
| let end = *offsets.get(1)? as usize; | |||||||||||||||||
| Some((keys, values, start..end)) | |||||||||||||||||
| } | |||||||||||||||||
|
|
|||||||||||||||||
| /// Returns the effective value of the `key` option, or `None` when the map does not carry it. | |||||||||||||||||
| /// | |||||||||||||||||
| /// Spark reads CSV options through a `CaseInsensitiveMap`, so `SEP` selects the same option as | |||||||||||||||||
| /// `sep` and a later case-variant of a key SHADOWS an earlier one — the LAST match wins. A null | |||||||||||||||||
| /// value yields `None` rather than an empty string. | |||||||||||||||||
| pub(super) fn find_option<'a>(map: &'a MapArray, key: &str) -> Option<&'a str> { | |||||||||||||||||
| let (keys, values, entries) = first_row_entries(map)?; | |||||||||||||||||
| keys.iter() | |||||||||||||||||
| .zip(values.iter()) | |||||||||||||||||
| .take(entries.end) | |||||||||||||||||
| .skip(entries.start) | |||||||||||||||||
| .filter_map(|(entry_key, entry_value)| match entry_key { | |||||||||||||||||
| Some(entry_key) if entry_key.eq_ignore_ascii_case(key) => Some(entry_value), | |||||||||||||||||
| _ => None, | |||||||||||||||||
| }) | |||||||||||||||||
| .next_back() | |||||||||||||||||
| .flatten() | |||||||||||||||||
| } | |||||||||||||||||
|
|
|||||||||||||||||
| /// Rejects a structurally invalid options map — a NULL key or value — the way Spark does during | |||||||||||||||||
| /// map conversion, before any input row is evaluated. | |||||||||||||||||
| /// | |||||||||||||||||
| /// This is the ONLY part of option handling that is eager: Spark fails the call for a NULL entry | |||||||||||||||||
| /// even when the input is NULL and even under a key it does not know, because the failure happens | |||||||||||||||||
| /// while building the options, not while parsing a row. Value validation, by contrast, is lazy | |||||||||||||||||
| /// (see [`validate_options`]). | |||||||||||||||||
| pub(super) fn reject_null_entries(map: &MapArray, function: CsvFunction) -> Result<()> { | |||||||||||||||||
| let Some((keys, values, entries)) = first_row_entries(map) else { | |||||||||||||||||
| return Ok(()); | |||||||||||||||||
| }; | |||||||||||||||||
| for (key, value) in keys | |||||||||||||||||
| .iter() | |||||||||||||||||
| .zip(values.iter()) | |||||||||||||||||
| .take(entries.end) | |||||||||||||||||
| .skip(entries.start) | |||||||||||||||||
| { | |||||||||||||||||
| if key.is_none() || value.is_none() { | |||||||||||||||||
| return exec_err!( | |||||||||||||||||
| "Failed preparing of the function `{}` for call. Please, double check function's arguments.", | |||||||||||||||||
| function.name() | |||||||||||||||||
| ); | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| Ok(()) | |||||||||||||||||
| } | |||||||||||||||||
|
|
|||||||||||||||||
| /// Rejects an invalid VALUE for any CSV option present in `map`. | |||||||||||||||||
| /// | |||||||||||||||||
| /// Read-driven, matching Spark: each canonical option is validated by READING its effective value | |||||||||||||||||
| /// from the case-insensitive map (last case-variant wins, via [`find_option`]), NOT by scanning raw | |||||||||||||||||
| /// entries. A shadowed value — an earlier case-variant, or the losing side of an alias such as | |||||||||||||||||
| /// `delimiter` when `sep` is present — is never read and therefore never validated. | |||||||||||||||||
| /// | |||||||||||||||||
| /// Spark builds `CSVOptions` eagerly relative to other options — so `from_csv` rejects a bad | |||||||||||||||||
| /// `quoteAll` and `to_csv` rejects a bad `header` — but LAZILY relative to the input: a bad value is | |||||||||||||||||
| /// not seen when every input row is NULL. Callers therefore invoke this only when the input has at | |||||||||||||||||
| /// least one non-null row; the structural NULL-entry check ([`reject_null_entries`]) runs | |||||||||||||||||
| /// unconditionally instead. Unknown keys are ignored, also matching Spark. | |||||||||||||||||
| pub(super) fn validate_options(map: &MapArray, function: CsvFunction) -> Result<()> { | |||||||||||||||||
| for option in BOOL_OPTIONS { | |||||||||||||||||
| if let Some(value) = find_option(map, option) { | |||||||||||||||||
| parse_bool_option(Some(value), option, false)?; | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| // Read with Scala's `String.toBoolean`, which quotes the value instead of naming the option. | |||||||||||||||||
| for option in TO_BOOLEAN_OPTIONS { | |||||||||||||||||
| if let Some(value) = find_option(map, option) | |||||||||||||||||
| && !fold_eq(value, "true") | |||||||||||||||||
| && !fold_eq(value, "false") | |||||||||||||||||
| { | |||||||||||||||||
| return exec_err!("For input string: \"{value}\""); | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| // Spark measures the value with Java's `String.length`, which counts UTF-16 code units: | |||||||||||||||||
| // an emoji is two units and is rejected, while `ñ` is one and is accepted. | |||||||||||||||||
| for option in SINGLE_CHAR_OPTIONS { | |||||||||||||||||
| if let Some(value) = find_option(map, option) | |||||||||||||||||
| && value.encode_utf16().count() > 1 | |||||||||||||||||
| { | |||||||||||||||||
| return exec_err!("{option} cannot be more than one character."); | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| // Spark reads `sep`, falling back to `delimiter`; only the value it would read is validated. | |||||||||||||||||
| if let Some(value) = find_option(map, "sep").or_else(|| find_option(map, "delimiter")) | |||||||||||||||||
| && value.is_empty() | |||||||||||||||||
| { | |||||||||||||||||
| return exec_err!("Delimiter cannot be empty"); | |||||||||||||||||
| } | |||||||||||||||||
| for option in INT_OPTIONS { | |||||||||||||||||
| if let Some(value) = find_option(map, option) | |||||||||||||||||
| && value.parse::<i32>().is_err() | |||||||||||||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Is it possible to use Java/Scala numeric parsing here?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Good catch — I measured it, and there is a real divergence, though a narrow one:
Whitespace, sign and overflow already match. The only gap is Unicode decimal digits: Java's Added a test for it ( @sail-bug
Scenario: from_csv accepts a Unicode decimal digit in an integer option
SELECT from_csv('1', 'a INT', map('maxColumns', '٥')) AS result
→ | {1} | |
|||||||||||||||||
| { | |||||||||||||||||
| return exec_err!("{option} should be an integer. Found {value}."); | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| // `inputBufferSize` (`.toInt`) and `samplingRatio` (`.toDouble`) both report the raw value. | |||||||||||||||||
| for option in TO_INT_OPTIONS { | |||||||||||||||||
| if let Some(value) = find_option(map, option) | |||||||||||||||||
| && value.parse::<i32>().is_err() | |||||||||||||||||
| { | |||||||||||||||||
| return exec_err!("For input string: \"{value}\""); | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| if let Some(value) = find_option(map, SAMPLING_RATIO_OPTION) | |||||||||||||||||
| && value.parse::<f64>().is_err() | |||||||||||||||||
| { | |||||||||||||||||
| return exec_err!("For input string: \"{value}\""); | |||||||||||||||||
| } | |||||||||||||||||
| if let Some(value) = find_option(map, UNESCAPED_QUOTE_HANDLING_OPTION) | |||||||||||||||||
| && !UNESCAPED_QUOTE_HANDLING_VALUES | |||||||||||||||||
| .iter() | |||||||||||||||||
| .any(|x| fold_eq(value, x)) | |||||||||||||||||
| { | |||||||||||||||||
| // Spark upper-cases the value before the lookup, so it surfaces the failure of the | |||||||||||||||||
| // underlying parser's `valueOf` with the upper-cased name. | |||||||||||||||||
| let value = value.to_uppercase(); | |||||||||||||||||
| return exec_err!( | |||||||||||||||||
| "No enum constant com.univocity.parsers.csv.UnescapedQuoteHandling.{value}" | |||||||||||||||||
| ); | |||||||||||||||||
| } | |||||||||||||||||
| // The writer requires at most one Unicode character, which Spark measures as at most two UTF-16 | |||||||||||||||||
| // code units — so `ab` (two units) is accepted but `abc` (three) is not. The reader ignores it. | |||||||||||||||||
| if function == CsvFunction::To | |||||||||||||||||
| && let Some(value) = find_option(map, LINE_SEP_OPTION) | |||||||||||||||||
| { | |||||||||||||||||
| if value.is_empty() { | |||||||||||||||||
| return exec_err!("requirement failed: '{LINE_SEP_OPTION}' cannot be an empty string."); | |||||||||||||||||
| } | |||||||||||||||||
| if value.encode_utf16().count() > 2 { | |||||||||||||||||
| return exec_err!( | |||||||||||||||||
| "requirement failed: '{LINE_SEP_OPTION}' can contain only 1 character." | |||||||||||||||||
| ); | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| // An unrecognized mode is not an error: Spark warns and falls back to PERMISSIVE. Only | |||||||||||||||||
| // `from_csv` rejects `DROPMALFORMED`. | |||||||||||||||||
| if function == CsvFunction::From | |||||||||||||||||
| && let Some(value) = find_option(map, MODE_OPTION) | |||||||||||||||||
| && fold_eq(value, DROP_MALFORMED_MODE) | |||||||||||||||||
| { | |||||||||||||||||
| return exec_err!( | |||||||||||||||||
| "The function `from_csv` doesn't support the {DROP_MALFORMED_MODE} mode. Acceptable modes are PERMISSIVE and FAILFAST." | |||||||||||||||||
| ); | |||||||||||||||||
| } | |||||||||||||||||
| // Spark's `CSVOptions` validates the charset against a fixed set, matched case-insensitively | |||||||||||||||||
| // (Spark lower-cases with `Locale.ROOT`), and rejects everything else — narrower than Java's. | |||||||||||||||||
| for option in CHARSET_OPTIONS { | |||||||||||||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I think we want to do something like this: I believe that Spark treats |
|||||||||||||||||
| if let Some(value) = find_option(map, option) | |||||||||||||||||
| && !VALID_CHARSETS.iter().any(|c| value.eq_ignore_ascii_case(c)) | |||||||||||||||||
| { | |||||||||||||||||
| return exec_err!( | |||||||||||||||||
| "[INVALID_PARAMETER_VALUE.CHARSET] The value of parameter(s) `charset` in `CSVOptions` is invalid: expects one of the iso-8859-1, us-ascii, utf-16, utf-16be, utf-16le, utf-32, utf-8, but got {value}." | |||||||||||||||||
| ); | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
| // Spark validates `compression` against its registered codecs, case-insensitively. | |||||||||||||||||
| if let Some(value) = find_option(map, COMPRESSION_OPTION) | |||||||||||||||||
| && !VALID_CODECS.iter().any(|c| value.eq_ignore_ascii_case(c)) | |||||||||||||||||
| { | |||||||||||||||||
| return exec_err!( | |||||||||||||||||
| "[CODEC_NOT_AVAILABLE.WITH_AVAILABLE_CODECS_SUGGESTION] The codec {value} is not available. Available codecs are bzip2, deflate, uncompressed, snappy, none, lz4, gzip." | |||||||||||||||||
| ); | |||||||||||||||||
| } | |||||||||||||||||
| Ok(()) | |||||||||||||||||
| } | |||||||||||||||||
|
|
|||||||||||||||||
| /// Case-insensitive comparison matching Java's `String.equalsIgnoreCase`, which folds non-ASCII | |||||||||||||||||
| /// letters too — Spark accepts `multiLine='falſe'` because `ſ` (U+017F) upper-cases to `S`. The | |||||||||||||||||
| /// ASCII fast path keeps the common value allocation-free; only a non-ASCII value pays the fold. | |||||||||||||||||
| fn fold_eq(value: &str, target: &str) -> bool { | |||||||||||||||||
| value.eq_ignore_ascii_case(target) | |||||||||||||||||
| || (!value.is_ascii() && value.to_uppercase() == target.to_uppercase()) | |||||||||||||||||
| } | |||||||||||||||||
|
|
|||||||||||||||||
| /// Parses a boolean CSV option the way Spark's `CSVOptions.getBool` does: it lower-cases with | |||||||||||||||||
| /// `Locale.ROOT` and compares to `true`/`false`. That is ASCII-only folding — `ſ` (U+017F) does | |||||||||||||||||
| /// NOT lower-case to `s`, so `header='falſe'` is rejected, unlike the `toBoolean` options, which | |||||||||||||||||
| /// upper-case-fold and accept it. `eq_ignore_ascii_case`, not [`fold_eq`], matches this. | |||||||||||||||||
| pub(super) fn parse_bool_option( | |||||||||||||||||
| value: Option<&str>, | |||||||||||||||||
| option_name: &str, | |||||||||||||||||
| default: bool, | |||||||||||||||||
| ) -> Result<bool> { | |||||||||||||||||
| match value { | |||||||||||||||||
| None => Ok(default), | |||||||||||||||||
| Some(value) if value.eq_ignore_ascii_case("true") => Ok(true), | |||||||||||||||||
| Some(value) if value.eq_ignore_ascii_case("false") => Ok(false), | |||||||||||||||||
| Some(_) => exec_err!("{option_name} flag can be true or false."), | |||||||||||||||||
| } | |||||||||||||||||
| } | |||||||||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think Spark also validates extension, charset, time zone, compression codec, etc. But this can be followup work.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Agreed —
encoding/charset,timeZoneandcompressionare tracked with@sail-bugscenarios so CI keeps an eye on them, and they're on the followup list.One caveat on
extensionspecifically: I measured it and it looks like a Spark quirk rather than validation we should mirror. It's a file-writer option (added by SPARK-50616 for the CSV DataSource writer), not an expression option. Passing it toto_csvacceptsab/abc/abcdbut throwsINTERNAL_ERRORon a non-letter value likea1— andINTERNAL_ERRORis Spark's "this shouldn't happen" class, not a user-facing validation. So I'd lean toward leaving Sail ignoring it as an unknown option on the expression path, rather than replicating it. Happy to revisit if you'd rather match it.