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 9c86638f629 [fix](topn) Handle empty rowid fetch RPC failures (#66443)
9c86638f629 is described below

commit 9c86638f629a0b8017d62aa5b83d66fb3110d9e8
Author: HappenLee <[email protected]>
AuthorDate: Thu Aug 6 11:36:57 2026 +0800

    [fix](topn) Handle empty rowid fetch RPC failures (#66443)
    
    Problem Summary: TopN lazy materialization sends reusable rowid-fetch
    RPC requests to every target backend. A backend request can contain no
    row IDs for a batch, but a transport failure for that no-op request
    previously aborted the query. If such a failure left no response blocks,
    response merging could also read a missing or stale response. Ignore
    transport failures only when the request has no rows, clear reused
    responses before dispatch, and skip deserialization for empty
    per-relation requests. Requests containing rows and all backend-reported
    errors still fail.
    
    ### Release note
    
    Fix TopN lazy materialization queries that fail when an empty
    rowid-fetch RPC encounters a transport error.
---
 be/src/exec/operator/materialization_opertor.cpp   | 57 ++++++++++++++--------
 be/src/exec/operator/materialization_opertor.h     |  1 +
 .../operator/materialization_shared_state_test.cpp | 41 ++++++++++++++--
 3 files changed, 75 insertions(+), 24 deletions(-)

diff --git a/be/src/exec/operator/materialization_opertor.cpp 
b/be/src/exec/operator/materialization_opertor.cpp
index 7ac146aa75d..5628faceecc 100644
--- a/be/src/exec/operator/materialization_opertor.cpp
+++ b/be/src/exec/operator/materialization_opertor.cpp
@@ -280,6 +280,12 @@ Status 
MaterializationSharedState::merge_multi_response(RuntimeProfile* profile)
         // Phase 1: Deserialize the i-th response block from every BE into 
block_maps.
         // Each BE's response.blocks(i) corresponds to the i-th relation's 
fetched columns.
         for (auto& [backend_id, rpc_struct] : rpc_struct_map) {
+            const auto request_row_count = 
rpc_struct.request.request_block_descs(i).row_id_size();
+            // An empty request may have no response block when its RPC fails. 
It cannot
+            // contribute rows to block_order_results, so there is nothing to 
deserialize.
+            if (request_row_count == 0) {
+                continue;
+            }
             Block partial_block;
             size_t uncompressed_size = 0;
             int64_t uncompressed_time = 0;
@@ -293,13 +299,12 @@ Status 
MaterializationSharedState::merge_multi_response(RuntimeProfile* profile)
             // refer 'if (!id_file_map)' in RowIdStorageReader::read_by_rowids.
             // 2. Report error in any case where the row count doesn't match, 
even if it's not empty,
             //    since that indicates a bug in BE's row fetching logic or 
serialization logic.
-            if (rpc_struct.request.request_block_descs(i).row_id_size() != 
partial_block.rows()) {
+            if (request_row_count != partial_block.rows()) {
                 return Status::InternalError(
                         fmt::format("merge_multi_response, "
                                     "backend_id {} returned block with row 
count {} not match "
                                     "request row id count {}",
-                                    backend_id, partial_block.rows(),
-                                    
rpc_struct.request.request_block_descs(i).row_id_size()));
+                                    backend_id, partial_block.rows(), 
request_row_count));
             }
             if (rpc_struct.response.blocks(i).has_profile()) {
                 auto response_profile =
@@ -307,8 +312,7 @@ Status 
MaterializationSharedState::merge_multi_response(RuntimeProfile* profile)
                 _update_profile_info(backend_id, response_profile.get());
             }
 
-            // Only insert non-empty blocks. A BE may return an empty block if
-            // request.request_block_descs(i).row_id_size() is 0
+            // Only insert non-empty blocks.
             if (!partial_block.is_empty_column()) {
                 // Reset row cursor to 0 — we'll consume rows from this block 
sequentially.
                 block_maps[backend_id] = 
std::make_pair(std::move(partial_block), 0);
@@ -369,6 +373,29 @@ Status 
MaterializationSharedState::merge_multi_response(RuntimeProfile* profile)
     return Status::OK();
 }
 
+Status MaterializationSharedState::validate_rpc_results(int node_id) {
+    for (auto& [backend_id, rpc_struct] : rpc_struct_map) {
+        if (rpc_struct.cntl->Failed()) {
+            if (count_request_rows(rpc_struct.request) > 0) {
+                return Status::InternalError(
+                        "Failed to send brpc request, error_text=" + 
rpc_struct.cntl->ErrorText() +
+                        " Materialization Sink node id:" + 
std::to_string(node_id) +
+                        " target_backend_id:" + std::to_string(backend_id));
+            }
+            rpc_struct.cntl->Reset();
+            continue;
+        }
+        if (rpc_struct.response.status().status_code() != 0) {
+            Status st = Status::create(rpc_struct.response.status());
+            st.append(fmt::format(", Backend:{}, Materialization Sink node 
id:{}", backend_id,
+                                  node_id));
+            return st;
+        }
+        rpc_struct.cntl->Reset();
+    }
+    return Status::OK();
+}
+
 void MaterializationSharedState::_update_profile_info(int64_t backend_id,
                                                       RuntimeProfile* 
response_profile) {
     if (!backend_profile_info_string.contains(backend_id)) {
@@ -601,6 +628,9 @@ Status MaterializationOperator::push(RuntimeState* state, 
Block* in_block, bool
         MonotonicStopWatch rpc_timer(true);
         for (auto& [backend_id, rpc_struct] : 
local_state._materialization_state.rpc_struct_map) {
             auto* callback = brpc::NewCallback(fetch_callback, &counter);
+            // The response object is reused across batches. Clear it so an 
ignored failure for
+            // an empty request cannot expose a response left by an earlier 
batch.
+            rpc_struct.response.Clear();
             rpc_struct.cntl->set_timeout_ms(state->execution_timeout() * 1000);
             // send brpc request
             rpc_struct.stub->multiget_data_v2(rpc_struct.cntl.get(), 
&rpc_struct.request,
@@ -611,22 +641,7 @@ Status MaterializationOperator::push(RuntimeState* state, 
Block* in_block, bool
             local_state._max_rpc_timer->set(time);
         }
 
-        for (auto& [backend_id, rpc_struct] : 
local_state._materialization_state.rpc_struct_map) {
-            if (rpc_struct.cntl->Failed()) {
-                std::string error_text =
-                        "Failed to send brpc request, error_text=" + 
rpc_struct.cntl->ErrorText() +
-                        " Materialization Sink node id:" + 
std::to_string(node_id()) +
-                        " target_backend_id:" + std::to_string(backend_id);
-                return Status::InternalError(error_text);
-            }
-            if (rpc_struct.response.status().status_code() != 0) {
-                Status st = Status::create(rpc_struct.response.status());
-                st.append(fmt::format(", Backend:{}, Materialization Sink node 
id:{}", backend_id,
-                                      node_id()));
-                return st;
-            }
-            rpc_struct.cntl->Reset();
-        }
+        
RETURN_IF_ERROR(local_state._materialization_state.validate_rpc_results(node_id()));
 
         if (local_state._materialization_state.need_merge_block) {
             SCOPED_TIMER(local_state._merge_response_timer);
diff --git a/be/src/exec/operator/materialization_opertor.h 
b/be/src/exec/operator/materialization_opertor.h
index 889d94a11ac..bc6618b1bb0 100644
--- a/be/src/exec/operator/materialization_opertor.h
+++ b/be/src/exec/operator/materialization_opertor.h
@@ -46,6 +46,7 @@ public:
     Status init_multi_requests(const TMaterializationNode& tnode, 
RuntimeState* state);
     Status create_muiltget_result(const Columns& columns, bool eos);
 
+    Status validate_rpc_results(int node_id);
     Status merge_multi_response(RuntimeProfile* profile);
     void get_block(Block* block);
 
diff --git a/be/test/exec/operator/materialization_shared_state_test.cpp 
b/be/test/exec/operator/materialization_shared_state_test.cpp
index b1950482a70..fb3c31493a6 100644
--- a/be/test/exec/operator/materialization_shared_state_test.cpp
+++ b/be/test/exec/operator/materialization_shared_state_test.cpp
@@ -60,6 +60,8 @@ protected:
         _shared_state->rpc_struct_map[_backend_id2] = FetchRpcStruct();
         
_shared_state->rpc_struct_map[_backend_id1].request.add_request_block_descs();
         
_shared_state->rpc_struct_map[_backend_id2].request.add_request_block_descs();
+        _shared_state->rpc_struct_map[_backend_id1].cntl = 
std::make_unique<brpc::Controller>();
+        _shared_state->rpc_struct_map[_backend_id2].cntl = 
std::make_unique<brpc::Controller>();
     }
 
     std::shared_ptr<MaterializationSharedState> _shared_state;
@@ -69,6 +71,39 @@ protected:
     int64_t _backend_id2;
 };
 
+TEST_F(MaterializationSharedStateTest, 
TestRpcFailureWithoutFetchRowsIsIgnored) {
+    auto& rpc_struct = _shared_state->rpc_struct_map[_backend_id1];
+    rpc_struct.cntl->SetFailed("injected connection failure");
+
+    Status st = _shared_state->validate_rpc_results(10);
+
+    EXPECT_TRUE(st.ok()) << st.to_string();
+    EXPECT_FALSE(rpc_struct.cntl->Failed());
+}
+
+TEST_F(MaterializationSharedStateTest, TestRpcFailureWithFetchRowsFails) {
+    auto& rpc_struct = _shared_state->rpc_struct_map[_backend_id1];
+    add_request_row(rpc_struct.request.mutable_request_block_descs(0), 0, 1);
+    rpc_struct.cntl->SetFailed("injected connection failure");
+
+    Status st = _shared_state->validate_rpc_results(10);
+
+    EXPECT_FALSE(st.ok());
+    EXPECT_NE(st.to_string().find("target_backend_id:1001"), 
std::string::npos);
+}
+
+TEST_F(MaterializationSharedStateTest, TestNonOkResponseWithoutFetchRowsFails) 
{
+    auto& rpc_struct = _shared_state->rpc_struct_map[_backend_id1];
+    Status::InternalError("injected remote failure")
+            .to_protobuf(rpc_struct.response.mutable_status());
+
+    Status st = _shared_state->validate_rpc_results(10);
+
+    EXPECT_FALSE(st.ok());
+    EXPECT_NE(st.to_string().find("injected remote failure"), 
std::string::npos);
+    EXPECT_NE(st.to_string().find("Backend:1001"), std::string::npos);
+}
+
 TEST_F(MaterializationSharedStateTest, TestCreateMultiGetResult) {
     // Create test columns for rowids
     Columns columns;
@@ -453,7 +488,7 @@ TEST_F(MaterializationSharedStateTest, 
TestMergeMultiResponseStaleBlockMaps) {
     
_shared_state->rpc_struct_map[_backend_id1].request.add_request_block_descs();
     
_shared_state->rpc_struct_map[_backend_id2].request.add_request_block_descs();
 
-    // --- Build BE_1's response: blocks[0]=1 row (INT), blocks[1]=empty ---
+    // --- Build BE_1's response: blocks[0]=1 row (INT), no block for relation 
1 ---
     {
         add_request_row(
                 
_shared_state->rpc_struct_map[_backend_id1].request.mutable_request_block_descs(0),
@@ -473,8 +508,8 @@ TEST_F(MaterializationSharedStateTest, 
TestMergeMultiResponseStaleBlockMaps) {
         ASSERT_TRUE(rel0_block.serialize(0, pb0, &us, &cs, &ct, 
CompressionTypePB::LZ4).ok());
         _shared_state->response_blocks[0] = rel0_block.clone_empty();
 
-        // blocks[1]: empty (BE_1 has no data for relation 1)
-        response.add_blocks();
+        // BE_1 has no data for relation 1. This also models an ignored 
transport failure for
+        // an empty RPC: there is no response block to deserialize for that 
relation.
 
         _shared_state->rpc_struct_map[_backend_id1].response = 
std::move(response);
     }


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

Reply via email to