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]