This is an automated email from the ASF dual-hosted git repository.
sollhui pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 379f648196b [fix](scan) Fix lost and duplicated rows when splitting
CSV/JSON on multi-character line delimiters (#68539)
379f648196b is described below
commit 379f648196bd13010a290a69586be237c9d77a8c
Author: hui lai <[email protected]>
AuthorDate: Tue Sep 29 10:53:34 2026 +0800
[fix](scan) Fix lost and duplicated rows when splitting CSV/JSON on
multi-character line delimiters (#68539)
### What problem does this PR solve?
When a plain-text CSV/JSON file is scanned in parallel, each non-first
split starts
a little before its offset and discards its first line; the preceding
split reads
past its end to finish its last line. This only works if both splits
agree on where
each line starts. Two cases broke that agreement:
1. **Split inside a multi-character delimiter (JSON).** The JSON readers
looked back
only one byte, so a split inside e.g. `\r\n` could not see the full
delimiter and
also discarded the next complete record. Result: silently lost rows, 0
filtered rows.
(CSV was fixed in #53374; JSON was not.)
2. **Self-overlapping delimiters (CSV and JSON).** For delimiters whose
prefix equals
their suffix (e.g. `\n\n`, `||`, `abab`), a run such as `|||` or `||||`
can be split
into delimiters in more than one way. The preceding split parses from
the file start,
while the next split parses from its lookbehind point, so they may pick
different
delimiter occurrences. Example with `||` and `A|||B||C` split at offset
4: the
unsplit read returns `A, |B, C`, but the two splits return `A, |B, B,
C`. With an
empty record (`A||||B||C`), the row count is right but `B` becomes `|B`.
### How is it fixed?
- Both JSON readers (legacy and format-v2) now look back by the
delimiter length,
clamped at file start, and extend the range size so the original split
end is kept.
- New `NewPlainTextLineReader::skip_split_prefix()` replaces "skip
exactly one line":
- Ordinary delimiters: unchanged behavior, no extra I/O.
- Self-overlapping delimiters: scan backwards to the nearest byte that
cannot be part
of any delimiter, then replay greedy matching from just after it and
skip every line
that starts before the split offset. Any scan that passes such a byte
produces the
same delimiter positions afterwards, so both splits agree with the
unsplit read.
Backward probes start at 1 KiB and double up to 64 KiB without
re-reading.
- Used by legacy/v2 JSON readers and by legacy/v2 plain (non-enclosed)
CSV readers,
for both row materialization and COUNT pushdown. First-split
header/`skip_lines`/BOM
handling is unchanged.
**Out of scope:** enclosed CSV and Hive text/CSV readers keep their
current behavior
(a delimiter sync point cannot recover quote/escape state). Tracked in
#xxxxx.
---
be/src/format/csv/csv_reader.cpp | 8 +
.../file_reader/new_plain_text_line_reader.cpp | 77 +++++++++
.../file_reader/new_plain_text_line_reader.h | 6 +
be/src/format/json/new_json_reader.cpp | 15 +-
be/src/format/json/new_json_reader.h | 5 +-
be/src/format_v2/delimited_text/csv_reader.cpp | 2 +
.../delimited_text/delimited_text_reader.cpp | 14 ++
.../delimited_text/delimited_text_reader.h | 2 +
be/src/format_v2/json/json_reader.cpp | 15 +-
be/src/format_v2/json/json_reader.h | 4 +-
.../new_plain_text_line_reader_test.cpp | 133 ++++++++++++++++
.../format_v2/delimited_text/csv_reader_test.cpp | 160 +++++++++++++++++++
be/test/format_v2/json/json_reader_test.cpp | 175 +++++++++++++++++++++
13 files changed, 598 insertions(+), 18 deletions(-)
diff --git a/be/src/format/csv/csv_reader.cpp b/be/src/format/csv/csv_reader.cpp
index b8f0be49bfe..234b3596f6f 100644
--- a/be/src/format/csv/csv_reader.cpp
+++ b/be/src/format/csv/csv_reader.cpp
@@ -416,6 +416,14 @@ Status CsvReader::_do_get_next_block(Block* block, size_t*
read_rows, bool* eof)
bool success = false;
bool is_remove_bom = false;
+ if (_range.start_offset != 0 && _skip_lines > 0 && _enclose == 0 &&
+ _file_format_type == TFileFormatType::FORMAT_CSV_PLAIN) {
+ auto* text_reader =
assert_cast<NewPlainTextLineReader*>(_line_reader.get());
+ RETURN_IF_ERROR(text_reader->skip_split_prefix(_range.start_offset,
_line_delimiter,
+ &_line_reader_eof,
_io_ctx));
+ _skip_lines = 0;
+ is_remove_bom = true;
+ }
if (_push_down_agg_type == TPushAggOp::type::COUNT) {
while (rows < batch_size && !_line_reader_eof) {
const uint8_t* ptr = nullptr;
diff --git a/be/src/format/file_reader/new_plain_text_line_reader.cpp
b/be/src/format/file_reader/new_plain_text_line_reader.cpp
index 25a7a0a7ac2..cb18ad5f21b 100644
--- a/be/src/format/file_reader/new_plain_text_line_reader.cpp
+++ b/be/src/format/file_reader/new_plain_text_line_reader.cpp
@@ -25,6 +25,7 @@
#include <immintrin.h>
#endif
#include <algorithm>
+#include <array>
#include <cstddef>
#include <cstring>
#include <ostream>
@@ -287,6 +288,82 @@ inline bool NewPlainTextLineReader::update_eof() {
return _eof;
}
+Status NewPlainTextLineReader::skip_split_prefix(size_t split_start, const
std::string& delimiter,
+ bool* eof, const
io::IOContext* io_ctx,
+ size_t* skipped_lines) {
+ DCHECK_EQ(_total_read_bytes, 0);
+ DCHECK_EQ(_output_buf_limit, 0);
+ bool overlaps = false;
+ if (delimiter.size() > 1) {
+ // Compute the KMP prefix function once for this split in linear time.
A nonempty
+ // proper prefix that is also a suffix of the whole delimiter permits
overlapping matches.
+ std::vector<size_t> prefix_lengths(delimiter.size());
+ for (size_t i = 1, matched = 0; i < delimiter.size(); ++i) {
+ while (matched > 0 && delimiter[i] != delimiter[matched]) {
+ matched = prefix_lengths[matched - 1];
+ }
+ if (delimiter[i] == delimiter[matched]) {
+ ++matched;
+ }
+ prefix_lengths[i] = matched;
+ }
+ overlaps = prefix_lengths.back() > 0;
+ }
+
+ if (overlaps && _decompressor == nullptr) {
+ // A fixed lookbehind can start in an overlapping delimiter chain
(e.g. three newlines
+ // with a two-newline delimiter). Find a byte that cannot belong to
any delimiter, then
+ // replay greedy matches from immediately after it. Keep scratch
bounded even for long
+ // delimiter runs; ordinary delimiters do not need this extra I/O.
+ std::array<bool, 256> delimiter_bytes {};
+ for (unsigned char byte : delimiter) {
+ delimiter_bytes[byte] = true;
+ }
+ constexpr size_t max_lookbehind_size = 64 * 1024;
+ std::vector<char> buffer(1024);
+ size_t sync_offset = _current_offset;
+ bool synchronized = false;
+ while (sync_offset > 0 && !synchronized) {
+ const size_t length = std::min(sync_offset, buffer.size());
+ const size_t offset = sync_offset - length;
+ size_t bytes_read = 0;
+ RETURN_IF_ERROR(_file_reader->read_at(offset, Slice(buffer.data(),
length), &bytes_read,
+ io_ctx));
+ if (bytes_read != length) {
+ return Status::IOError("Short read while aligning text split
at offset {}",
+ split_start);
+ }
+ sync_offset = offset;
+ for (size_t i = length; i > 0; --i) {
+ if (!delimiter_bytes[static_cast<unsigned char>(buffer[i -
1])]) {
+ sync_offset = offset + i;
+ synchronized = true;
+ break;
+ }
+ }
+ if (!synchronized && sync_offset > 0) {
+ // Extend backward into new bytes; do not reread the already
searched suffix.
+ buffer.resize(std::min(buffer.size() * 2,
max_lookbehind_size));
+ }
+ }
+ _min_length += _current_offset - sync_offset;
+ _current_offset = sync_offset;
+ }
+
+ const size_t prefix_length = split_start - _current_offset;
+ const uint8_t* line = nullptr;
+ size_t size = 0;
+ size_t skipped = 0;
+ do {
+ RETURN_IF_ERROR(read_line(&line, &size, eof, io_ctx));
+ skipped += !*eof;
+ } while (!*eof && _total_read_bytes < prefix_length);
+ if (skipped_lines != nullptr) {
+ *skipped_lines = skipped;
+ }
+ return Status::OK();
+}
+
// extend input buf if necessary only when _more_input_bytes > 0
void NewPlainTextLineReader::extend_input_buf() {
DCHECK(_more_input_bytes > 0);
diff --git a/be/src/format/file_reader/new_plain_text_line_reader.h
b/be/src/format/file_reader/new_plain_text_line_reader.h
index ed7f80493b0..72c090556b6 100644
--- a/be/src/format/file_reader/new_plain_text_line_reader.h
+++ b/be/src/format/file_reader/new_plain_text_line_reader.h
@@ -248,6 +248,12 @@ public:
Status read_line(const uint8_t** ptr, size_t* size, bool* eof,
const io::IOContext* io_ctx) override;
+ // Called before the first read of a non-first plain-text split. Discard
records owned by the
+ // preceding split, preserving greedy delimiter matching. Not suitable for
enclosed CSV:
+ // finding a delimiter synchronization point does not recover quote/escape
state.
+ Status skip_split_prefix(size_t split_start, const std::string& delimiter,
bool* eof,
+ const io::IOContext* io_ctx, size_t*
skipped_lines = nullptr);
+
inline TextLineReaderCtxPtr text_line_reader_ctx() { return
_line_reader_ctx; }
void close() override;
diff --git a/be/src/format/json/new_json_reader.cpp
b/be/src/format/json/new_json_reader.cpp
index 1aa19574b39..33d17efdf9b 100644
--- a/be/src/format/json/new_json_reader.cpp
+++ b/be/src/format/json/new_json_reader.cpp
@@ -151,6 +151,8 @@ NewJsonReader::NewJsonReader(RuntimeProfile* profile, const
TFileScanRangeParams
_init_file_description();
}
+NewJsonReader::~NewJsonReader() = default;
+
void NewJsonReader::_init_system_properties() {
if (_range.__isset.file_type) {
// for compatibility
@@ -265,9 +267,8 @@ Status NewJsonReader::_do_get_next_block(Block* block,
size_t* read_rows, bool*
while (block->rows() < batch_size && !_reader_eof && (block->bytes() <
max_block_bytes)) {
if (UNLIKELY(_read_json_by_line && _skip_first_line)) {
- size_t size = 0;
- const uint8_t* line_ptr = nullptr;
- RETURN_IF_ERROR(_line_reader->read_line(&line_ptr, &size,
&_reader_eof, _io_ctx));
+
RETURN_IF_ERROR(_line_reader->skip_split_prefix(_range.start_offset,
_line_delimiter,
+ &_reader_eof,
_io_ctx));
_skip_first_line = false;
continue;
}
@@ -487,7 +488,9 @@ void
json_reader_detail::pop_back_last_inserted_value(Block& block, size_t colum
Status NewJsonReader::_open_file_reader(bool need_schema) {
int64_t start_offset = _range.start_offset;
if (start_offset != 0) {
- start_offset -= 1;
+ // Include the whole delimiter when the split starts inside it, so
skipping the first
+ // partial line cannot discard the next complete JSON record.
+ start_offset -= std::min<int64_t>(start_offset,
_line_delimiter_length);
}
_current_offset = start_offset;
@@ -524,8 +527,8 @@ Status NewJsonReader::_open_file_reader(bool need_schema) {
Status NewJsonReader::_open_line_reader() {
int64_t size = _range.size;
if (_range.start_offset != 0) {
- // When we fetch range doesn't start from 0, size will += 1.
- size += 1;
+ // Preserve the original range end after moving the start backwards.
+ size += _range.start_offset - _current_offset;
_skip_first_line = true;
} else {
_skip_first_line = false;
diff --git a/be/src/format/json/new_json_reader.h
b/be/src/format/json/new_json_reader.h
index b975433c34f..029af5adc27 100644
--- a/be/src/format/json/new_json_reader.h
+++ b/be/src/format/json/new_json_reader.h
@@ -61,6 +61,7 @@ struct IOContext;
struct ScannerCounter;
class Block;
class IColumn;
+class NewPlainTextLineReader;
namespace json_reader_detail {
Status append_null_for_malformed_json(Block& block);
@@ -89,7 +90,7 @@ public:
const TFileRangeDesc& range, const
std::vector<SlotDescriptor*>& file_slot_descs,
size_t batch_size, io::IOContext* io_ctx,
std::shared_ptr<io::IOContext> io_ctx_holder = nullptr);
- ~NewJsonReader() override = default;
+ ~NewJsonReader() override;
Status init_reader(
const std::unordered_map<std::string, VExprContextSPtr>&
col_default_value_ctx,
@@ -209,7 +210,7 @@ private:
const std::vector<SlotDescriptor*>& _file_slot_descs;
io::FileReaderSPtr _file_reader;
- std::unique_ptr<LineReader> _line_reader;
+ std::unique_ptr<NewPlainTextLineReader> _line_reader;
bool _reader_eof;
std::unique_ptr<Decompressor> _decompressor;
TFileCompressType::type _file_compress_type;
diff --git a/be/src/format_v2/delimited_text/csv_reader.cpp
b/be/src/format_v2/delimited_text/csv_reader.cpp
index 5bfd57346b9..bc4026138f1 100644
--- a/be/src/format_v2/delimited_text/csv_reader.cpp
+++ b/be/src/format_v2/delimited_text/csv_reader.cpp
@@ -141,6 +141,7 @@ Status CsvReader::_create_decompressor() {
}
Status CsvReader::_create_line_reader() {
+ _align_split_prefix = false;
if (is_csv_text_format(_file_format_type)) {
std::shared_ptr<TextLineReaderContextIf> text_line_reader_ctx;
if (_hive_csv_parser) {
@@ -148,6 +149,7 @@ Status CsvReader::_create_line_reader() {
} else if (_enclose == 0) {
text_line_reader_ctx = std::make_shared<PlainTextLineReaderCtx>(
_line_delimiter, _line_delimiter.size(), _keep_cr);
+ _align_split_prefix = _file_description->range_start_offset != 0;
} else {
const size_t col_sep_num =
_source_file_slot_descs.size() > 1 ?
_source_file_slot_descs.size() - 1 : 0;
diff --git a/be/src/format_v2/delimited_text/delimited_text_reader.cpp
b/be/src/format_v2/delimited_text/delimited_text_reader.cpp
index 63486d174ef..fa185600cc7 100644
--- a/be/src/format_v2/delimited_text/delimited_text_reader.cpp
+++ b/be/src/format_v2/delimited_text/delimited_text_reader.cpp
@@ -32,6 +32,7 @@
#include "core/data_type/data_type_nullable.h"
#include "core/data_type/data_type_string.h"
#include "core/data_type/data_type_struct.h"
+#include "format/file_reader/new_plain_text_line_reader.h"
#include "format/line_reader.h"
#include "format_v2/column_mapper.h"
#include "format_v2/materialized_reader_util.h"
@@ -555,6 +556,19 @@ Status DelimitedTextReader::_open_file() {
Status DelimitedTextReader::_read_next_line(Slice* line, bool* eof) {
DORIS_CHECK(line != nullptr);
DORIS_CHECK(eof != nullptr);
+ if (_align_split_prefix) {
+ SCOPED_TIMER(_text_profile.read_line_time);
+ DCHECK_EQ(_skip_lines, 1);
+ size_t skipped_lines = 0;
+ auto* text_reader =
assert_cast<NewPlainTextLineReader*>(_line_reader.get());
+
RETURN_IF_ERROR(text_reader->skip_split_prefix(_file_description->range_start_offset,
+ _line_delimiter,
&_line_reader_eof,
+ _io_ctx.get(),
&skipped_lines));
+ _align_split_prefix = false;
+ _skip_lines = 0;
+ _bom_removed = true;
+ update_counter(_text_profile.skipped_lines, skipped_lines);
+ }
while (true) {
const uint8_t* ptr = nullptr;
size_t size = 0;
diff --git a/be/src/format_v2/delimited_text/delimited_text_reader.h
b/be/src/format_v2/delimited_text/delimited_text_reader.h
index dff27980c9c..472f21f59c9 100644
--- a/be/src/format_v2/delimited_text/delimited_text_reader.h
+++ b/be/src/format_v2/delimited_text/delimited_text_reader.h
@@ -156,6 +156,8 @@ protected:
int64_t _start_offset = 0;
int64_t _size = -1;
int _skip_lines = 0;
+ // Enabled only by readers using plain, quote-independent delimiter
matching.
+ bool _align_split_prefix = false;
char _escape = 0;
bool _line_reader_eof = false;
bool _bom_removed = false;
diff --git a/be/src/format_v2/json/json_reader.cpp
b/be/src/format_v2/json/json_reader.cpp
index caf4205fa61..78601f90584 100644
--- a/be/src/format_v2/json/json_reader.cpp
+++ b/be/src/format_v2/json/json_reader.cpp
@@ -314,10 +314,8 @@ Status JsonReader::get_block(Block* file_block, size_t*
rows, bool* eof) {
while (file_block->rows() < batch_size && !_reader_eof &&
file_block->bytes() < max_block_bytes) {
if (_read_json_by_line && _skip_first_line) {
- size_t skipped_size = 0;
- const uint8_t* skipped_line = nullptr;
- RETURN_IF_ERROR(_line_reader->read_line(&skipped_line,
&skipped_size, &_reader_eof,
- _io_ctx.get()));
+ RETURN_IF_ERROR(_line_reader->skip_split_prefix(
+ _reader_range.start_offset, _line_delimiter, &_reader_eof,
_io_ctx.get()));
_skip_first_line = false;
continue;
}
@@ -445,7 +443,9 @@ TFileRangeDesc JsonReader::_json_range() const {
Status JsonReader::_open_file_reader() {
_current_offset = _reader_range.start_offset;
if (_current_offset != 0) {
- --_current_offset;
+ // Include the whole delimiter when the split starts inside it, so
skipping the first
+ // partial line cannot discard the next complete JSON record.
+ _current_offset -= std::min<int64_t>(_current_offset,
_line_delimiter_length);
}
if (_scan_params->file_type == TFileType::FILE_STREAM) {
if (!_stream_load_id.has_value()) {
@@ -478,9 +478,8 @@ Status JsonReader::_create_decompressor() {
Status JsonReader::_create_line_reader() {
int64_t size = _reader_range.size;
if (_reader_range.start_offset != 0) {
- // Start one byte earlier and discard the first partial line, matching
split semantics used
- // by text readers.
- ++size;
+ // Preserve the original range end after moving the start backwards.
+ size += _reader_range.start_offset - _current_offset;
_skip_first_line = true;
} else {
_skip_first_line = false;
diff --git a/be/src/format_v2/json/json_reader.h
b/be/src/format_v2/json/json_reader.h
index c7346cb1d66..f5e6614b2da 100644
--- a/be/src/format_v2/json/json_reader.h
+++ b/be/src/format_v2/json/json_reader.h
@@ -35,7 +35,7 @@
namespace doris {
class Decompressor;
-class LineReader;
+class NewPlainTextLineReader;
class SlotDescriptor;
class IColumn;
} // namespace doris
@@ -188,7 +188,7 @@ private:
io::FileReaderSPtr _physical_file_reader;
std::unique_ptr<Decompressor> _decompressor;
- std::unique_ptr<LineReader> _line_reader;
+ std::unique_ptr<NewPlainTextLineReader> _line_reader;
int64_t _current_offset = 0;
bool _reader_eof = false;
bool _skip_first_line = false;
diff --git a/be/test/format/file_reader/new_plain_text_line_reader_test.cpp
b/be/test/format/file_reader/new_plain_text_line_reader_test.cpp
index 93d02067863..3d8ddd57f0c 100644
--- a/be/test/format/file_reader/new_plain_text_line_reader_test.cpp
+++ b/be/test/format/file_reader/new_plain_text_line_reader_test.cpp
@@ -21,8 +21,141 @@
#include <gtest/gtest.h>
+#include <algorithm>
+#include <cstring>
+#include <utility>
+
+#include "io/fs/file_reader.h"
+
namespace doris {
+namespace {
+class RecordingSplitFileReader : public io::FileReader {
+public:
+ explicit RecordingSplitFileReader(std::string content) :
_content(std::move(content)) {}
+ Status close() override {
+ _closed = true;
+ return Status::OK();
+ }
+ const io::Path& path() const override { return _path; }
+ size_t size() const override { return _content.size(); }
+ bool closed() const override { return _closed; }
+ int64_t mtime() const override { return 0; }
+
+ std::vector<std::pair<size_t, size_t>> requests;
+
+protected:
+ Status read_at_impl(size_t offset, Slice result, size_t* bytes_read,
+ const io::IOContext*) override {
+ requests.emplace_back(offset, result.size);
+ *bytes_read = std::min(result.size, _content.size() - std::min(offset,
_content.size()));
+ if (*bytes_read > 0) {
+ std::memcpy(result.mutable_data(), _content.data() + offset,
*bytes_read);
+ }
+ return Status::OK();
+ }
+
+private:
+ io::Path _path {"split-prefix-test"};
+ std::string _content;
+ bool _closed = false;
+};
+} // namespace
+
+TEST(PlainTextSplitPrefixTest, DelimiterOverlapControlsBackwardProbes) {
+ const std::vector<std::pair<std::string, bool>> delimiters = {
+ {"\n", false},
+ {"abc", false},
+ {"aaaaab", false},
+ {"||", true},
+ {"aba", true},
+ {"abcab", true},
+ {"ababcabab", true},
+ {std::string(100 * 1024 - 1, 'a') + "b", false},
+ {std::string(100 * 1024, 'a'), true}};
+ for (const auto& [delimiter, overlaps] : delimiters) {
+ SCOPED_TRACE(::testing::Message() << "delimiter length=" <<
delimiter.size()
+ << ", prefix=" <<
delimiter.substr(0, 16));
+ const std::string first = "xxx";
+ const std::string content = first + delimiter + "second" + delimiter +
"third";
+ const size_t split = first.size() + delimiter.size();
+ auto file = std::make_shared<RecordingSplitFileReader>(content);
+ RuntimeProfile profile("split_prefix");
+ NewPlainTextLineReader reader(
+ &profile, file, nullptr,
+ std::make_shared<PlainTextLineReaderCtx>(delimiter,
delimiter.size(), false),
+ content.size() - first.size(), first.size());
+ bool eof = false;
+ size_t skipped_lines = 0;
+ ASSERT_TRUE(reader.skip_split_prefix(split, delimiter, &eof, nullptr,
&skipped_lines).ok());
+ ASSERT_FALSE(eof);
+ EXPECT_EQ(skipped_lines, 1);
+ ASSERT_FALSE(file->requests.empty());
+ // Only a self-overlapping delimiter requires a probe before the
initial read offset.
+ EXPECT_EQ(file->requests.front().first < first.size(), overlaps);
+ const uint8_t* line = nullptr;
+ size_t size = 0;
+ ASSERT_TRUE(reader.read_line(&line, &size, &eof, nullptr).ok());
+ EXPECT_EQ(std::string(reinterpret_cast<const char*>(line), size),
"second");
+ }
+}
+
+TEST(PlainTextSplitPrefixTest, NearbySynchronizationUsesSmallProbe) {
+ const std::string content = std::string(8192, 'x') + "a|||b||c";
+ const size_t split = 8192 + 4;
+ auto file = std::make_shared<RecordingSplitFileReader>(content);
+ RuntimeProfile profile("split_prefix");
+ NewPlainTextLineReader reader(&profile, file, nullptr,
+
std::make_shared<PlainTextLineReaderCtx>("||", 2, false),
+ content.size() - split + 2, split - 2);
+ bool eof = false;
+ size_t skipped_lines = 0;
+ ASSERT_TRUE(reader.skip_split_prefix(split, "||", &eof, nullptr,
&skipped_lines).ok());
+ ASSERT_FALSE(eof);
+ ASSERT_GE(file->requests.size(), 2);
+ EXPECT_EQ(file->requests.front().second, 1024);
+ // The next read replays forward from the synchronization point; no larger
probe was needed.
+ EXPECT_GT(file->requests[1].first, file->requests[0].first);
+ EXPECT_EQ(skipped_lines, 2);
+ const uint8_t* line = nullptr;
+ size_t size = 0;
+ ASSERT_TRUE(reader.read_line(&line, &size, &eof, nullptr).ok());
+ EXPECT_EQ(std::string(reinterpret_cast<const char*>(line), size), "c");
+}
+
+TEST(PlainTextSplitPrefixTest, LongRunGrowsProbesWithoutRereading) {
+ const std::string prefix = "a" + std::string(256 * 1024 + 1, '|');
+ const std::string content = prefix + "b||c";
+ const size_t split = prefix.size();
+ auto file = std::make_shared<RecordingSplitFileReader>(content);
+ RuntimeProfile profile("split_prefix");
+ NewPlainTextLineReader reader(&profile, file, nullptr,
+
std::make_shared<PlainTextLineReaderCtx>("||", 2, false),
+ content.size() - split + 2, split - 2);
+ bool eof = false;
+ ASSERT_TRUE(reader.skip_split_prefix(split, "||", &eof, nullptr).ok());
+ ASSERT_FALSE(eof);
+ size_t previous_offset = split - 2;
+ size_t probe_size = 1024;
+ size_t probes = 0;
+ for (const auto& [offset, length] : file->requests) {
+ if (offset >= previous_offset) {
+ break; // Forward replay has begun.
+ }
+ EXPECT_EQ(offset + length, previous_offset);
+ EXPECT_EQ(length, std::min(previous_offset, probe_size));
+ previous_offset = offset;
+ probe_size = std::min(probe_size * 2, size_t {64 * 1024});
+ ++probes;
+ }
+ EXPECT_GT(probes, 6);
+ EXPECT_EQ(previous_offset, 0);
+ const uint8_t* line = nullptr;
+ size_t size = 0;
+ ASSERT_TRUE(reader.read_line(&line, &size, &eof, nullptr).ok());
+ EXPECT_EQ(std::string(reinterpret_cast<const char*>(line), size), "c");
+}
+
// Base test class for text line reader tests
class PlainTextLineReaderTest : public testing::Test {
protected:
diff --git a/be/test/format_v2/delimited_text/csv_reader_test.cpp
b/be/test/format_v2/delimited_text/csv_reader_test.cpp
index c04e17e07f6..8c537b1e60e 100644
--- a/be/test/format_v2/delimited_text/csv_reader_test.cpp
+++ b/be/test/format_v2/delimited_text/csv_reader_test.cpp
@@ -39,13 +39,16 @@
#include "core/data_type/data_type_number.h"
#include "core/data_type/data_type_string.h"
#include "core/data_type/data_type_struct.h"
+#include "exec/scan/scanner.h"
#include "exprs/vexpr.h"
#include "exprs/vexpr_context.h"
+#include "format/csv/csv_reader.h"
#include "format_v2/column_mapper.h"
#include "io/io_common.h"
#include "runtime/runtime_profile.h"
#include "testutil/desc_tbl_builder.h"
#include "testutil/mock/mock_runtime_state.h"
+#include "testutil/scoped_temp_dir.h"
#include "util/debug_points.h"
#include "util/defer_op.h"
@@ -324,6 +327,163 @@ VExprContextSPtr prepared_conjunct(RuntimeState* state,
const VExprSPtr& expr) {
return context;
}
+class PlainCsvSplitTest : public testing::TestWithParam<bool> {
+protected:
+ void read_range(const std::string& content, const std::string& delimiter,
int64_t start,
+ int64_t size, bool count_only, std::vector<std::string>*
values,
+ size_t* total_rows, int header_mode = 0) {
+ const auto path = (_dir.path() / "split.csv").string();
+ std::ofstream(path, std::ios::binary) << content;
+ auto params = csv_scan_params();
+ params.__set_compress_type(TFileCompressType::PLAIN);
+ params.__set_column_idxs({0});
+ params.file_attributes.__isset.header_type = false;
+ params.file_attributes.text_params.__set_line_delimiter(delimiter);
+ if (header_mode == 1) {
+ params.file_attributes.__set_header_type(BeConsts::CSV_WITH_NAMES);
+ } else if (header_mode == 2) {
+
params.file_attributes.__set_header_type(BeConsts::CSV_WITH_NAMES_AND_TYPES);
+ } else if (header_mode == 3) {
+ params.file_attributes.__set_skip_lines(2);
+ }
+ ObjectPool pool;
+ auto type = make_nullable(std::make_shared<DataTypeString>());
+ std::vector<SlotDescriptor*> slots {make_test_slot(&pool, 0, 0, type,
"id")};
+ MockRuntimeState state;
+ state._batch_size = 2;
+ RuntimeProfile profile("plain_csv_split_test");
+ auto read_blocks = [&](auto&& next_block) {
+ bool eof = false;
+ while (!eof) {
+ Block block;
+ block.insert({type->create_column(), type, "id"});
+ size_t rows = 0;
+ auto status = next_block(&block, &rows, &eof);
+ ASSERT_TRUE(status.ok()) << status;
+ ASSERT_EQ(rows, block.rows());
+ *total_rows += rows;
+ if (!count_only) {
+ for (size_t row = 0; row < rows; ++row) {
+
ASSERT_FALSE(is_null_at(*block.get_by_position(0).column, row));
+ values->push_back(
+
nullable_string_at(*block.get_by_position(0).column, row));
+ }
+ }
+ }
+ };
+ if (GetParam()) {
+ auto reader = create_reader(path, ¶ms, slots, &state,
&profile, start, size);
+ auto request = std::make_shared<FileScanRequest>();
+ request->local_positions.emplace(LocalColumnId(0), LocalIndex(0));
+ ASSERT_TRUE(reader->open(request).ok());
+ if (count_only) {
+ FileAggregateRequest aggregate_request;
+ aggregate_request.agg_type = TPushAggOp::type::COUNT;
+ FileAggregateResult result;
+ ASSERT_TRUE(reader->get_aggregate_result(aggregate_request,
&result).ok());
+ *total_rows += result.count;
+ } else {
+ read_blocks([&](Block* block, size_t* rows, bool* eof) {
+ return reader->get_block(block, rows, eof);
+ });
+ }
+ } else {
+ TFileRangeDesc range;
+ range.__set_path(path);
+ range.__set_start_offset(start);
+ range.__set_size(size);
+ range.__set_file_size(content.size());
+ ScannerCounter counter;
+ auto reader = ::doris::CsvReader::create_unique(
+ &state, &profile, &counter, params, range, slots,
state.batch_size(), nullptr);
+ ASSERT_TRUE(reader->init_reader(true).ok());
+ if (count_only) {
+ reader->set_push_down_agg_type(TPushAggOp::type::COUNT);
+ }
+ read_blocks([&](Block* block, size_t* rows, bool* eof) {
+ return reader->get_next_block(block, rows, eof);
+ });
+ EXPECT_EQ(counter.num_rows_filtered, 0);
+ }
+ }
+
+ doris::test::ScopedTempDirectory _dir {"doris_plain_csv_split_test"};
+};
+
+TEST_P(PlainCsvSplitTest, EveryByteBoundaryMatchesUnsplitRowsAndCount) {
+ for (const std::string delimiter : {"\n", "\r\n", "ABCDE", "||", "aba"}) {
+ const std::string extra = delimiter == "||" ? "|" : delimiter == "aba"
? "ba" : "";
+ for (bool trailing : {false, true}) {
+ const std::string content =
+ "1" + delimiter + extra + "2" + delimiter + "3" +
(trailing ? delimiter : "");
+ const auto file_size = static_cast<int64_t>(content.size());
+ const std::vector<std::string> expected {"1", extra + "2", "3"};
+ for (bool count_only : {false, true}) {
+ SCOPED_TRACE(testing::Message() << "delimiter=" << delimiter
<< ", trailing="
+ << trailing << ", count=" <<
count_only);
+ std::vector<std::string> unsplit;
+ size_t unsplit_rows = 0;
+ ASSERT_NO_FATAL_FAILURE(read_range(content, delimiter, 0,
file_size, count_only,
+ &unsplit, &unsplit_rows));
+ ASSERT_EQ(unsplit_rows, expected.size());
+ if (!count_only) {
+ ASSERT_EQ(unsplit, expected);
+ }
+ for (int64_t split = 1; split < file_size; ++split) {
+ SCOPED_TRACE(testing::Message() << "split=" << split);
+ std::vector<std::string> values;
+ size_t rows = 0;
+ ASSERT_NO_FATAL_FAILURE(
+ read_range(content, delimiter, 0, split,
count_only, &values, &rows));
+ ASSERT_NO_FATAL_FAILURE(read_range(content, delimiter,
split, file_size - split,
+ count_only, &values,
&rows));
+ ASSERT_EQ(rows, expected.size());
+ if (!count_only) {
+ ASSERT_EQ(values, expected);
+ }
+ }
+ std::vector<std::string> values;
+ size_t rows = 0;
+ for (int64_t start = 0; start < file_size; ++start) {
+ ASSERT_NO_FATAL_FAILURE(
+ read_range(content, delimiter, start, 1,
count_only, &values, &rows));
+ }
+ EXPECT_EQ(rows, expected.size());
+ if (!count_only) {
+ EXPECT_EQ(values, expected);
+ }
+ }
+ }
+ }
+}
+
+TEST_P(PlainCsvSplitTest, FirstSplitStillHonorsHeadersAndSkipLines) {
+ for (int header_mode : {1, 2, 3}) {
+ const std::string header = header_mode == 1 ? "id||" : "id||String||";
+ const std::string content = header + "1|||2||3";
+ const auto split = static_cast<int64_t>(header.size() + 4);
+ for (bool count_only : {false, true}) {
+ SCOPED_TRACE(testing::Message()
+ << "header_mode=" << header_mode << ", count=" <<
count_only);
+ std::vector<std::string> values;
+ size_t rows = 0;
+ ASSERT_NO_FATAL_FAILURE(
+ read_range(content, "||", 0, split, count_only, &values,
&rows, header_mode));
+ ASSERT_NO_FATAL_FAILURE(read_range(content, "||", split,
content.size() - split,
+ count_only, &values, &rows,
header_mode));
+ EXPECT_EQ(rows, 3);
+ if (!count_only) {
+ EXPECT_EQ(values, (std::vector<std::string> {"1", "|2", "3"}));
+ }
+ }
+ }
+}
+
+INSTANTIATE_TEST_SUITE_P(LegacyAndV2, PlainCsvSplitTest, testing::Bool(),
+ [](const testing::TestParamInfo<bool>& info) {
+ return info.param ? "V2" : "Legacy";
+ });
+
class CsvV2ReaderTest : public testing::Test {
public:
void SetUp() override {
diff --git a/be/test/format_v2/json/json_reader_test.cpp
b/be/test/format_v2/json/json_reader_test.cpp
index c55bb3f512e..f2f48ba6367 100644
--- a/be/test/format_v2/json/json_reader_test.cpp
+++ b/be/test/format_v2/json/json_reader_test.cpp
@@ -38,8 +38,10 @@
#include "core/data_type/data_type_number.h"
#include "core/data_type/data_type_string.h"
#include "core/data_type/data_type_struct.h"
+#include "exec/scan/scanner.h"
#include "exprs/vexpr.h"
#include "exprs/vexpr_context.h"
+#include "format/json/new_json_reader.h"
#include "format_v2/column_data.h"
#include "io/io_common.h"
#include "runtime/descriptors.h"
@@ -301,6 +303,179 @@ VExprContextSPtr prepared_conjunct(RuntimeState* state,
const VExprSPtr& expr) {
} // namespace
+// Exercise both readers with the same physical splits and compare the
complete record sequence,
+// since a row-count assertion alone can hide one lost record and one
duplicated record.
+class JsonReaderSplitTest : public testing::TestWithParam<bool> {
+protected:
+ void read_range(const std::filesystem::path& path, const std::string&
delimiter, int64_t start,
+ int64_t size, std::vector<int32_t>* ids) {
+ auto params = json_scan_params();
+ params.file_attributes.text_params.__set_line_delimiter(delimiter);
+ auto range = file_range(path);
+ range.__set_start_offset(start);
+ range.__set_size(size);
+ ObjectPool pool;
+ auto type = make_nullable(std::make_shared<DataTypeInt32>());
+ std::vector<SlotDescriptor*> slots {make_test_slot(&pool, 0, 0, type,
"id")};
+ RuntimeProfile profile("json_split_test");
+ MockRuntimeState state;
+ state._batch_size = 2;
+
+ auto read_blocks = [&](auto&& next_block) {
+ bool eof = false;
+ while (!eof) {
+ Block block;
+ block.insert({type->create_column(), type, "id"});
+ size_t rows = 0;
+ auto status = next_block(&block, &rows, &eof);
+ ASSERT_TRUE(status.ok()) << status;
+ ASSERT_EQ(rows, block.rows());
+ const auto& nullable =
+ assert_cast<const
ColumnNullable&>(*block.get_by_position(0).column);
+ const auto& column = assert_cast<const
ColumnInt32&>(nullable.get_nested_column());
+ for (size_t row = 0; row < rows; ++row) {
+ ASSERT_FALSE(nullable.is_null_at(row));
+ ids->push_back(column.get_element(row));
+ }
+ }
+ };
+
+ if (GetParam()) {
+ auto properties = std::make_shared<io::FileSystemProperties>();
+ properties->system_type = TFileType::FILE_LOCAL;
+ auto desc = file_description(path.string());
+ desc->range_start_offset = start;
+ desc->range_size = size;
+ JsonReader reader(properties, desc, nullptr, &profile, ¶ms,
range, slots);
+ ASSERT_TRUE(reader.init(&state).ok());
+ auto request = std::make_shared<FileScanRequest>();
+ request->local_positions.emplace(LocalColumnId(0), LocalIndex(0));
+ ASSERT_TRUE(reader.open(request).ok());
+ read_blocks([&](Block* block, size_t* rows, bool* eof) {
+ return reader.get_block(block, rows, eof);
+ });
+ } else {
+ ScannerCounter counter;
+ bool scanner_eof = false;
+ auto reader =
+ NewJsonReader::create_unique(&state, &profile, &counter,
params, range, slots,
+ &scanner_eof,
state.batch_size(), nullptr);
+ ASSERT_TRUE(reader->init_reader({}, true).ok());
+ read_blocks([&](Block* block, size_t* rows, bool* eof) {
+ return reader->get_next_block(block, rows, eof);
+ });
+ EXPECT_EQ(counter.num_rows_filtered, 0);
+ }
+ }
+};
+
+TEST_P(JsonReaderSplitTest, EveryByteBoundaryPreservesRecords) {
+ // Includes single-byte delimiters, CRLF, UTF-8 bytes, and a delimiter
longer than a record
+ // to cover starts smaller than the amount of lookbehind. Test EOF with
and without a delimiter.
+ for (const std::string delimiter : {"\n", "\r\n", "ABCDE", "\xE2\x98\x83",
"ABCDEFGHIJKLM"}) {
+ for (bool trailing_delimiter : {false, true}) {
+ SCOPED_TRACE(testing::Message()
+ << "delimiter=" << delimiter << ", trailing=" <<
trailing_delimiter);
+ std::string content =
+ R"({"id":1})" + delimiter + R"({"id":2})" + delimiter +
R"({"id":3})";
+ if (trailing_delimiter) {
+ content += delimiter;
+ }
+ const auto path = write_json_file("split_boundaries.json",
content);
+ const auto file_size = static_cast<int64_t>(content.size());
+ const std::vector<int32_t> expected {1, 2, 3};
+ std::vector<int32_t> unsplit;
+ ASSERT_NO_FATAL_FAILURE(read_range(path, delimiter, 0, file_size,
&unsplit));
+ ASSERT_EQ(unsplit, expected);
+ for (int64_t split = 1; split < file_size; ++split) {
+ SCOPED_TRACE(testing::Message() << "split=" << split);
+ std::vector<int32_t> ids;
+ ASSERT_NO_FATAL_FAILURE(read_range(path, delimiter, 0, split,
&ids));
+ ASSERT_NO_FATAL_FAILURE(
+ read_range(path, delimiter, split, file_size - split,
&ids));
+ ASSERT_EQ(ids, expected);
+ }
+ }
+ }
+}
+
+TEST_P(JsonReaderSplitTest, OverlappingDelimitersPreserveRecords) {
+ for (const std::string delimiter : {"\n\n", "\r\n\r\n", " \t "}) {
+ for (bool trailing_delimiter : {false, true}) {
+ SCOPED_TRACE(testing::Message()
+ << "delimiter=" << delimiter << ", trailing=" <<
trailing_delimiter);
+ // The delimiter followed by its prefix has overlapping matches.
For "\n\n", a split
+ // at byte 11 previously returned {1, 2, 2, 3}: the two splits
matched different pairs
+ // of newlines in the three-newline run before id=2.
+ std::string content = R"({"id":1})" + delimiter +
+ delimiter.substr(0, delimiter.size() / 2) +
R"({"id":2})" +
+ delimiter + R"({"id":3})";
+ if (trailing_delimiter) {
+ content += delimiter;
+ }
+ const auto path = write_json_file("overlapping_delimiters.json",
content);
+ const auto file_size = static_cast<int64_t>(content.size());
+ const std::vector<int32_t> expected {1, 2, 3};
+ std::vector<int32_t> unsplit;
+ ASSERT_NO_FATAL_FAILURE(read_range(path, delimiter, 0, file_size,
&unsplit));
+ ASSERT_EQ(unsplit, expected);
+ for (int64_t split = 1; split < file_size; ++split) {
+ SCOPED_TRACE(testing::Message() << "split=" << split);
+ std::vector<int32_t> ids;
+ ASSERT_NO_FATAL_FAILURE(read_range(path, delimiter, 0, split,
&ids));
+ ASSERT_NO_FATAL_FAILURE(
+ read_range(path, delimiter, split, file_size - split,
&ids));
+ ASSERT_EQ(ids, expected);
+ }
+ std::vector<int32_t> ids;
+ for (int64_t start = 0; start < file_size; ++start) {
+ ASSERT_NO_FATAL_FAILURE(read_range(path, delimiter, start, 1,
&ids));
+ }
+ EXPECT_EQ(ids, expected);
+ }
+ }
+}
+
+TEST_P(JsonReaderSplitTest, OverlappingDelimiterRunCrossesLookbehindBuffers) {
+ const std::string delimiter = "\n\n";
+ // An odd run longer than the alignment scratch buffer must be replayed
from its true start.
+ const std::string prefix = R"({"id":1})" + std::string(64 * 1024 + 3,
'\n');
+ const std::string content = prefix + R"({"id":2})" + delimiter +
R"({"id":3})";
+ const auto path = write_json_file("long_overlapping_delimiters.json",
content);
+ const auto file_size = static_cast<int64_t>(content.size());
+ const auto split = static_cast<int64_t>(prefix.size());
+ const std::vector<int32_t> expected {1, 2, 3};
+ std::vector<int32_t> ids;
+ ASSERT_NO_FATAL_FAILURE(read_range(path, delimiter, 0, split, &ids));
+ ASSERT_NO_FATAL_FAILURE(read_range(path, delimiter, split, file_size -
split, &ids));
+ EXPECT_EQ(ids, expected);
+}
+
+TEST_P(JsonReaderSplitTest, FourRangesPreserveEveryRecord) {
+ const std::string delimiter = "ABCDE";
+ std::string content;
+ std::vector<int32_t> expected;
+ // Equal-width, distinct IDs retain deterministic byte boundaries while
detecting duplicates.
+ for (int32_t id = 100; id < 503; ++id) {
+ content += "{\"id\":" + std::to_string(id) + "}" + delimiter;
+ expected.push_back(id);
+ }
+ const auto path = write_json_file("four_ranges.json", content);
+ const auto file_size = static_cast<int64_t>(content.size());
+ const int64_t bytes_per_range = file_size / 4 + 1;
+ std::vector<int32_t> ids;
+ for (int64_t start = 0; start < file_size; start += bytes_per_range) {
+ ASSERT_NO_FATAL_FAILURE(read_range(path, delimiter, start,
+ std::min(bytes_per_range, file_size
- start), &ids));
+ }
+ EXPECT_EQ(ids, expected);
+}
+
+INSTANTIATE_TEST_SUITE_P(LegacyAndV2, JsonReaderSplitTest, testing::Bool(),
+ [](const testing::TestParamInfo<bool>& info) {
+ return info.param ? "V2" : "Legacy";
+ });
+
TEST(JsonReaderTest, ReadsRequestedColumnsInFileScanRequestOrder) {
ObjectPool pool;
auto slots = build_slots(&pool);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]