From 765ffcb08d9b11c254bbaced0486e411f05ef856 Mon Sep 17 00:00:00 2001 From: Zehua Zou Date: Tue, 1 Sep 2026 14:31:21 +0800 Subject: [PATCH 1/2] allow CSV reader to ignore extra columns --- cpp/src/arrow/csv/options.h | 2 + cpp/src/arrow/csv/parser.cc | 62 +++++++++++++++++++++++-------- cpp/src/arrow/csv/parser_test.cc | 13 +++++++ cpp/src/arrow/csv/reader_test.cc | 21 +++++++++++ cpp/src/arrow/dataset/file_csv.cc | 10 +++-- 5 files changed, 88 insertions(+), 20 deletions(-) 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..4dc3415875f 100644 --- a/cpp/src/arrow/csv/parser.cc +++ b/cpp/src/arrow/csv/parser.cc @@ -287,7 +287,23 @@ class BlockParserImpl { DCHECK_GT(data_end, data); - auto FinishField = [&]() { values_writer->FinishField(parsed_writer); }; + bool ignoring_extra_field = false; + + auto IsExtraField = [&]() { + return batch_.num_cols_ >= 0 && options_.ignore_extra_columns && + num_cols >= batch_.num_cols_; + }; + auto StartField = [&](bool quoted) { + ignoring_extra_field = IsExtraField(); + if (!ignoring_extra_field) { + values_writer->StartField(quoted); + } + }; + auto FinishField = [&]() { + if (!IsExtraField()) { + values_writer->FinishField(parsed_writer); + } + }; values_writer->BeginLine(); parsed_writer->BeginLine(); @@ -314,7 +330,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 +344,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(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 +375,9 @@ class BlockParserImpl { goto AbortLine; } c = *data++; - parsed_writer->PushFieldChar(c); + if (!ignoring_extra_field) { + parsed_writer->PushFieldChar(c); + } goto InField; } if (ARROW_PREDICT_FALSE(c == options_.delimiter)) { @@ -376,13 +395,16 @@ class BlockParserImpl { goto LineEnd; } } - parsed_writer->PushFieldChar(c); + if (!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(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 +423,9 @@ class BlockParserImpl { goto AbortLine; } c = *data++; - parsed_writer->PushFieldChar(c); + if (!ignoring_extra_field) { + parsed_writer->PushFieldChar(c); + } goto InQuotedField; } if (ARROW_PREDICT_FALSE(c == options_.quote_char)) { @@ -414,7 +438,9 @@ class BlockParserImpl { goto InField; } } - parsed_writer->PushFieldChar(c); + if (!ignoring_extra_field) { + parsed_writer->PushFieldChar(c); + } goto InQuotedField; FieldEnd: @@ -436,11 +462,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 (!options_.ignore_extra_columns || num_cols < batch_.num_cols_) { return HandleInvalidRow(values_writer, parsed_writer, start, data, num_cols, out_data); } @@ -466,9 +492,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_; } @@ -479,7 +506,8 @@ class BlockParserImpl { template 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,7 +523,9 @@ class BlockParserImpl { return data; } // No special chars - data_writer->PushFieldWord(word); + if (!ignoring_extra_field) { + data_writer->PushFieldWord(word); + } data += sizeof(WordType); } } diff --git a/cpp/src/arrow/csv/parser_test.cc b/cpp/src/arrow/csv/parser_test.cc index c03f492ed27..a15e8ed04e3 100644 --- a/cpp/src/arrow/csv/parser_test.cc +++ b/cpp/src/arrow/csv/parser_test.cc @@ -298,6 +298,19 @@ TEST(BlockParser, PadShortRows) { ASSERT_EQ(last_row_missing, std::vector({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, 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(std::make_shared("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(std::make_shared("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> GetOrderedColumnNames( int32_t max_num_rows = read_options.skip_rows + 1; std::optional 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 CsvFileFormat::IsSupported(const FileSource& source) const { From 0f0c8c323ee03c8f64dc3643c872ff287cfcf9a5 Mon Sep 17 00:00:00 2001 From: Zehua Zou Date: Fri, 4 Sep 2026 16:03:19 +0800 Subject: [PATCH 2/2] minor performance optimizations --- cpp/src/arrow/csv/parser.cc | 25 +++++++++++++++---------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/cpp/src/arrow/csv/parser.cc b/cpp/src/arrow/csv/parser.cc index 4dc3415875f..8283718acc1 100644 --- a/cpp/src/arrow/csv/parser.cc +++ b/cpp/src/arrow/csv/parser.cc @@ -290,17 +290,20 @@ class BlockParserImpl { bool ignoring_extra_field = false; auto IsExtraField = [&]() { - return batch_.num_cols_ >= 0 && options_.ignore_extra_columns && + return options_.ignore_extra_columns && batch_.num_cols_ >= 0 && num_cols >= batch_.num_cols_; }; auto StartField = [&](bool quoted) { - ignoring_extra_field = IsExtraField(); - if (!ignoring_extra_field) { - values_writer->StartField(quoted); + if (ARROW_PREDICT_TRUE(!ignoring_extra_field)) { + if (ARROW_PREDICT_FALSE(IsExtraField())) { + ignoring_extra_field = true; + } else { + values_writer->StartField(quoted); + } } }; auto FinishField = [&]() { - if (!IsExtraField()) { + if (ARROW_PREDICT_TRUE(!ignoring_extra_field)) { values_writer->FinishField(parsed_writer); } }; @@ -375,7 +378,7 @@ class BlockParserImpl { goto AbortLine; } c = *data++; - if (!ignoring_extra_field) { + if (ARROW_PREDICT_TRUE(!ignoring_extra_field)) { parsed_writer->PushFieldChar(c); } goto InField; @@ -395,7 +398,7 @@ class BlockParserImpl { goto LineEnd; } } - if (!ignoring_extra_field) { + if (ARROW_PREDICT_TRUE(!ignoring_extra_field)) { parsed_writer->PushFieldChar(c); } goto InField; @@ -423,7 +426,7 @@ class BlockParserImpl { goto AbortLine; } c = *data++; - if (!ignoring_extra_field) { + if (ARROW_PREDICT_TRUE(!ignoring_extra_field)) { parsed_writer->PushFieldChar(c); } goto InQuotedField; @@ -438,7 +441,7 @@ class BlockParserImpl { goto InField; } } - if (!ignoring_extra_field) { + if (ARROW_PREDICT_TRUE(!ignoring_extra_field)) { parsed_writer->PushFieldChar(c); } goto InQuotedField; @@ -478,6 +481,8 @@ class BlockParserImpl { AbortLine: // Not a full line except perhaps if in final block if (is_final) { + // Handle an implicit trailing empty field after a delimiter. + ignoring_extra_field = IsExtraField(); goto LineEnd; } // Truncated line at end of block, rewind parsed state @@ -523,7 +528,7 @@ class BlockParserImpl { return data; } // No special chars - if (!ignoring_extra_field) { + if (ARROW_PREDICT_TRUE(!ignoring_extra_field)) { data_writer->PushFieldWord(word); } data += sizeof(WordType);