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 {

Reply via email to