This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 7b50e05cdbe branch-4.1: [fix](scan) Fix lost and duplicated rows when
splitting CSV/JSON on multi-character line delimiters #68539 (#68600)
7b50e05cdbe is described below
commit 7b50e05cdbe352207b7f99084d3e0085eee3e168
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Tue Sep 29 23:39:51 2026 +0800
branch-4.1: [fix](scan) Fix lost and duplicated rows when splitting
CSV/JSON on multi-character line delimiters #68539 (#68600)
Cherry-picked from #68539
Co-authored-by: hui lai <[email protected]>
---
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 6518ab42459..250307b3348 100644
--- a/be/src/format/csv/csv_reader.cpp
+++ b/be/src/format/csv/csv_reader.cpp
@@ -345,6 +345,14 @@ Status CsvReader::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 31cb50f4096..4dec574df88 100644
--- a/be/src/format/json/new_json_reader.cpp
+++ b/be/src/format/json/new_json_reader.cpp
@@ -152,6 +152,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
@@ -211,9 +213,8 @@ Status NewJsonReader::get_next_block(Block* block, size_t*
read_rows, bool* eof)
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;
}
@@ -445,7 +446,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;
@@ -482,8 +485,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 15ac8f14c41..1179591ad44 100644
--- a/be/src/format/json/new_json_reader.h
+++ b/be/src/format/json/new_json_reader.h
@@ -62,6 +62,7 @@ struct IOContext;
struct ScannerCounter;
class Block;
class IColumn;
+class NewPlainTextLineReader;
namespace json_reader_detail {
Status append_null_for_malformed_json(Block& block);
@@ -83,7 +84,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,
@@ -200,7 +201,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 bb6c55d3084..e5eca54df6f 100644
--- a/be/src/format_v2/delimited_text/csv_reader.cpp
+++ b/be/src/format_v2/delimited_text/csv_reader.cpp
@@ -128,11 +128,13 @@ 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 (_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 e8d9aec40af..1d94e244f64 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]