This is an automated email from the ASF dual-hosted git repository.
hello-stephen 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 db0b39a4f1a [fix](be) Make replica fault injection deterministic
(#65627)
db0b39a4f1a is described below
commit db0b39a4f1ab56811b9be9662d80d7c65435ab45
Author: shuke <[email protected]>
AuthorDate: Thu Jul 16 21:04:07 2026 +0800
[fix](be) Make replica fault injection deterministic (#65627)
Related PR: #47082
Problem Summary:
`StreamSinkFileWriter` previously implemented the one/two-replica
fault-injection debug points by skipping the first entries in each
writer's local stream order. Writers for the same tablet can have
different stream orders, so they could skip different destination
backends. A nominal one-replica failure could therefore affect two
replicas across writers and make the load lose quorum.
This PR derives the failed replica set from sorted destination backend
IDs. Every writer with the same replica set now skips the same
backend(s), independent of local stream order. It also adds unit
coverage for one- and two-replica injection with reordered stream lists.
---
be/src/io/fs/stream_sink_file_writer.cpp | 68 ++++++++++++++------------
be/src/io/fs/stream_sink_file_writer.h | 2 +
be/test/io/fs/stream_sink_file_writer_test.cpp | 56 +++++++++++++++++++--
3 files changed, 92 insertions(+), 34 deletions(-)
diff --git a/be/src/io/fs/stream_sink_file_writer.cpp
b/be/src/io/fs/stream_sink_file_writer.cpp
index 1316cd71fac..e4a4a4d932e 100644
--- a/be/src/io/fs/stream_sink_file_writer.cpp
+++ b/be/src/io/fs/stream_sink_file_writer.cpp
@@ -19,6 +19,8 @@
#include <gen_cpp/internal_service.pb.h>
+#include <algorithm>
+
#include "exec/sink/load_stream_stub.h"
#include "storage/olap_common.h"
#include "storage/rowset/beta_rowset_writer.h"
@@ -27,6 +29,28 @@
namespace doris::io {
+std::unordered_set<int64_t>
StreamSinkFileWriter::_get_fault_injection_failed_dst_ids() const {
+ size_t failed_replica_num = 0;
+
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_one_replica",
+ { failed_replica_num = 1; });
+
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_two_replica",
+ { failed_replica_num = 2; });
+
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_all_replica",
+ { failed_replica_num = _streams.size(); });
+ if (failed_replica_num == 0) {
+ return {};
+ }
+
+ std::vector<int64_t> dst_ids;
+ dst_ids.reserve(_streams.size());
+ for (const auto& stream : _streams) {
+ dst_ids.push_back(stream->dst_id());
+ }
+ std::sort(dst_ids.begin(), dst_ids.end());
+ dst_ids.resize(std::min(failed_replica_num, dst_ids.size()));
+ return std::unordered_set<int64_t>(dst_ids.begin(), dst_ids.end());
+}
+
void StreamSinkFileWriter::init(PUniqueId load_id, int64_t partition_id,
int64_t index_id,
int64_t tablet_id, int32_t segment_id,
FileType file_type) {
VLOG_DEBUG << "init stream writer, load id(" <<
UniqueId(load_id).to_string()
@@ -52,24 +76,16 @@ Status StreamSinkFileWriter::appendv(const Slice* data,
size_t data_cnt) {
<< ", data_length: " << bytes_req << "file_type" << _file_type;
std::span<const Slice> slices {data, data_cnt};
- size_t fault_injection_skipped_streams = 0;
+ auto fault_injection_failed_dst_ids =
_get_fault_injection_failed_dst_ids();
bool ok = false;
Status st;
for (auto& stream : _streams) {
-
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_one_replica",
{
- if (fault_injection_skipped_streams < 1) {
- fault_injection_skipped_streams++;
- continue;
- }
- });
-
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_two_replica",
{
- if (fault_injection_skipped_streams < 2) {
- fault_injection_skipped_streams++;
- continue;
- }
- });
-
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_all_replica",
- { continue; });
+ if (fault_injection_failed_dst_ids.contains(stream->dst_id())) {
+ LOG(INFO) << "fault injection skips segment data to backend " <<
stream->dst_id()
+ << ", load_id: " << print_id(_load_id) << ", index_id: "
<< _index_id
+ << ", tablet_id: " << _tablet_id << ", segment_id: " <<
_segment_id;
+ continue;
+ }
st = stream->append_data(_partition_id, _index_id, _tablet_id,
_segment_id, _bytes_appended,
slices, false, _file_type);
ok = ok || st.ok();
@@ -123,23 +139,15 @@ Status StreamSinkFileWriter::_finalize() {
VLOG_DEBUG << "writer finalize, load_id: " << print_id(_load_id) << ",
index_id: " << _index_id
<< ", tablet_id: " << _tablet_id << ", segment_id: " <<
_segment_id;
// TODO(zhengyu): update get_inverted_index_file_size into stat
- size_t fault_injection_skipped_streams = 0;
+ auto fault_injection_failed_dst_ids =
_get_fault_injection_failed_dst_ids();
bool ok = false;
for (auto& stream : _streams) {
-
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_one_replica",
{
- if (fault_injection_skipped_streams < 1) {
- fault_injection_skipped_streams++;
- continue;
- }
- });
-
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_two_replica",
{
- if (fault_injection_skipped_streams < 2) {
- fault_injection_skipped_streams++;
- continue;
- }
- });
-
DBUG_EXECUTE_IF("StreamSinkFileWriter.appendv.write_segment_failed_all_replica",
- { continue; });
+ if (fault_injection_failed_dst_ids.contains(stream->dst_id())) {
+ LOG(INFO) << "fault injection skips segment eos to backend " <<
stream->dst_id()
+ << ", load_id: " << print_id(_load_id) << ", index_id: "
<< _index_id
+ << ", tablet_id: " << _tablet_id << ", segment_id: " <<
_segment_id;
+ continue;
+ }
auto st = stream->append_data(_partition_id, _index_id, _tablet_id,
_segment_id,
_bytes_appended, {}, true, _file_type);
ok = ok || st.ok();
diff --git a/be/src/io/fs/stream_sink_file_writer.h
b/be/src/io/fs/stream_sink_file_writer.h
index f092319e7fa..74f7d4e091c 100644
--- a/be/src/io/fs/stream_sink_file_writer.h
+++ b/be/src/io/fs/stream_sink_file_writer.h
@@ -21,6 +21,7 @@
#include <gen_cpp/olap_common.pb.h>
#include <queue>
+#include <unordered_set>
#include "io/fs/file_writer.h"
#include "util/uid_util.h"
@@ -57,6 +58,7 @@ public:
private:
Status _finalize();
+ std::unordered_set<int64_t> _get_fault_injection_failed_dst_ids() const;
std::vector<std::shared_ptr<LoadStreamStub>> _streams;
PUniqueId _load_id;
diff --git a/be/test/io/fs/stream_sink_file_writer_test.cpp
b/be/test/io/fs/stream_sink_file_writer_test.cpp
index 11de185d8e2..35d9fb1ef68 100644
--- a/be/test/io/fs/stream_sink_file_writer_test.cpp
+++ b/be/test/io/fs/stream_sink_file_writer_test.cpp
@@ -20,10 +20,15 @@
#include <brpc/channel.h>
#include <brpc/server.h>
+#include <array>
+#include <unordered_map>
+
+#include "common/config.h"
#include "exec/sink/load_stream_stub.h"
#include "gtest/gtest_pred_impl.h"
#include "storage/olap_common.h"
#include "util/debug/leakcheck_disabler.h"
+#include "util/debug_points.h"
#include "util/faststring.h"
namespace doris {
@@ -47,13 +52,16 @@ const std::string DATA0 = "segment data";
const std::string DATA1 = "hello world";
static std::atomic<int64_t> g_num_request;
+static std::unordered_map<int64_t, int64_t> g_num_requests_by_dst;
class StreamSinkFileWriterTest : public testing::Test {
class MockStreamStub : public LoadStreamStub {
public:
- MockStreamStub(PUniqueId load_id, int64_t src_id)
+ MockStreamStub(PUniqueId load_id, int64_t src_id, int64_t dst_id)
: LoadStreamStub(load_id, src_id,
std::make_shared<IndexToTabletSchema>(),
- std::make_shared<IndexToEnableMoW>()) {};
+ std::make_shared<IndexToEnableMoW>()) {
+ _dst_id = dst_id;
+ }
virtual ~MockStreamStub() = default;
@@ -76,6 +84,7 @@ class StreamSinkFileWriterTest : public testing::Test {
EXPECT_EQ(0, offset);
}
g_num_request++;
+ g_num_requests_by_dst[_dst_id]++;
return Status::OK();
}
};
@@ -88,15 +97,30 @@ protected:
virtual void SetUp() {
_load_id.set_hi(LOAD_ID_HI);
_load_id.set_lo(LOAD_ID_LO);
+ std::array<int64_t, NUM_STREAM> dst_ids {103, 101, 102};
for (int src_id = 0; src_id < NUM_STREAM; src_id++) {
- _streams.emplace_back(new MockStreamStub(_load_id, src_id));
+ _streams.emplace_back(new MockStreamStub(_load_id, src_id,
dst_ids[src_id]));
}
+ _enable_debug_points = config::enable_debug_points;
+ g_num_requests_by_dst.clear();
}
- virtual void TearDown() {}
+ virtual void TearDown() {
+ DebugPoints::instance()->clear();
+ config::enable_debug_points = _enable_debug_points;
+ }
+
+ void write_with_streams(const
std::vector<std::shared_ptr<LoadStreamStub>>& streams) {
+ io::StreamSinkFileWriter writer(streams);
+ writer.init(_load_id, PARTITION_ID, INDEX_ID, TABLET_ID, SEGMENT_ID);
+ std::vector<Slice> slices {DATA0, DATA1};
+ CHECK_STATUS_OK(writer.appendv(&(*slices.begin()), slices.size()));
+ CHECK_STATUS_OK(writer.close());
+ }
PUniqueId _load_id;
std::vector<std::shared_ptr<LoadStreamStub>> _streams;
+ bool _enable_debug_points;
};
TEST_F(StreamSinkFileWriterTest, Test) {
@@ -111,4 +135,28 @@ TEST_F(StreamSinkFileWriterTest, Test) {
EXPECT_EQ(NUM_STREAM * 2, g_num_request);
}
+TEST_F(StreamSinkFileWriterTest, DeterministicOneReplicaFaultInjection) {
+ config::enable_debug_points = true;
+
DebugPoints::instance()->add("StreamSinkFileWriter.appendv.write_segment_failed_one_replica");
+
+ write_with_streams(_streams);
+ write_with_streams({_streams[2], _streams[0], _streams[1]});
+
+ EXPECT_EQ(0, g_num_requests_by_dst[101]);
+ EXPECT_EQ(4, g_num_requests_by_dst[102]);
+ EXPECT_EQ(4, g_num_requests_by_dst[103]);
+}
+
+TEST_F(StreamSinkFileWriterTest, DeterministicTwoReplicaFaultInjection) {
+ config::enable_debug_points = true;
+
DebugPoints::instance()->add("StreamSinkFileWriter.appendv.write_segment_failed_two_replica");
+
+ write_with_streams(_streams);
+ write_with_streams({_streams[2], _streams[0], _streams[1]});
+
+ EXPECT_EQ(0, g_num_requests_by_dst[101]);
+ EXPECT_EQ(0, g_num_requests_by_dst[102]);
+ EXPECT_EQ(4, g_num_requests_by_dst[103]);
+}
+
} // namespace doris
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]