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 e4fb2c4dc50 [fix](cloud) Exclude decommissioning BE from BE selection
(#65221) (#66076)
e4fb2c4dc50 is described below
commit e4fb2c4dc503f88eeea97838a0f437dafda2839c
Author: deardeng <[email protected]>
AuthorDate: Thu Jul 30 09:08:41 2026 +0800
[fix](cloud) Exclude decommissioning BE from BE selection (#65221) (#66076)
pick from https://github.com/apache/doris/pull/65221
Problem Summary: Cloud decommissioning backends can still be alive
before they reach the final decommissioned state. These backends could
still be selected by CloudReplica routing, routine load assignment, and
Kafka proxy backend selection. This change treats decommissioning
backends the same as decommissioned backends in those selection paths,
while preserving current master behavior such as direct multi-replica
hashing and routine load blacklist stale-backend cleanup.
Cloud backend selection avoids routing new replica, routine load, and
Kafka proxy work to decommissioning backends.
(cherry picked from commit a6ab604907e752d71f1359d5de9b84a7ca170a77)
### What problem does this PR solve?
Issue Number: close #xxx
Related PR: #xxx
Problem Summary:
### Release note
None
### Check List (For Author)
- Test <!-- At least one of them must be included. -->
- [ ] Regression test
- [ ] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason <!-- Add your reason? -->
- Behavior changed:
- [ ] No.
- [ ] Yes. <!-- Explain the behavior change -->
- Does this need documentation?
- [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
https://github.com/apache/doris-website/pull/1214 -->
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->
---
.../apache/doris/cloud/catalog/CloudReplica.java | 37 ++++++++++++++++++----
1 file changed, 31 insertions(+), 6 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudReplica.java
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudReplica.java
index 0924406c50b..8f8da824f97 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudReplica.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudReplica.java
@@ -107,6 +107,21 @@ public class CloudReplica extends Replica implements
GsonPostProcessable {
return Env.getCurrentColocateIndex().isColocateTableNoLock(tableId);
}
+ private boolean isDecommissioningOrDecommissioned(Backend be) {
+ boolean decommissioning = be.isDecommissioning();
+ boolean decommissioned = be.isDecommissioned();
+ if ((decommissioning || decommissioned) && LOG.isDebugEnabled()) {
+ LOG.debug("backend {} is filtered by decommission state,
decommissioning={}, decommissioned={}, "
+ + "replica info {}",
+ be.getId(), decommissioning, decommissioned, this);
+ }
+ return decommissioning || decommissioned;
+ }
+
+ private boolean isQueryAvailableAndNotDecommissioning(Backend be) {
+ return be != null && be.isQueryAvailable() &&
!isDecommissioningOrDecommissioned(be);
+ }
+
public long getColocatedBeId(String clusterId) throws
ComputeGroupException {
List<Backend> clusterBackends = ((CloudSystemInfoService)
Env.getCurrentSystemInfo())
.getBackendsByClusterId(clusterId);
@@ -131,7 +146,7 @@ public class CloudReplica extends Replica implements
GsonPostProcessable {
List<Backend> decommissionAvailBes = new ArrayList<>();
for (Backend be : bes) {
if (be.isAlive()) {
- if (be.isDecommissioned()) {
+ if (isDecommissioningOrDecommissioned(be)) {
decommissionAvailBes.add(be);
} else {
availableBes.add(be);
@@ -160,7 +175,7 @@ public class CloudReplica extends Replica implements
GsonPostProcessable {
.collect(Collectors.toList());
long index = getIndexByBeNum(hashCode.asLong() + idx,
beAliveOrDeadShort.size());
Backend be = beAliveOrDeadShort.get((int) index);
- if (be.isAlive() && !be.isDecommissioned()) {
+ if (be.isAlive() && !isDecommissioningOrDecommissioned(be)) {
return be.getId();
}
}
@@ -294,6 +309,10 @@ public class CloudReplica extends Replica implements
GsonPostProcessable {
int indexRand = rand.nextInt(Config.cloud_replica_num);
List<Long> res = hashReplicaToBes(clusterId, false,
Config.cloud_replica_num);
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("rehash multi replica backend, clusterId {}, replica
info {}, indexRand {}, hashedBes {}",
+ clusterId, this, indexRand, res);
+ }
if (res.size() < indexRand + 1) {
if (res.isEmpty()) {
return -1;
@@ -307,14 +326,14 @@ public class CloudReplica extends Replica implements
GsonPostProcessable {
// use primaryClusterToBackend, if find be normal
Backend be = getPrimaryBackend(clusterId, false);
- if (be != null && be.isQueryAvailable()) {
+ if (isQueryAvailableAndNotDecommissioning(be)) {
return be.getId();
}
if (!Config.enable_immediate_be_assign) {
// use secondaryClusterToBackends, if find be normal
be = getSecondaryBackend(clusterId);
- if (be != null && be.isQueryAvailable()) {
+ if (isQueryAvailableAndNotDecommissioning(be)) {
return be.getId();
}
}
@@ -326,6 +345,12 @@ public class CloudReplica extends Replica implements
GsonPostProcessable {
// be abnormal, rehash it. configure settings to different maps
long pickBeId = hashReplicaToBe(clusterId, false);
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("rehash replica backend, clusterId {}, pickedBeId {},
immediateAssign {}, replica info {}, "
+ + "primaryBackend {}, secondaryBackend {}",
+ clusterId, pickBeId, Config.enable_immediate_be_assign,
this, getPrimaryBackend(clusterId, false),
+ getSecondaryBackend(clusterId));
+ }
if (Config.enable_immediate_be_assign) {
updateClusterToPrimaryBe(clusterId, pickBeId);
} else {
@@ -384,7 +409,7 @@ public class CloudReplica extends Replica implements
GsonPostProcessable {
List<Backend> decommissionAvailBes = new ArrayList<>();
for (Backend be : clusterBes) {
if (be.isQueryAvailable() && !be.isSmoothUpgradeSrc()) {
- if (be.isDecommissioned()) {
+ if (isDecommissioningOrDecommissioned(be)) {
decommissionAvailBes.add(be);
} else {
availableBes.add(be);
@@ -450,7 +475,7 @@ public class CloudReplica extends Replica implements
GsonPostProcessable {
// be core or restart must in heartbeat_interval_second
if ((be.isAlive() || missTimeMs <=
Config.heartbeat_interval_second * 1000L)
&& !be.isSmoothUpgradeSrc()) {
- if (be.isDecommissioned()) {
+ if (isDecommissioningOrDecommissioned(be)) {
decommissionAvailBes.add(be);
} else {
availableBes.add(be);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]