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, &params, 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, &params, 
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]

Reply via email to