Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions cpp/src/arrow/csv/options.h
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
62 changes: 46 additions & 16 deletions cpp/src/arrow/csv/parser.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand All @@ -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;
Expand All @@ -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;
Expand All @@ -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)) {
Expand All @@ -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;
Expand All @@ -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)) {
Expand All @@ -414,7 +438,9 @@ class BlockParserImpl {
goto InField;
}
}
parsed_writer->PushFieldChar(c);
if (!ignoring_extra_field) {
parsed_writer->PushFieldChar(c);
}
goto InQuotedField;

FieldEnd:
Expand All @@ -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);
}
Expand All @@ -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_;
}
Expand All @@ -479,7 +506,8 @@ class BlockParserImpl {
template <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;

Expand All @@ -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);
}
}
Expand Down
13 changes: 13 additions & 0 deletions cpp/src/arrow/csv/parser_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -298,6 +298,19 @@ 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, EmptyHeader) {
// Cannot infer number of columns
uint32_t out_size;
Expand Down
21 changes: 21 additions & 0 deletions cpp/src/arrow/csv/reader_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
10 changes: 6 additions & 4 deletions cpp/src/arrow/dataset/file_csv.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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 {
Expand Down
Loading