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]

Reply via email to