This is an automated email from the ASF dual-hosted git repository.

liaoxin01 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 0818e61e4ed [fix](load) Avoid false quorum failures with distributed 
commit reports (#68275)
0818e61e4ed is described below

commit 0818e61e4edf570d71270271419452088220b2a2
Author: Refrain <[email protected]>
AuthorDate: Fri Oct 9 14:44:54 2026 +0800

    [fix](load) Avoid false quorum failures with distributed commit reports 
(#68275)
    
    ### What problem does this PR solve?
    
    Issue Number: N/A
    
    Related PR: apache/doris#60953
    
    Problem Summary: With enable_memtable_on_sink_node enabled, destinations
    report final tablet results to their last closing streams, which can
    belong
    to different source BEs. A source receiving one failed replica and no
    local
    success records incorrectly rejects an INSERT even when two replicas
    succeed.
    Count distinct failed backend IDs together with version-gap backends and
    reject only when known failures exceed the tolerated count. Continue
    sending
    successful commit records to FE for its aggregate quorum check.
    
    Add table-scoped fault injection to force results to different sources,
    three unit tests, and an INSERT regression with generated expected
    output.
    On master, save session variables through SELECT @@ to support
    experimental
    variables. Keep test tables and compare row counts and both set
    differences.
    
    ### Release note
    
    Fix false INSERT quorum failures when replica results are reported to
    different source BEs while preserving version-gap handling.
---
 be/src/exec/sink/writer/vtablet_writer_v2.cpp      |  36 ++++----
 be/src/load/channel/load_stream.cpp                |  62 +++++++++++++
 be/src/load/channel/load_stream.h                  |   2 +
 be/test/exec/sink/vtablet_writer_v2_test.cpp       |  74 +++++++++++++++
 .../test_insert_quorum_split_reports.out           |   6 ++
 .../test_insert_quorum_split_reports.groovy        | 101 +++++++++++++++++++++
 6 files changed, 265 insertions(+), 16 deletions(-)

diff --git a/be/src/exec/sink/writer/vtablet_writer_v2.cpp 
b/be/src/exec/sink/writer/vtablet_writer_v2.cpp
index 8a5fe58500a..13a38943d09 100644
--- a/be/src/exec/sink/writer/vtablet_writer_v2.cpp
+++ b/be/src/exec/sink/writer/vtablet_writer_v2.cpp
@@ -30,6 +30,7 @@
 #include <ranges>
 #include <string>
 #include <unordered_map>
+#include <utility>
 
 #include "common/compiler_util.h" // IWYU pragma: keep
 #include "common/logging.h"
@@ -1072,15 +1073,15 @@ void VTabletWriterV2::_calc_tablets_to_commit() {
 
 Status VTabletWriterV2::_create_commit_info(std::vector<TTabletCommitInfo>& 
tablet_commit_infos,
                                             std::shared_ptr<LoadStreamMap> 
load_stream_map) {
-    // Track per-tablet non-gap success count and failure reasons
-    std::unordered_map<int64_t, int> success_tablets_replica;
-    std::unordered_set<int64_t> failed_tablets;
+    // Commit results may be reported to different sources. Only reject a 
tablet when
+    // known failures make quorum impossible; FE checks the aggregated commit 
info.
+    std::unordered_map<int64_t, std::unordered_set<int64_t>> failed_tablets;
     std::unordered_map<int64_t, Status> failed_reason;
     load_stream_map->for_each([&](int64_t dst_id, LoadStreamStubs& streams) {
         size_t num_success_tablets = 0;
         size_t num_failed_tablets = 0;
         for (auto [tablet_id, reason] : streams.failed_tablets()) {
-            failed_tablets.insert(tablet_id);
+            failed_tablets[tablet_id].insert(dst_id);
             failed_reason[tablet_id] = reason;
             num_failed_tablets++;
         }
@@ -1089,25 +1090,28 @@ Status 
VTabletWriterV2::_create_commit_info(std::vector<TTabletCommitInfo>& tabl
             commit_info.tabletId = tablet_id;
             commit_info.backendId = dst_id;
             tablet_commit_infos.emplace_back(std::move(commit_info));
-            // Only count non-gap backends toward success
-            auto gap_it = _tablet_version_gap_backends.find(tablet_id);
-            if (gap_it == _tablet_version_gap_backends.end() ||
-                gap_it->second.find(dst_id) == gap_it->second.end()) {
-                success_tablets_replica[tablet_id]++;
-            }
             num_success_tablets++;
         }
         LOG(INFO) << "streams to dst_id: " << dst_id << ", success tablets: " 
<< num_success_tablets
                   << ", failed tablets: " << num_failed_tablets;
     });
 
-    for (auto tablet_id : failed_tablets) {
-        int succ_count = success_tablets_replica[tablet_id];
-        int required = _load_required_replicas_num(tablet_id);
-        if (succ_count < required) {
+    for (auto& [tablet_id, failed_backends] : failed_tablets) {
+        // Version-gap replicas cannot contribute to quorum, even if this 
write succeeds.
+        // Count a backend only once when it also reported a write failure.
+        if (auto gap_it = _tablet_version_gap_backends.find(tablet_id);
+            gap_it != _tablet_version_gap_backends.end()) {
+            failed_backends.insert(gap_it->second.begin(), 
gap_it->second.end());
+        }
+        auto [total_replicas_num, load_required_replicas_num] = 
_tablet_replica_info[tablet_id];
+        int max_failed_replicas = total_replicas_num == 0
+                                          ? (_num_replicas - 1) / 2
+                                          : total_replicas_num - 
load_required_replicas_num;
+        if (std::cmp_greater(failed_backends.size(), max_failed_replicas)) {
             LOG(INFO) << "tablet " << tablet_id
-                      << " failed on majority backends (success=" << succ_count
-                      << ", required=" << required << "): " << 
failed_reason[tablet_id];
+                      << " failed on majority backends (failed=" << 
failed_backends.size()
+                      << ", max_failed=" << max_failed_replicas
+                      << "): " << failed_reason[tablet_id];
             return Status::InternalError("tablet {} failed on majority 
backends: {}", tablet_id,
                                          failed_reason[tablet_id]);
         }
diff --git a/be/src/load/channel/load_stream.cpp 
b/be/src/load/channel/load_stream.cpp
index 64c8beb0a2a..c3665d3abfe 100644
--- a/be/src/load/channel/load_stream.cpp
+++ b/be/src/load/channel/load_stream.cpp
@@ -792,12 +792,74 @@ void LoadStream::_dispatch(StreamId id, const 
PStreamHeader& hdr, butil::IOBuf*
     } break;
     case PStreamHeader::CLOSE_LOAD: {
         DBUG_EXECUTE_IF("LoadStream.close_load.block", DBUG_BLOCK);
+        DBUG_EXECUTE_IF("LoadStream.close_load.force_last_source", {
+            if (_schema->table_id() == dp->param<int64_t>("table_id", -1)) {
+                MonotonicStopWatch wait_timer;
+                wait_timer.start();
+                while (true) {
+                    bool ready = false;
+                    bool single_source = false;
+                    {
+                        std::lock_guard lock_guard(_lock);
+                        if (_debug_last_close_src_id < 0) {
+                            int opened_streams = _close_load_cnt;
+                            for (const auto& [_, count] : _open_streams) {
+                                opened_streams += count;
+                            }
+                            // Wait for every source to open before selecting 
a stable
+                            // min/max source ID. No CLOSE_LOAD is counted 
before selection.
+                            if (opened_streams == _total_streams) {
+                                single_source = _open_streams.size() < 2;
+                                if (!single_source) {
+                                    const bool pick_max = 
dp->param<bool>("pick_max", true);
+                                    _debug_last_close_src_id = 
_open_streams.begin()->first;
+                                    for (const auto& [src_id, _] : 
_open_streams) {
+                                        if ((pick_max && src_id > 
_debug_last_close_src_id) ||
+                                            (!pick_max && src_id < 
_debug_last_close_src_id)) {
+                                            _debug_last_close_src_id = src_id;
+                                        }
+                                    }
+                                }
+                            }
+                        }
+                        ready = _debug_last_close_src_id >= 0 &&
+                                (hdr.src_id() != _debug_last_close_src_id ||
+                                 _open_streams.size() == 1);
+                    }
+                    if (ready) {
+                        break;
+                    }
+                    if (single_source ||
+                        wait_timer.elapsed_time() / 1000000 >=
+                                dp->param<int64_t>("wait_ms", 30000) ||
+                        !DebugPoints::instance()->is_enable(DP_NAME)) {
+                        LOG(WARNING) << "cannot force last CLOSE_LOAD source 
(requires multiple "
+                                        "sources), "
+                                     << *this;
+                        // Closing without EOS makes the load fail instead of 
silently
+                        // passing a test that did not establish the requested 
ordering.
+                        brpc::StreamClose(id);
+                        return;
+                    }
+                    bthread_usleep(1000);
+                }
+            }
+        });
         std::vector<int64_t> success_tablet_ids;
         FailedTablets failed_tablets;
         std::vector<PTabletID> tablets_to_commit(hdr.tablets().begin(), 
hdr.tablets().end());
         // Step 1: count this CLOSE_LOAD and, if this is the last one, commit. 
Under _lock.
         bool all_received =
                 close(hdr.src_id(), tablets_to_commit, &success_tablet_ids, 
&failed_tablets);
+        DBUG_EXECUTE_IF("LoadStream.close_load.force_last_source", {
+            if (all_received && _schema->table_id() == 
dp->param<int64_t>("table_id", -1) &&
+                dp->param<bool>("require_failure", false) && 
failed_tablets.empty()) {
+                LOG(WARNING) << "expected a failed tablet in the final 
CLOSE_LOAD result, "
+                             << *this;
+                brpc::StreamClose(id);
+                return;
+            }
+        });
         // Step 2: send THIS stream's EOS (network IO, must be outside _lock). 
A stream
         // must not be StreamClose'd before its own EOS is delivered, 
otherwise the
         // sender sees on_closed without EOS and reports "Stream closed 
without EOS".
diff --git a/be/src/load/channel/load_stream.h 
b/be/src/load/channel/load_stream.h
index 4c8865e1603..3e883a39c69 100644
--- a/be/src/load/channel/load_stream.h
+++ b/be/src/load/channel/load_stream.h
@@ -186,6 +186,8 @@ private:
     std::unordered_map<int64_t, IndexStreamSharedPtr> _index_streams_map;
     int32_t _total_streams = 0;
     int32_t _close_load_cnt = 0;
+    // Only used by close_load.force_last_source, protected by _lock.
+    int64_t _debug_last_close_src_id = -1;
     std::atomic<int32_t> _close_rpc_cnt = 0;
     std::vector<PTabletID> _tablets_to_commit;
     bthread::Mutex _lock;
diff --git a/be/test/exec/sink/vtablet_writer_v2_test.cpp 
b/be/test/exec/sink/vtablet_writer_v2_test.cpp
index 759ca06b51f..c2d43c3995a 100644
--- a/be/test/exec/sink/vtablet_writer_v2_test.cpp
+++ b/be/test/exec/sink/vtablet_writer_v2_test.cpp
@@ -414,6 +414,80 @@ TEST_F(TestVTabletWriterV2, fail_one) {
     ASSERT_EQ(tablet_commit_infos.size(), 5);
 }
 
+TEST_F(TestVTabletWriterV2, commit_info_with_results_split_across_sources) {
+    UniqueId load_id;
+    auto first_source = std::make_shared<LoadStreamMap>(load_id, src_id, 1, 1, 
nullptr);
+    auto second_source = std::make_shared<LoadStreamMap>(load_id, src_id + 1, 
1, 1, nullptr);
+    add_stream(first_source, 1001, {1}, {});
+    add_stream(first_source, 1002, {1}, {});
+    add_stream(first_source, 1003, {}, {{1, Status::InternalError("write 
failed")}});
+    add_stream(second_source, 1001, {}, {});
+    add_stream(second_source, 1002, {}, {});
+    add_stream(second_source, 1003, {}, {{1, Status::InternalError("write 
failed")}});
+
+    auto first_writer = create_vtablet_writer();
+    auto second_writer = create_vtablet_writer();
+    std::vector<TTabletCommitInfo> first_commit_infos;
+    std::vector<TTabletCommitInfo> second_commit_infos;
+    ASSERT_TRUE(first_writer->_create_commit_info(first_commit_infos, 
first_source).ok());
+    // The second source received the failed replica's final result, but the 
two
+    // successful replicas reported to the first source. It must not reject 
the load.
+    ASSERT_TRUE(second_writer->_create_commit_info(second_commit_infos, 
second_source).ok());
+    ASSERT_EQ(first_commit_infos.size(), 2);
+    ASSERT_TRUE(second_commit_infos.empty());
+    for (const auto& info : first_commit_infos) {
+        EXPECT_EQ(info.tabletId, 1);
+        EXPECT_TRUE(info.backendId == 1001 || info.backendId == 1002);
+    }
+}
+
+TEST_F(TestVTabletWriterV2, 
write_failure_and_version_gap_exceed_failure_quorum) {
+    for (bool gap_replica_reported_success : {false, true}) {
+        SCOPED_TRACE(gap_replica_reported_success);
+        UniqueId load_id;
+        auto load_stream_map = std::make_shared<LoadStreamMap>(load_id, 
src_id, 1, 1, nullptr);
+        add_stream(load_stream_map, 1001, {1}, {});
+        add_stream(load_stream_map, 1002, {}, {{1, 
Status::InternalError("write failed")}});
+        add_stream(
+                load_stream_map, 1003,
+                gap_replica_reported_success ? std::vector<int64_t> {1} : 
std::vector<int64_t> {},
+                {});
+
+        auto writer = create_vtablet_writer();
+        writer->_tablet_version_gap_backends[1].insert(1003);
+        std::vector<TTabletCommitInfo> commit_infos;
+        auto st = writer->_create_commit_info(commit_infos, load_stream_map);
+        // A version-gap replica is unavailable regardless of where its result 
is
+        // reported. Together with the write failure, only one valid replica 
remains.
+        ASSERT_FALSE(st.ok());
+        EXPECT_NE(st.to_string().find("failed on majority backends"), 
std::string::npos);
+        EXPECT_NE(st.to_string().find("write failed"), std::string::npos);
+        EXPECT_EQ(commit_infos.size(), gap_replica_reported_success ? 2 : 1);
+    }
+}
+
+TEST_F(TestVTabletWriterV2, 
version_gap_and_duplicate_write_failure_count_once) {
+    UniqueId load_id;
+    auto load_stream_map = std::make_shared<LoadStreamMap>(load_id, src_id, 2, 
1, nullptr);
+    add_stream(load_stream_map, 1001, {1}, {});
+    auto failed_streams = load_stream_map->get_or_create(1002);
+    failed_streams->mark_open();
+    for (const auto& stream : failed_streams->streams()) {
+        stream->add_failed_tablet(1, Status::InternalError("write failed"));
+    }
+    add_stream(load_stream_map, 1003, {}, {});
+
+    auto writer = create_vtablet_writer();
+    writer->_tablet_version_gap_backends[1].insert(1002);
+    std::vector<TTabletCommitInfo> commit_infos;
+    // Only backend 1002 is unavailable. Backend 1003 may have reported 
success to
+    // another source, so a missing local success must not cause a quorum 
failure.
+    ASSERT_TRUE(writer->_create_commit_info(commit_infos, 
load_stream_map).ok());
+    ASSERT_EQ(commit_infos.size(), 1);
+    EXPECT_EQ(commit_infos[0].tabletId, 1);
+    EXPECT_EQ(commit_infos[0].backendId, 1001);
+}
+
 TEST_F(TestVTabletWriterV2, fail_one_duplicate) {
     UniqueId load_id;
     std::vector<TTabletCommitInfo> tablet_commit_infos;
diff --git 
a/regression-test/data/fault_injection_p0/test_insert_quorum_split_reports.out 
b/regression-test/data/fault_injection_p0/test_insert_quorum_split_reports.out
new file mode 100644
index 00000000000..a6a02a18cee
--- /dev/null
+++ 
b/regression-test/data/fault_injection_p0/test_insert_quorum_split_reports.out
@@ -0,0 +1,6 @@
+-- This file is automatically generated. You should know what you did if you 
want to edit this
+-- !rows --
+2048
+
+-- !mismatches --
+0
diff --git 
a/regression-test/suites/fault_injection_p0/test_insert_quorum_split_reports.groovy
 
b/regression-test/suites/fault_injection_p0/test_insert_quorum_split_reports.groovy
new file mode 100644
index 00000000000..0b6e80408fa
--- /dev/null
+++ 
b/regression-test/suites/fault_injection_p0/test_insert_quorum_split_reports.groovy
@@ -0,0 +1,101 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+import org.apache.doris.regression.util.NodeType
+
+suite("test_insert_quorum_split_reports", "nonConcurrent") {
+    if (isCloudMode()) {
+        return
+    }
+    def backendIPs = [:]
+    def backendPorts = [:]
+    getBackendIpHttpPort(backendIPs, backendPorts)
+    if (backendIPs.size() < 3) {
+        return
+    }
+
+    def closePoint = "LoadStream.close_load.force_last_source"
+    def failurePoint = "TabletStream.add_segment.unknown_segid"
+    def debugPoint = GetDebugPoint()
+    def failedBackend = null
+    def sessionVariables = ["enable_memtable_on_sink_node", 
"enable_local_shuffle",
+            "parallel_pipeline_task_num", "load_stream_per_node", 
"query_timeout"]
+    def savedVariables = sessionVariables.collectEntries { name ->
+        [(name): sql("select @@${name}")[0][0]]
+    }
+    try {
+        sql "set enable_memtable_on_sink_node = true"
+        sql "set enable_local_shuffle = false"
+        sql "set parallel_pipeline_task_num = 1"
+        sql "set load_stream_per_node = 1"
+        sql "set query_timeout = 60"
+        sql "drop table if exists insert_quorum_split_reports_source"
+        sql "drop table if exists insert_quorum_split_reports_target"
+        // Spread the scan over multiple BEs. The close injection fails the 
load
+        // explicitly if execution nevertheless opens streams from only one 
source.
+        sql """create table insert_quorum_split_reports_source (k int, v 
bigint)
+                duplicate key(k) distributed by hash(k) buckets 
${backendIPs.size() * 4}
+                properties("replication_num" = "1")"""
+        sql """create table insert_quorum_split_reports_target (k int, v 
bigint)
+                duplicate key(k) distributed by hash(k) buckets 5
+                properties("replication_num" = "3")"""
+        sql """insert into insert_quorum_split_reports_source
+                select number, number * 11 + 5 from numbers("number" = 
"2048")"""
+
+        def replicas = sql_return_maparray "show tablets from 
insert_quorum_split_reports_target"
+        failedBackend = replicas.collect { it.BackendId.toString() 
}.unique().min { it.toLong() }
+        def tableId = 
getTableId("insert_quorum_split_reports_target").toString()
+        // All healthy destinations report to the largest source ID. The failed
+        // destination reports to the smallest, which receives no healthy 
results.
+        debugPoint.enableDebugPointForAllBEs(closePoint,
+                [table_id: tableId, pick_max: "true", wait_ms: "30000", 
timeout: "90"])
+        debugPoint.enableDebugPoint(backendIPs[failedBackend], 
backendPorts[failedBackend] as int,
+                NodeType.BE, closePoint,
+                [table_id: tableId, pick_max: "false", require_failure: "true",
+                 wait_ms: "30000", timeout: "90"])
+        debugPoint.enableDebugPoint(backendIPs[failedBackend], 
backendPorts[failedBackend] as int,
+                NodeType.BE, failurePoint, [timeout: "90"])
+
+        // Before the fix, the source receiving the failed destination's final
+        // report rejects the INSERT with success=0, required=2.
+        sql "insert into insert_quorum_split_reports_target select * from 
insert_quorum_split_reports_source"
+
+        // require_failure checks the actual final response, avoiding a race 
with
+        // background replica repair when checking the FE's failed-version 
metadata.
+        qt_rows "select count(*) from insert_quorum_split_reports_target"
+        qt_mismatches """select count(*) from (
+                (select k, v from insert_quorum_split_reports_source
+                 except select k, v from insert_quorum_split_reports_target)
+                union all
+                (select k, v from insert_quorum_split_reports_target
+                 except select k, v from insert_quorum_split_reports_source)
+                ) differences"""
+    } finally {
+        try {
+            if (failedBackend != null) {
+                debugPoint.disableDebugPoint(backendIPs[failedBackend], 
backendPorts[failedBackend] as int,
+                        NodeType.BE, failurePoint)
+            }
+        } finally {
+            try {
+                debugPoint.disableDebugPointForAllBEs(closePoint)
+            } finally {
+                savedVariables.each { name, value -> sql "set ${name} = 
'${value}'" }
+            }
+        }
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to