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 ff30c725d2 GH-50925: [C++] Allow CSV reader to pad rows with missing
trailing fields (#50926)
ff30c725d2 is described below
commit ff30c725d28f1e6299a225064f8e81f0941a8b0d
Author: Zehua Zou <[email protected]>
AuthorDate: Tue Sep 1 10:20:51 2026 +0800
GH-50925: [C++] Allow CSV reader to pad rows with missing trailing fields
(#50926)
### Rationale for this change
Some CSV files omit trailing optional fields. The C++ CSV reader currently
rejects these short rows, requiring callers to preprocess or discard them.
### What changes are included in this PR?
Add an option `pad_short_rows` in CSV `struct ParseOptions` that pads
missing trailing fields with nulls.
### Are these changes tested?
Yes.
### Are there any user-facing changes?
Add an option `pad_short_rows` in CSV `struct ParseOptions`.
* GitHub Issue: #50925
Authored-by: Zehua Zou <[email protected]>
Signed-off-by: Zehua Zou <[email protected]>
---
cpp/src/arrow/csv/converter.cc | 15 ++++---
cpp/src/arrow/csv/options.h | 2 +
cpp/src/arrow/csv/parser.cc | 7 ++++
cpp/src/arrow/csv/parser.h | 27 ++++++++++--
cpp/src/arrow/csv/parser_benchmark.cc | 2 +-
cpp/src/arrow/csv/parser_test.cc | 78 +++++++++++++++++++++++++++--------
cpp/src/arrow/csv/reader.cc | 4 +-
cpp/src/arrow/csv/reader_test.cc | 46 +++++++++++++++++++++
cpp/src/arrow/dataset/file_csv.cc | 20 +++++++--
9 files changed, 168 insertions(+), 33 deletions(-)
diff --git a/cpp/src/arrow/csv/converter.cc b/cpp/src/arrow/csv/converter.cc
index bb59d02cd2..0b7b9f5c8a 100644
--- a/cpp/src/arrow/csv/converter.cc
+++ b/cpp/src/arrow/csv/converter.cc
@@ -534,8 +534,9 @@ class NullConverter : public ConcreteConverter {
int32_t col_index) override {
NullBuilder builder(pool_);
- auto visit = [&](const uint8_t* data, uint32_t size, bool quoted) ->
Status {
- if (ARROW_PREDICT_TRUE(decoder_.IsNull(data, size, quoted))) {
+ auto visit = [&](const uint8_t* data, uint32_t size, bool quoted,
+ bool missing) -> Status {
+ if (missing || ARROW_PREDICT_TRUE(decoder_.IsNull(data, size, quoted))) {
return builder.AppendNull();
} else {
return GenericConversionError(type_, data, size);
@@ -573,8 +574,9 @@ class PrimitiveConverter : public ConcreteConverter {
BuilderType builder(type_, pool_);
RETURN_NOT_OK(PresizeBuilder(parser, &builder));
- auto visit = [&](const uint8_t* data, uint32_t size, bool quoted) ->
Status {
- if (decoder_.IsNull(data, size, quoted /* quoted */)) {
+ auto visit = [&](const uint8_t* data, uint32_t size, bool quoted,
+ bool missing) -> Status {
+ if (missing || decoder_.IsNull(data, size, quoted /* quoted */)) {
return builder.AppendNull();
}
value_type value{};
@@ -617,8 +619,9 @@ class TypedDictionaryConverter : public
ConcreteDictionaryConverter {
BuilderType builder(value_type_, pool_);
RETURN_NOT_OK(PresizeBuilder(parser, &builder));
- auto visit = [&](const uint8_t* data, uint32_t size, bool quoted) ->
Status {
- if (decoder_.IsNull(data, size, quoted /* quoted */)) {
+ auto visit = [&](const uint8_t* data, uint32_t size, bool quoted,
+ bool missing) -> Status {
+ if (missing || decoder_.IsNull(data, size, quoted /* quoted */)) {
return builder.AppendNull();
}
if (ARROW_PREDICT_FALSE(builder.dictionary_length() > max_cardinality_))
{
diff --git a/cpp/src/arrow/csv/options.h b/cpp/src/arrow/csv/options.h
index f0b923d0f3..5d83f9cb49 100644
--- a/cpp/src/arrow/csv/options.h
+++ b/cpp/src/arrow/csv/options.h
@@ -61,6 +61,8 @@ struct ARROW_EXPORT ParseOptions {
bool ignore_empty_lines = true;
/// A handler function for rows which do not have the correct number of
columns
InvalidRowHandler invalid_row_handler;
+ /// Whether rows with fewer columns than expected are padded with nulls.
+ bool pad_short_rows = 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 bed94d1970..6a6138cec8 100644
--- a/cpp/src/arrow/csv/parser.cc
+++ b/cpp/src/arrow/csv/parser.cc
@@ -433,6 +433,13 @@ class BlockParserImpl {
if (ARROW_PREDICT_FALSE(num_cols != batch_.num_cols_)) {
if (batch_.num_cols_ == -1) {
batch_.num_cols_ = num_cols;
+ } 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 */);
+ FinishField();
+ ++num_cols;
+ }
} else {
return HandleInvalidRow(values_writer, parsed_writer, start, data,
num_cols,
out_data);
diff --git a/cpp/src/arrow/csv/parser.h b/cpp/src/arrow/csv/parser.h
index c73e52ce83..3632ab9c75 100644
--- a/cpp/src/arrow/csv/parser.h
+++ b/cpp/src/arrow/csv/parser.h
@@ -22,6 +22,7 @@
#include <cstdint>
#include <memory>
#include <string_view>
+#include <utility>
#include <vector>
#include "arrow/buffer.h"
@@ -71,6 +72,7 @@ class ARROW_EXPORT DataBatch {
using detail::ParsedValueDesc;
int32_t batch_row = 0;
+ size_t missing_index = 0;
for (size_t buf_index = 0; buf_index < values_buffers_.size();
++buf_index) {
const auto& values_buffer = values_buffers_[buf_index];
const auto values = reinterpret_cast<const
ParsedValueDesc*>(values_buffer->data());
@@ -80,7 +82,16 @@ class ARROW_EXPORT DataBatch {
auto start = values[pos].offset;
auto stop = values[pos + 1].offset;
auto quoted = values[pos + 1].quoted;
- Status status = visit(parsed_ + start, stop - start, quoted);
+ const bool row_has_missing_fields =
+ missing_index < missing_fields_.size() &&
+ missing_fields_[missing_index].row == batch_row;
+ const bool missing =
+ row_has_missing_fields &&
+ col_index >= missing_fields_[missing_index].first_missing_column;
+ if (row_has_missing_fields) {
+ ++missing_index;
+ }
+ Status status = visit(parsed_ + start, stop - start, quoted, missing);
if (ARROW_PREDICT_FALSE(!status.ok())) {
return DecorateWithRowNumber(std::move(status), first_row,
batch_row);
}
@@ -98,11 +109,15 @@ class ARROW_EXPORT DataBatch {
const auto start_pos =
static_cast<int32_t>(values_buffer->size() / sizeof(ParsedValueDesc)) -
num_cols_ - 1;
+ const bool last_row_has_missing_fields =
+ !missing_fields_.empty() && missing_fields_.back().row == num_rows_ -
1;
for (int32_t col_index = 0; col_index < num_cols_; ++col_index) {
auto start = values[start_pos + col_index].offset;
auto stop = values[start_pos + col_index + 1].offset;
auto quoted = values[start_pos + col_index + 1].quoted;
- ARROW_RETURN_NOT_OK(visit(parsed_ + start, stop - start, quoted));
+ const bool missing = last_row_has_missing_fields &&
+ col_index >=
missing_fields_.back().first_missing_column;
+ ARROW_RETURN_NOT_OK(visit(parsed_ + start, stop - start, quoted,
missing));
}
return Status::OK();
}
@@ -138,6 +153,12 @@ class ARROW_EXPORT DataBatch {
// Record the current num_rows_ each time a row is skipped
std::vector<int32_t> skipped_rows_;
+ // Record the first missing column for rows padded with nulls
+ struct MissingFieldRange {
+ int32_t row;
+ int32_t first_missing_column;
+ };
+ std::vector<MissingFieldRange> missing_fields_;
friend class ::arrow::csv::BlockParserImpl;
};
@@ -206,7 +227,7 @@ class ARROW_EXPORT BlockParser {
/// \brief Visit parsed values in a column
///
/// The signature of the visitor is
- /// Status(const uint8_t* data, uint32_t size, bool quoted)
+ /// Status(const uint8_t* data, uint32_t size, bool quoted, bool missing)
template <typename Visitor>
Status VisitColumn(int32_t col_index, Visitor&& visit) const {
return parsed_batch().VisitColumn(col_index, first_row_num(),
diff --git a/cpp/src/arrow/csv/parser_benchmark.cc
b/cpp/src/arrow/csv/parser_benchmark.cc
index fa2d5f76d4..afe82025f7 100644
--- a/cpp/src/arrow/csv/parser_benchmark.cc
+++ b/cpp/src/arrow/csv/parser_benchmark.cc
@@ -168,7 +168,7 @@ static void BenchmarkCSVParsing(benchmark::State& state,
// NOLINT non-const re
// vary depending on the parser's internal data structures.
bool dummy_quoted = false;
uint32_t dummy_size = 0;
- auto visit = [&](const uint8_t* data, uint32_t size, bool quoted) {
+ auto visit = [&](const uint8_t* data, uint32_t size, bool quoted, bool
missing) {
dummy_size += size;
dummy_quoted ^= quoted;
return Status::OK();
diff --git a/cpp/src/arrow/csv/parser_test.cc b/cpp/src/arrow/csv/parser_test.cc
index 719f13d65f..c03f492ed2 100644
--- a/cpp/src/arrow/csv/parser_test.cc
+++ b/cpp/src/arrow/csv/parser_test.cc
@@ -91,7 +91,9 @@ void GetColumn(const BlockParser& parser, int32_t col_index,
std::vector<std::string>* out, std::vector<bool>* out_quoted =
nullptr) {
std::vector<std::string> values;
std::vector<bool> quoted_values;
- auto visit = [&](const uint8_t* data, uint32_t size, bool quoted) -> Status {
+ auto visit = [&](const uint8_t* data, uint32_t size, bool quoted,
+ bool missing) -> Status {
+ EXPECT_FALSE(missing);
values.push_back(std::string(reinterpret_cast<const char*>(data), size));
if (out_quoted) {
quoted_values.push_back(quoted);
@@ -109,7 +111,9 @@ void GetLastRow(const BlockParser& parser,
std::vector<std::string>* out,
std::vector<bool>* out_quoted = nullptr) {
std::vector<std::string> values;
std::vector<bool> quoted_values;
- auto visit = [&](const uint8_t* data, uint32_t size, bool quoted) -> Status {
+ auto visit = [&](const uint8_t* data, uint32_t size, bool quoted,
+ bool missing) -> Status {
+ EXPECT_FALSE(missing);
values.push_back(std::string(reinterpret_cast<const char*>(data), size));
if (out_quoted) {
quoted_values.push_back(quoted);
@@ -264,6 +268,36 @@ TEST(BlockParser, Basics) {
}
}
+TEST(BlockParser, PadShortRows) {
+ auto options = ParseOptions::Defaults();
+ options.pad_short_rows = true;
+
+ BlockParser parser(options, /*num_cols=*/3);
+ AssertParseOk(parser, "1,2\n3,4,5\n");
+ AssertColumnEq(parser, 0, {"1", "3"});
+ AssertColumnEq(parser, 1, {"2", "4"});
+ std::vector<std::string> values;
+ std::vector<bool> missing;
+ ASSERT_OK(parser.VisitColumn(
+ 2, [&](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>({"", "5"}));
+ ASSERT_EQ(missing, std::vector<bool>({true, false}));
+
+ BlockParser last_row_parser(options, /*num_cols=*/3);
+ AssertParseOk(last_row_parser, "1,2\n");
+ std::vector<bool> last_row_missing;
+ ASSERT_OK(last_row_parser.VisitLastRow(
+ [&](const uint8_t*, uint32_t, bool, bool is_missing) -> Status {
+ last_row_missing.push_back(is_missing);
+ return Status::OK();
+ }));
+ ASSERT_EQ(last_row_missing, std::vector<bool>({false, false, true}));
+}
+
TEST(BlockParser, EmptyHeader) {
// Cannot infer number of columns
uint32_t out_size;
@@ -884,10 +918,12 @@ TEST(BlockParser, RowNumberAppendedToError) {
BlockParser parser(options, -1, 0);
ASSERT_NO_FATAL_FAILURE(AssertParseOk(parser, csv));
int row = 0;
- auto status = parser.VisitColumn(
- 0, [row](const uint8_t* data, uint32_t size, bool quoted) mutable ->
Status {
- return ++row == 2 ? Status::Invalid("Bad value") : Status::OK();
- });
+ auto status = parser.VisitColumn(0,
+ [row](const uint8_t* data, uint32_t size,
+ bool quoted, bool missing) mutable
-> Status {
+ return ++row == 2 ?
Status::Invalid("Bad value")
+ : Status::OK();
+ });
EXPECT_RAISES_WITH_MESSAGE_THAT(Invalid, testing::HasSubstr("Row #1: Bad
value"),
status);
}
@@ -896,10 +932,12 @@ TEST(BlockParser, RowNumberAppendedToError) {
BlockParser parser(options, -1, 100);
ASSERT_NO_FATAL_FAILURE(AssertParseOk(parser, csv));
int row = 0;
- auto status = parser.VisitColumn(
- 0, [row](const uint8_t* data, uint32_t size, bool quoted) mutable ->
Status {
- return ++row == 3 ? Status::Invalid("Bad value") : Status::OK();
- });
+ auto status = parser.VisitColumn(0,
+ [row](const uint8_t* data, uint32_t size,
+ bool quoted, bool missing) mutable
-> Status {
+ return ++row == 3 ?
Status::Invalid("Bad value")
+ : Status::OK();
+ });
EXPECT_RAISES_WITH_MESSAGE_THAT(Invalid, testing::HasSubstr("Row #102: Bad
value"),
status);
}
@@ -909,10 +947,12 @@ TEST(BlockParser, RowNumberAppendedToError) {
BlockParser parser(options, -1, -1);
ASSERT_NO_FATAL_FAILURE(AssertParseOk(parser, csv));
int row = 0;
- auto status = parser.VisitColumn(
- 0, [row](const uint8_t* data, uint32_t size, bool quoted) mutable ->
Status {
- return ++row == 3 ? Status::Invalid("Bad value") : Status::OK();
- });
+ auto status = parser.VisitColumn(0,
+ [row](const uint8_t* data, uint32_t size,
+ bool quoted, bool missing) mutable
-> Status {
+ return ++row == 3 ?
Status::Invalid("Bad value")
+ : Status::OK();
+ });
EXPECT_RAISES_WITH_MESSAGE_THAT(Invalid,
testing::Not(testing::HasSubstr("Row")),
status);
}
@@ -926,10 +966,12 @@ TEST(BlockParser, RowNumberAppendedToError) {
BlockParser parser(opts, /*num_cols=*/2, /*first_row=*/1);
ASSERT_NO_FATAL_FAILURE(AssertParseOk(parser,
"a,b,c\nd,e\nf,g\nh\ni\nj,k\nl\n"));
int row = 0;
- auto status = parser.VisitColumn(
- 0, [row](const uint8_t* data, uint32_t size, bool quoted) mutable ->
Status {
- return ++row == 3 ? Status::Invalid("Bad value") : Status::OK();
- });
+ auto status = parser.VisitColumn(0,
+ [row](const uint8_t* data, uint32_t size,
+ bool quoted, bool missing) mutable
-> Status {
+ return ++row == 3 ?
Status::Invalid("Bad value")
+ : Status::OK();
+ });
EXPECT_RAISES_WITH_MESSAGE_THAT(Invalid, testing::HasSubstr("Row #6: Bad
value"),
status);
diff --git a/cpp/src/arrow/csv/reader.cc b/cpp/src/arrow/csv/reader.cc
index 09e8290ba2..1ae9d27e6b 100644
--- a/cpp/src/arrow/csv/reader.cc
+++ b/cpp/src/arrow/csv/reader.cc
@@ -622,7 +622,9 @@ class ReaderMixin {
column_names_ = GenerateColumnNames(parser.num_cols());
} else {
// Read column names from header row
- auto visit = [&](const uint8_t* data, uint32_t size, bool quoted) ->
Status {
+ auto visit = [&](const uint8_t* data, uint32_t size, bool quoted,
+ bool missing) -> Status {
+ DCHECK(!missing);
column_names_.emplace_back(reinterpret_cast<const char*>(data),
size);
return Status::OK();
};
diff --git a/cpp/src/arrow/csv/reader_test.cc b/cpp/src/arrow/csv/reader_test.cc
index 9e0f4804b0..2493cb6627 100644
--- a/cpp/src/arrow/csv/reader_test.cc
+++ b/cpp/src/arrow/csv/reader_test.cc
@@ -29,6 +29,7 @@
#include <utility>
#include <vector>
+#include "arrow/array/array_dict.h"
#include "arrow/csv/options.h"
#include "arrow/csv/test_common.h"
#include "arrow/io/interfaces.h"
@@ -38,6 +39,7 @@
#include "arrow/testing/future_util.h"
#include "arrow/testing/gtest_util.h"
#include "arrow/util/async_generator.h"
+#include "arrow/util/checked_cast.h"
#include "arrow/util/future.h"
#include "arrow/util/thread_pool.h"
@@ -618,5 +620,49 @@ TEST(ReaderTests, DefaultColumnTypeAllStringsNoHeader) {
ASSERT_TRUE(table->Equals(*expected_table));
}
+TEST(ReaderTests, ShortRows) {
+ auto input =
+
std::make_shared<io::BufferReader>(std::make_shared<Buffer>("a,b,c\n1,2\n3,,"));
+ auto read_options = ReadOptions::Defaults();
+ read_options.block_size = 8;
+ auto parse_options = ParseOptions::Defaults();
+ parse_options.pad_short_rows = true;
+ auto convert_options = ConvertOptions::Defaults();
+ convert_options.default_column_type = utf8();
+ convert_options.null_values.clear();
+ convert_options.strings_can_be_null = false;
+
+ ASSERT_OK_AND_ASSIGN(auto reader,
+ TableReader::Make(io::default_io_context(), input,
read_options,
+ parse_options, convert_options));
+ ASSERT_OK_AND_ASSIGN(auto table, reader->Read());
+
+ auto expected_schema =
+ schema({field("a", utf8()), field("b", utf8()), field("c", utf8())});
+ auto expected_table = TableFromJSON(expected_schema, {R"([{"a":"1", "b":"2",
"c":null},
+ {"a":"3", "b":"", "c":""}])"});
+ 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();
+ read_options.autogenerate_column_names = true;
+ auto parse_options = ParseOptions::Defaults();
+ parse_options.pad_short_rows = true;
+ auto convert_options = ConvertOptions::Defaults();
+ convert_options.column_types["f0"] = int64();
+ convert_options.column_types["f1"] = dictionary(int32(), utf8());
+ ASSERT_OK_AND_ASSIGN(auto reader,
+ TableReader::Make(io::default_io_context(), input,
read_options,
+ parse_options, convert_options));
+ ASSERT_OK_AND_ASSIGN(auto table, reader->Read());
+ ASSERT_TRUE(table->column(0)->chunk(0)->Equals(*ArrayFromJSON(int64(), "[1,
2]")));
+ const auto& dict_array =
+ internal::checked_cast<const
DictionaryArray&>(*table->column(1)->chunk(0));
+ ASSERT_TRUE(dict_array.indices()->Equals(*ArrayFromJSON(int32(), "[0,
null]")));
+ ASSERT_TRUE(dict_array.dictionary()->Equals(*ArrayFromJSON(utf8(),
"[\"10\"]")));
+}
+
} // namespace csv
} // namespace arrow
diff --git a/cpp/src/arrow/dataset/file_csv.cc
b/cpp/src/arrow/dataset/file_csv.cc
index c0d85581f6..079103fa79 100644
--- a/cpp/src/arrow/dataset/file_csv.cc
+++ b/cpp/src/arrow/dataset/file_csv.cc
@@ -19,6 +19,7 @@
#include <algorithm>
#include <memory>
+#include <optional>
#include <sstream>
#include <string>
#include <unordered_set>
@@ -161,7 +162,16 @@ Result<std::vector<std::string>> GetOrderedColumnNames(
uint32_t parsed_size = 0;
int32_t max_num_rows = read_options.skip_rows + 1;
- csv::BlockParser parser(pool, parse_options, /*num_cols=*/-1,
/*first_row=*/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.
+ inspection_parse_options.emplace(parse_options);
+ inspection_parse_options->pad_short_rows = false;
+ parser_options = &*inspection_parse_options;
+ }
+ csv::BlockParser parser(pool, *parser_options, /*num_cols=*/-1,
/*first_row=*/1,
max_num_rows);
RETURN_NOT_OK(parser.Parse(std::string_view{first_block}, &parsed_size));
@@ -188,8 +198,9 @@ Result<std::vector<std::string>> GetOrderedColumnNames(
return column_names;
}
- RETURN_NOT_OK(
- parser.VisitLastRow([&](const uint8_t* data, uint32_t size, bool quoted)
-> Status {
+ RETURN_NOT_OK(parser.VisitLastRow(
+ [&](const uint8_t* data, uint32_t size, bool quoted, bool missing) ->
Status {
+ DCHECK(!missing);
std::string_view view{reinterpret_cast<const char*>(data), size};
column_names.emplace_back(view);
return Status::OK();
@@ -367,7 +378,8 @@ bool CsvFileFormat::Equals(const FileFormat& format) const {
parse_options.escaping == other_parse_options.escaping &&
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.ignore_empty_lines ==
other_parse_options.ignore_empty_lines &&
+ parse_options.pad_short_rows == other_parse_options.pad_short_rows;
}
Result<bool> CsvFileFormat::IsSupported(const FileSource& source) const {