This is an automated email from the ASF dual-hosted git repository.
HuaHuaY pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow.git
The following commit(s) were added to refs/heads/main by this push:
new 68243c4ba6e GH-50967: [C++] Allow CSV reader to ignore extra columns
in rows with more columns (#51118)
68243c4ba6e is described below
commit 68243c4ba6e3a26751b0c74f8769d80f9b452bd8
Author: Zehua Zou <[email protected]>
AuthorDate: Thu Sep 10 14:28:27 2026 +0800
GH-50967: [C++] Allow CSV reader to ignore extra columns in rows with more
columns (#51118)
### Rationale for this change
Currently, the C++ CSV reader rejects rows with more columns than expected.
We can allow users to ignore the extra values instead of throwing exception.
### What changes are included in this PR?
Add an option `ignore_extra_columns` in CSV `struct ParseOptions` that
ignores extra columns.
### Are these changes tested?
Yes.
### Are there any user-facing changes?
Add an option `ignore_extra_columns` in CSV `struct ParseOptions`.
* GitHub Issue: #50967
Authored-by: Zehua Zou <[email protected]>
Signed-off-by: Zehua Zou <[email protected]>
---
cpp/src/arrow/csv/lexing_internal.h | 26 +++++++
cpp/src/arrow/csv/options.h | 2 +
cpp/src/arrow/csv/parser.cc | 133 +++++++++++++++++++++++-------------
cpp/src/arrow/csv/parser_test.cc | 33 +++++++++
cpp/src/arrow/csv/reader_test.cc | 21 ++++++
cpp/src/arrow/dataset/file_csv.cc | 10 +--
6 files changed, 173 insertions(+), 52 deletions(-)
diff --git a/cpp/src/arrow/csv/lexing_internal.h
b/cpp/src/arrow/csv/lexing_internal.h
index b45af7a370d..d5ca120eb58 100644
--- a/cpp/src/arrow/csv/lexing_internal.h
+++ b/cpp/src/arrow/csv/lexing_internal.h
@@ -20,6 +20,8 @@
#include <cstdint>
#include <cstring>
#include <string_view>
+#include <type_traits>
+#include <utility>
#include "arrow/csv/options.h"
#include "arrow/util/simd.h"
@@ -35,6 +37,30 @@ class SpecializedOptions {
static constexpr bool escaping = Escaping;
};
+/// Convert runtime boolean options into template arguments for a callable.
+template <bool... CompiledBools, typename Fn, typename... Rest>
+decltype(auto) DispatchBool(Fn&& fn, Rest... rest)
+ requires requires {
+ std::forward<Fn>(fn)
+ .template operator()<CompiledBools..., std::is_convertible_v<Rest,
bool>...>();
+ }
+{
+ if constexpr (sizeof...(Rest) == 0) {
+ // All runtime booleans have been appended to the compile-time pack.
+ return std::forward<Fn>(fn).template operator()<CompiledBools...>();
+ } else {
+ // Split off the next runtime boolean, append its value to the
compile-time pack,
+ // and recursively dispatch the remaining booleans.
+ return [&](bool head, auto... tail) -> decltype(auto) {
+ if (head) {
+ return DispatchBool<CompiledBools..., true>(std::forward<Fn>(fn),
tail...);
+ } else {
+ return DispatchBool<CompiledBools..., false>(std::forward<Fn>(fn),
tail...);
+ }
+ }(rest...);
+ }
+}
+
//
// Bulk filters for packed character matching.
// These filters allow checking multiple CSV bytes at once for specific
diff --git a/cpp/src/arrow/csv/options.h b/cpp/src/arrow/csv/options.h
index 5d83f9cb491..41d0b63dacb 100644
--- a/cpp/src/arrow/csv/options.h
+++ b/cpp/src/arrow/csv/options.h
@@ -63,6 +63,8 @@ struct ARROW_EXPORT ParseOptions {
InvalidRowHandler invalid_row_handler;
/// Whether rows with fewer columns than expected are padded with nulls.
bool pad_short_rows = false;
+ /// Whether rows with more columns than expected should ignore the extra
columns.
+ bool ignore_extra_columns = false;
/// Create parsing options with default values
static ParseOptions Defaults();
diff --git a/cpp/src/arrow/csv/parser.cc b/cpp/src/arrow/csv/parser.cc
index 6a6138cec85..20f984ffeb6 100644
--- a/cpp/src/arrow/csv/parser.cc
+++ b/cpp/src/arrow/csv/parser.cc
@@ -59,6 +59,14 @@ Status MismatchingColumns(const InvalidRow& row) {
inline bool IsControlChar(uint8_t c) { return c < ' '; }
+template <bool IgnoreExtraColumns>
+constexpr bool ShouldWrite([[maybe_unused]] bool ignoring_extra_field) {
+ if constexpr (IgnoreExtraColumns) {
+ return !ignoring_extra_field;
+ }
+ return true;
+}
+
// A helper class allocating the buffer for parsed values and writing into it
// without any further resizes, except at the end.
class PresizedDataWriter {
@@ -276,8 +284,8 @@ class BlockParserImpl {
return MismatchingColumns(row);
}
- template <typename SpecializedOptions, bool UseBulkFilter, typename
ValueDescWriter,
- typename DataWriter, typename BulkFilter>
+ template <typename SpecializedOptions, bool UseBulkFilter, bool
IgnoreExtraColumns,
+ typename ValueDescWriter, typename DataWriter, typename BulkFilter>
Status ParseLine(ValueDescWriter* values_writer, DataWriter* parsed_writer,
const char* data, const char* data_end, bool is_final,
const char** out_data, const BulkFilter& bulk_filter) {
@@ -287,7 +295,28 @@ class BlockParserImpl {
DCHECK_GT(data_end, data);
- auto FinishField = [&]() { values_writer->FinishField(parsed_writer); };
+ bool ignoring_extra_field = false;
+
+ auto IsExtraField = [&]() {
+ if constexpr (!IgnoreExtraColumns) {
+ return false;
+ }
+ return batch_.num_cols_ >= 0 && num_cols >= batch_.num_cols_;
+ };
+ auto StartField = [&](bool quoted) {
+ if
(ARROW_PREDICT_TRUE(ShouldWrite<IgnoreExtraColumns>(ignoring_extra_field))) {
+ if (ARROW_PREDICT_FALSE(IsExtraField())) {
+ ignoring_extra_field = true;
+ } else {
+ values_writer->StartField(quoted);
+ }
+ }
+ };
+ auto FinishField = [&]() {
+ if
(ARROW_PREDICT_TRUE(ShouldWrite<IgnoreExtraColumns>(ignoring_extra_field))) {
+ values_writer->FinishField(parsed_writer);
+ }
+ };
values_writer->BeginLine();
parsed_writer->BeginLine();
@@ -314,7 +343,7 @@ class BlockParserImpl {
// At the start of a field
if (*data == options_.delimiter) {
// Empty cells are very common in some files, shortcut them
- values_writer->StartField(false /* quoted */);
+ StartField(false /* quoted */);
FinishField();
++data;
++num_cols;
@@ -328,17 +357,18 @@ class BlockParserImpl {
if (SpecializedOptions::quoting &&
ARROW_PREDICT_FALSE(*data == options_.quote_char)) {
++data;
- values_writer->StartField(true /* quoted */);
+ StartField(true /* quoted */);
goto InQuotedField;
} else {
- values_writer->StartField(false /* quoted */);
+ StartField(false /* quoted */);
goto InField;
}
InField:
// Inside a non-quoted part of a field
if (UseBulkFilter) {
- const char* bulk_end = RunBulkFilter(parsed_writer, data, data_end,
bulk_filter);
+ const char* bulk_end = RunBulkFilter<IgnoreExtraColumns>(
+ parsed_writer, data, data_end, bulk_filter, ignoring_extra_field);
if (ARROW_PREDICT_FALSE(bulk_end == nullptr)) {
if (is_final) {
data = data_end;
@@ -358,7 +388,9 @@ class BlockParserImpl {
goto AbortLine;
}
c = *data++;
- parsed_writer->PushFieldChar(c);
+ if
(ARROW_PREDICT_TRUE(ShouldWrite<IgnoreExtraColumns>(ignoring_extra_field))) {
+ parsed_writer->PushFieldChar(c);
+ }
goto InField;
}
if (ARROW_PREDICT_FALSE(c == options_.delimiter)) {
@@ -376,13 +408,16 @@ class BlockParserImpl {
goto LineEnd;
}
}
- parsed_writer->PushFieldChar(c);
+ if
(ARROW_PREDICT_TRUE(ShouldWrite<IgnoreExtraColumns>(ignoring_extra_field))) {
+ parsed_writer->PushFieldChar(c);
+ }
goto InField;
InQuotedField:
// Inside a quoted part of a field
if (UseBulkFilter) {
- const char* bulk_end = RunBulkFilter(parsed_writer, data, data_end,
bulk_filter);
+ const char* bulk_end = RunBulkFilter<IgnoreExtraColumns>(
+ parsed_writer, data, data_end, bulk_filter, ignoring_extra_field);
if (ARROW_PREDICT_FALSE(bulk_end == nullptr)) {
if (is_final) {
data = data_end;
@@ -401,7 +436,9 @@ class BlockParserImpl {
goto AbortLine;
}
c = *data++;
- parsed_writer->PushFieldChar(c);
+ if
(ARROW_PREDICT_TRUE(ShouldWrite<IgnoreExtraColumns>(ignoring_extra_field))) {
+ parsed_writer->PushFieldChar(c);
+ }
goto InQuotedField;
}
if (ARROW_PREDICT_FALSE(c == options_.quote_char)) {
@@ -414,7 +451,9 @@ class BlockParserImpl {
goto InField;
}
}
- parsed_writer->PushFieldChar(c);
+ if
(ARROW_PREDICT_TRUE(ShouldWrite<IgnoreExtraColumns>(ignoring_extra_field))) {
+ parsed_writer->PushFieldChar(c);
+ }
goto InQuotedField;
FieldEnd:
@@ -436,11 +475,11 @@ class BlockParserImpl {
} else if (options_.pad_short_rows && num_cols < batch_.num_cols_) {
batch_.missing_fields_.push_back({batch_.num_rows_, num_cols});
while (num_cols < batch_.num_cols_) {
- values_writer->StartField(false /* quoted */);
+ StartField(false /* quoted */);
FinishField();
++num_cols;
}
- } else {
+ } else if (!IgnoreExtraColumns || num_cols < batch_.num_cols_) {
return HandleInvalidRow(values_writer, parsed_writer, start, data,
num_cols,
out_data);
}
@@ -452,6 +491,10 @@ class BlockParserImpl {
AbortLine:
// Not a full line except perhaps if in final block
if (is_final) {
+ if constexpr (IgnoreExtraColumns) {
+ // Handle an implicit trailing empty field after a delimiter.
+ ignoring_extra_field = IsExtraField();
+ }
goto LineEnd;
}
// Truncated line at end of block, rewind parsed state
@@ -466,9 +509,10 @@ class BlockParserImpl {
batch_.num_cols_ = 1;
}
// Record as row of empty (null?) values
- while (num_cols++ < batch_.num_cols_) {
- values_writer->StartField(false /* quoted */);
+ while (num_cols < batch_.num_cols_) {
+ StartField(false /* quoted */);
FinishField();
+ ++num_cols;
}
++batch_.num_rows_;
}
@@ -476,10 +520,11 @@ class BlockParserImpl {
return Status::OK();
}
- template <typename DataWriter, typename SpecializedBulkFilter>
+ template <bool IgnoreExtraColumns, typename DataWriter, typename
SpecializedBulkFilter>
const char* RunBulkFilter(DataWriter* data_writer, const char* data,
const char* data_end,
- const SpecializedBulkFilter& bulk_filter) {
+ const SpecializedBulkFilter& bulk_filter,
+ bool ignoring_extra_field) {
while (true) {
using WordType = typename SpecializedBulkFilter::WordType;
@@ -495,13 +540,15 @@ class BlockParserImpl {
return data;
}
// No special chars
- data_writer->PushFieldWord(word);
+ if
(ARROW_PREDICT_TRUE(ShouldWrite<IgnoreExtraColumns>(ignoring_extra_field))) {
+ data_writer->PushFieldWord(word);
+ }
data += sizeof(WordType);
}
}
- template <typename SpecializedOptions, typename ValueDescWriter, typename
DataWriter,
- typename BulkFilter>
+ template <typename SpecializedOptions, bool IgnoreExtraColumns,
+ typename ValueDescWriter, typename DataWriter, typename BulkFilter>
Status ParseChunk(ValueDescWriter* values_writer, DataWriter* parsed_writer,
const char* data, const char* data_end, bool is_final,
int32_t rows_in_chunk, const char** out_data, bool*
finished_parsing,
@@ -512,9 +559,9 @@ class BlockParserImpl {
if (use_bulk_filter_) {
while (data < data_end && batch_.num_rows_ < num_rows_deadline) {
const char* line_end = data;
- RETURN_NOT_OK((ParseLine<SpecializedOptions, true>(values_writer,
parsed_writer,
- data, data_end,
is_final,
- &line_end,
bulk_filter)));
+ RETURN_NOT_OK((ParseLine<SpecializedOptions, true, IgnoreExtraColumns>(
+ values_writer, parsed_writer, data, data_end, is_final, &line_end,
+ bulk_filter)));
RETURN_NOT_OK(values_writer->status());
if (line_end == data) {
// Cannot parse any further
@@ -526,9 +573,9 @@ class BlockParserImpl {
} else {
while (data < data_end && batch_.num_rows_ < num_rows_deadline) {
const char* line_end = data;
- RETURN_NOT_OK((ParseLine<SpecializedOptions, false>(values_writer,
parsed_writer,
- data, data_end,
is_final,
- &line_end,
bulk_filter)));
+ RETURN_NOT_OK((ParseLine<SpecializedOptions, false,
IgnoreExtraColumns>(
+ values_writer, parsed_writer, data, data_end, is_final, &line_end,
+ bulk_filter)));
RETURN_NOT_OK(values_writer->status());
if (line_end == data) {
// Cannot parse any further
@@ -559,7 +606,7 @@ class BlockParserImpl {
return Status::OK();
}
- template <typename SpecializedOptions>
+ template <typename SpecializedOptions, bool IgnoreExtraColumns>
Status ParseSpecialized(const std::vector<std::string_view>& views, bool
is_final,
uint32_t* out_size) {
internal::PreferredBulkFilterType<SpecializedOptions>
bulk_filter(options_);
@@ -598,9 +645,9 @@ class BlockParserImpl {
ARROW_ASSIGN_OR_RAISE(auto values_writer,
ResizableValueDescWriter::Make(pool_));
values_writer.Start(parsed_writer);
- RETURN_NOT_OK(ParseChunk<SpecializedOptions>(
+ RETURN_NOT_OK((ParseChunk<SpecializedOptions, IgnoreExtraColumns>(
&values_writer, &parsed_writer, data, data_end, is_final,
rows_in_chunk,
- &data, &finished_parsing, bulk_filter));
+ &data, &finished_parsing, bulk_filter)));
if (batch_.num_cols_ == -1) {
return ParseError("Empty CSV file or block: cannot infer number of
columns");
}
@@ -636,9 +683,9 @@ class BlockParserImpl {
PresizedValueDescWriter::Make(pool_, rows_in_chunk,
batch_.num_cols_));
values_writer.Start(parsed_writer);
- RETURN_NOT_OK(ParseChunk<SpecializedOptions>(
+ RETURN_NOT_OK((ParseChunk<SpecializedOptions, IgnoreExtraColumns>(
&values_writer, &parsed_writer, data, data_end, is_final,
rows_in_chunk,
- &data, &finished_parsing, bulk_filter));
+ &data, &finished_parsing, bulk_filter)));
}
DCHECK_GE(data, view.data());
DCHECK_LE(data, data_end);
@@ -679,23 +726,13 @@ class BlockParserImpl {
Status Parse(const std::vector<std::string_view>& data, bool is_final,
uint32_t* out_size) {
- if (options_.quoting) {
- if (options_.escaping) {
- return ParseSpecialized<internal::SpecializedOptions<true,
true>>(data, is_final,
+ return internal::DispatchBool(
+ [&]<bool Quoting, bool Escaping, bool IgnoreExtraColumns>() {
+ using SpecializedOptions = internal::SpecializedOptions<Quoting,
Escaping>;
+ return ParseSpecialized<SpecializedOptions,
IgnoreExtraColumns>(data, is_final,
out_size);
- } else {
- return ParseSpecialized<internal::SpecializedOptions<true,
false>>(data, is_final,
-
out_size);
- }
- } else {
- if (options_.escaping) {
- return ParseSpecialized<internal::SpecializedOptions<false,
true>>(data, is_final,
-
out_size);
- } else {
- return ParseSpecialized<internal::SpecializedOptions<false, false>>(
- data, is_final, out_size);
- }
- }
+ },
+ options_.quoting, options_.escaping, options_.ignore_extra_columns);
}
protected:
diff --git a/cpp/src/arrow/csv/parser_test.cc b/cpp/src/arrow/csv/parser_test.cc
index c03f492ed27..bd09f991183 100644
--- a/cpp/src/arrow/csv/parser_test.cc
+++ b/cpp/src/arrow/csv/parser_test.cc
@@ -298,6 +298,39 @@ TEST(BlockParser, PadShortRows) {
ASSERT_EQ(last_row_missing, std::vector<bool>({false, false, true}));
}
+TEST(BlockParser, IgnoreExtraColumns) {
+ auto options = ParseOptions::Defaults();
+ options.ignore_extra_columns = true;
+
+ BlockParser parser(options, /*num_cols=*/2);
+ AssertParseOk(parser, "a,\"b\",c,\nd,e\n");
+ AssertColumnsEq(parser, {{"a", "d"}, {"b", "e"}}, {{false, false}, {true,
false}});
+
+ BlockParser final_parser(options, /*num_cols=*/2);
+ AssertParseFinal(final_parser, "a,b,");
+ AssertColumnsEq(final_parser, {{"a"}, {"b"}});
+}
+
+TEST(BlockParser, PadAndIgnore) {
+ auto options = ParseOptions::Defaults();
+ options.pad_short_rows = true;
+ options.ignore_extra_columns = true;
+
+ BlockParser parser(options, /*num_cols=*/2);
+ AssertParseFinal(parser, "a,b,c\nd");
+ AssertColumnEq(parser, 0, {"a", "d"});
+ std::vector<std::string> values;
+ std::vector<bool> missing;
+ ASSERT_OK(parser.VisitColumn(
+ 1, [&](const uint8_t* data, uint32_t size, bool, bool is_missing) ->
Status {
+ values.emplace_back(reinterpret_cast<const char*>(data), size);
+ missing.push_back(is_missing);
+ return Status::OK();
+ }));
+ ASSERT_EQ(values, std::vector<std::string>({"b", ""}));
+ ASSERT_EQ(missing, std::vector<bool>({false, true}));
+}
+
TEST(BlockParser, EmptyHeader) {
// Cannot infer number of columns
uint32_t out_size;
diff --git a/cpp/src/arrow/csv/reader_test.cc b/cpp/src/arrow/csv/reader_test.cc
index 2493cb66271..009bbd5fc25 100644
--- a/cpp/src/arrow/csv/reader_test.cc
+++ b/cpp/src/arrow/csv/reader_test.cc
@@ -644,6 +644,27 @@ TEST(ReaderTests, ShortRows) {
ASSERT_TRUE(table->Equals(*expected_table));
}
+TEST(ReaderTests, IgnoreExtraColumns) {
+ auto input =
+
std::make_shared<io::BufferReader>(std::make_shared<Buffer>("a,b\n1,2,3\n4,5\n"));
+ auto parse_options = ParseOptions::Defaults();
+ parse_options.ignore_extra_columns = true;
+ auto convert_options = ConvertOptions::Defaults();
+ convert_options.default_column_type = int64();
+
+ ASSERT_OK_AND_ASSIGN(auto reader,
TableReader::Make(io::default_io_context(), input,
+ ReadOptions::Defaults(),
+ parse_options,
convert_options));
+ ASSERT_OK_AND_ASSIGN(auto table, reader->Read());
+
+ auto expected_schema = schema({field("a", int64()), field("b", int64())});
+ auto expected_table = TableFromJSON(expected_schema, {R"([
+ {"a":1, "b":2},
+ {"a":4, "b":5}
+ ])"});
+ ASSERT_TRUE(table->Equals(*expected_table));
+}
+
TEST(ReaderTests, ShortRowsTypedConverters) {
auto input =
std::make_shared<io::BufferReader>(std::make_shared<Buffer>("1,10\n2\n"));
auto read_options = ReadOptions::Defaults();
diff --git a/cpp/src/arrow/dataset/file_csv.cc
b/cpp/src/arrow/dataset/file_csv.cc
index 079103fa791..3b9e8d6ca20 100644
--- a/cpp/src/arrow/dataset/file_csv.cc
+++ b/cpp/src/arrow/dataset/file_csv.cc
@@ -164,11 +164,12 @@ Result<std::vector<std::string>> GetOrderedColumnNames(
int32_t max_num_rows = read_options.skip_rows + 1;
std::optional<csv::ParseOptions> inspection_parse_options;
const auto* parser_options = &parse_options;
- if (parse_options.pad_short_rows) {
- // Do not pad short rows while determining column names, since padding
cannot
- // synthesize missing names. Copy the parse options only when needed.
+ if (parse_options.pad_short_rows || parse_options.ignore_extra_columns) {
+ // Do not adjust row widths while determining column names: padding cannot
+ // synthesize missing names, and ignoring extra columns may discard
columns.
inspection_parse_options.emplace(parse_options);
inspection_parse_options->pad_short_rows = false;
+ inspection_parse_options->ignore_extra_columns = false;
parser_options = &*inspection_parse_options;
}
csv::BlockParser parser(pool, *parser_options, /*num_cols=*/-1,
/*first_row=*/1,
@@ -379,7 +380,8 @@ bool CsvFileFormat::Equals(const FileFormat& format) const {
parse_options.escape_char == other_parse_options.escape_char &&
parse_options.newlines_in_values ==
other_parse_options.newlines_in_values &&
parse_options.ignore_empty_lines ==
other_parse_options.ignore_empty_lines &&
- parse_options.pad_short_rows == other_parse_options.pad_short_rows;
+ parse_options.pad_short_rows == other_parse_options.pad_short_rows &&
+ parse_options.ignore_extra_columns ==
other_parse_options.ignore_extra_columns;
}
Result<bool> CsvFileFormat::IsSupported(const FileSource& source) const {