github-actions[bot] commented on code in PR #68275:
URL: https://github.com/apache/doris/pull/68275#discussion_r4120351063


##########
be/src/exec/sink/writer/vtablet_writer_v2.cpp:
##########
@@ -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)) {

Review Comment:
   [P1] Keep a global quorum check when there are no success records. This 
bound is per source: with three replicas and three source BEs, each destination 
can report a close-phase tablet failure to a different source. Each source sees 
one failure (within max_failed=1), returns OK, and emits no TTabletCommitInfo. 
FE derives tableToPartition only from commit infos, so an empty list skips 
tablet quorum validation; the transaction can become VISIBLE although no 
replica has the inserted rows. The same gap occurs with min_load_replica_num=1 
when two opened destinations fail and the third did not open. Please validate 
attempted tablets at FE or otherwise ensure a zero-success load fails, and add 
a distributed all-failure test.



##########
be/src/exec/sink/writer/vtablet_writer_v2.cpp:
##########
@@ -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());

Review Comment:
   [P2] Limit version-gap failures to replicas in this tablet's load location. 
FE sends `totalReplicaNum` from the configured allocation but 
`tabletVersionGapBackends` from every actual replica. During redundant-replica 
cleanup a three-replica tablet can have an extra CLONE replica B4 with a 
version gap; FE excludes B4 from load destinations. If B1/B2 succeed and B3 
fails, FE has the required two successes, but this union also counts B4 and the 
new `2 > max_failed(1)` check aborts the valid load. Intersect gap IDs with 
planned replica IDs (or use an actual-replica budget), and cover this 
extra-replica case in a test.



##########
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,

Review Comment:
   [P2] Verify that the selected failed backend actually receives the injected 
final tablet failure. `require_failure` only runs inside its `CLOSE_LOAD` 
handler. If `open_load_stream` to that BE fails from every source, the writer 
can still write to the other two replicas, no guard runs on this BE, and the 
INSERT plus row comparisons pass even with the old success-count code. A 
test-visible hit/result assertion would make this regression prove the 
split-report case.



##########
be/src/exec/sink/writer/vtablet_writer_v2.cpp:
##########
@@ -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

Review Comment:
   [P2] Make this fallback match the existing required-replica fallback for 
even replica counts. Local sinks share `LoadStreamMap`, but each writer owns 
`_tablet_replica_info`; only the writer that requested an automatic partition 
fills it, while a different last-closing sink can build the shared commit info 
with `{0,0}` metadata. With two replicas and `min_load_replica_num=1`, one 
success and one close failure should pass FE (and the old `(2+1)/2=1` local 
requirement), but `(2-1)/2=0` rejects it here. Derive the allowed failures from 
the old fallback requirement or share the tablet metadata, and test this 
multi-sink auto-partition case.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to