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 917e272c0e4 branch-4.1: [fix](be) Make replica fault injection
deterministic #65627 (#65727)
917e272c0e4 is described below
commit 917e272c0e4e512e1ac4bbaee549a850f6949c5c
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Fri Jul 17 08:35:59 2026 +0800
branch-4.1: [fix](be) Make replica fault injection deterministic #65627
(#65727)
Cherry-picked from #65627
Co-authored-by: shuke <[email protected]>
---
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]