Skip to content

HDDS-16628. RackAware placement should try other replicas' racks before giving up - #11403

Open
peterxcli wants to merge 5 commits into
apache:masterfrom
peterxcli:HDDS-16628
Open

peterxcli wants to merge 5 commits into
apache:masterfrom
peterxcli:HDDS-16628

Conversation

@peterxcli

@peterxcli peterxcli commented Oct 4, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

When a container needs a new replica, SCMContainerPlacementRackAware often tries to put it on the same rack as one of the existing replicas. It picks a few random nodes from that rack, and if none of them can be used after 3 tries, it gives up on the whole placement.

The catch is that decommissioned and in-maintenance datanodes are still part of the network topology. So if the first rack it tries has only such nodes left, all 3 tries are wasted there, and SCM never even looks at the racks of the other replicas or at the fallback. The container stays under-replicated, and a decommission waiting for it never finishes. We found this with the SCM simulation (HDDS-16627).

With this change:

  • Each rack gets its own 3 tries. If one rack doesn't work out, placement moves on to the next one, and then to the fallback.
  • chooseNode is rewritten as a short list of places to try, in order, instead of one loop that juggled several counters. This also makes the rule above easy to see.
  • The final error now says "No writable datanode with enough space was found", because the problem is not always disk space.

However, with the fix above, placement reaches the fallback much more often, and that exposed an older problem in it (found by @chungen0126 in review). The fallback stayed off the racks of the excluded nodes instead of the racks of the existing replicas. So an excluded node, usually a replica being decommissioned, ruled out its whole rack, while the racks that had just failed were tried again. The fallback now stays off the racks of the existing replicas.

What is the link to the Apache JIRA

https://issues.apache.org/jira/browse/HDDS-16628

How was this patch tested?

Three new tests in TestSCMContainerPlacementRackAware:

  • the first replica's rack has only decommissioned nodes left, so the new replica should go to the other replica's rack;
  • both replicas' racks have only decommissioned nodes left, so placement should fall back to a third rack;
  • the only rack without a replica also has an excluded node on it, and placement should still use that rack in the fallback.

Copilot AI balanced review requested due to automatic review settings October 4, 2026 14:47

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@peterxcli
peterxcli marked this pull request as ready for review October 5, 2026 03:45

assertEquals(1, datanodeDetails.size());
// After both racks fail, placement falls back to rack2, the only rack with usable nodes.
assertTrue(cluster.isSameParent(datanodes.get(10), datanodeDetails.get(0)));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we also assert metrics.getDatanodeChooseFallbackCount() is 1 here and 0 in the preceding test, since trying another affinity rack should not count as a fallback?

@chungen0126 chungen0126 left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @peterxcli for working on this.

I have a question regarding the fallback behavior when both affinityNodes and excludedNodes are not null, with fallback = true, see following:

In the method entry, we have the logic:

// When affinity node is null, in this case new node to be selected
// should be in different rack than used nodes rack.
// Exclude nodes should be just excluded from topology node selection,
// which is filled in excludedNodesForCapacity
// Used node rack should not be part of rack selection
// which is filled in excludedNodes
if (affinityNodes == null && excludedNodes != null) {
  for (DatanodeDetails node : excludedNodes) {
    excludedNodesForCapacity.add(node.getNetworkFullPath());
  }
  excludedNodes = usedNodes;
}

Since affinityNodes != null initially, this block is skipped, leaving excludedNodes as the original excluded list .

When all tries in Step 1 fail, we proceed to Step 2:

steps.add(new Step(null, RACK_LEVEL, affinityNodes != null));

Here step.affinityNode is null and step.ancestorGen is RACK_LEVEL. In this step:

  1. networkTopology.chooseRandom will exclude entire racks of the nodes in excludedNodes.
  2. However, because excludedNodes was never swapped with usedNodes, the topology will exclude the racks containing the original excludedNodes, rather than excluding the racks of usedNodes.

Should we perform the same swapping/handling between excludedNodes and usedNodes when transitioning to the fallback step?

@peterxcli
peterxcli requested a review from chungen0126 October 6, 2026 05:50
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants