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]