Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions hadoop-hdds/docs/content/feature/Reconfigurability.md
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,8 @@ ozone admin reconfig --service=[OM|SCM|DATANODE] --address=<ip:port|hostname:por
| `ozone.directory.deleting.service.interval` | `60s` | Directory deletion service run interval |
| `ozone.thread.number.dir.deletion` | `10` | Number of threads for directory deletion |
| `ozone.snapshot.filtering.service.interval` | `60s` | Snapshot SST filtering service run interval |
| `ozone.scm.nodes.<scmServiceId>` | - | Comma-separated HA SCM node IDs.|
| `ozone.scm.address.<scmServiceId>.<scmNodeId>` | - | SCM RPC address in HA mode.|

### Storage Container Manager (SCM)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@ public abstract class SCMFailoverProxyProviderBase<T> implements FailoverProxyPr
// scmNodeId -> ProxyInfo<rpcProxy>
private final Map<String, ProxyInfo<T>> scmProxies;
// scmNodeId -> SCMProxyInfo
private final Map<String, SCMProxyInfo> scmProxyInfoMap;
private Map<String, SCMProxyInfo> scmProxyInfoMap;
private List<String> scmNodeIds;

// As SCM Client is shared across threads, performFailOver()
Expand Down Expand Up @@ -122,7 +122,6 @@ public SCMFailoverProxyProviderBase(Class<T> protocol, ConfigurationSource conf,
this.scmVersion = RPC.getProtocolVersion(protocol);

this.scmProxies = new HashMap<>();
this.scmProxyInfoMap = new HashMap<>();
loadConfigs();

this.currentProxyIndex = 0;
Expand Down Expand Up @@ -174,9 +173,19 @@ synchronized void replaceProxyInfoForTest(String nodeId, SCMProxyInfo info) {

@VisibleForTesting
protected synchronized void loadConfigs() {
List<SCMNodeInfo> scmNodeInfoList = SCMNodeInfo.buildNodeInfo(conf);
scmNodeIds = new ArrayList<>();
ScmProxyConfig newConfig = buildConfigs();
scmNodeIds = newConfig.nodeIds;
scmProxyInfoMap = newConfig.proxyInfoMap;
}

/**
* Thread-safely resolves node addresses into an isolated holder.
* Fails atomically if any address is missing.
*/
private ScmProxyConfig buildConfigs() {
List<SCMNodeInfo> scmNodeInfoList = SCMNodeInfo.buildNodeInfo(conf);
List<String> newScmNodeIds = new ArrayList<>();
Map<String, SCMProxyInfo> newScmProxyInfoMap = new HashMap<>();

for (SCMNodeInfo scmNodeInfo : scmNodeInfoList) {
String protocolAddress = getProtocolAddress(scmNodeInfo);
Expand All @@ -188,15 +197,85 @@ protected synchronized void loadConfigs() {

String scmServiceId = scmNodeInfo.getServiceId();
String scmNodeId = scmNodeInfo.getNodeId();
scmNodeIds.add(scmNodeId);
newScmNodeIds.add(scmNodeId);
// Preserve the original config string so DNS can be re-resolved
// on connection failure when the SCM peer is rescheduled to a
// new IP (Kubernetes pod-IP-change recovery). See
// refreshProxyAddressIfChanged(String).
SCMProxyInfo scmProxyInfo = new SCMProxyInfo(scmServiceId, scmNodeId,
protocolAddr, protocolAddress);
scmProxyInfoMap.put(scmNodeId, scmProxyInfo);
newScmProxyInfoMap.put(scmNodeId, scmProxyInfo);
}
}

return new ScmProxyConfig(newScmNodeIds, newScmProxyInfoMap);
}
Comment on lines +208 to +212

@szetszwo szetszwo Sep 24, 2026 •

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.

The else is unnecessarily removed here and the comment is reformatted. Please revert them and minimize the change.

It took me 10 min to figure out such a simple change here.

+  private ScmProxyConfig buildConfigs() {
     List<SCMNodeInfo> scmNodeInfoList = SCMNodeInfo.buildNodeInfo(conf);
-    scmNodeIds = new ArrayList<>();
-
+    final List<String> newScmNodeIds = new ArrayList<>();
+    final Map<String, SCMProxyInfo> newScmProxyInfoMap = new HashMap<>();
 
     for (SCMNodeInfo scmNodeInfo : scmNodeInfoList) {
       String protocolAddress = getProtocolAddress(scmNodeInfo);
@@ -188,16 +199,96 @@ protected synchronized void loadConfigs() {
 
         String scmServiceId = scmNodeInfo.getServiceId();
         String scmNodeId = scmNodeInfo.getNodeId();
-        scmNodeIds.add(scmNodeId);
+        newScmNodeIds.add(scmNodeId);
         // Preserve the original config string so DNS can be re-resolved
         // on connection failure when the SCM peer is rescheduled to a
         // new IP (Kubernetes pod-IP-change recovery). See
         // refreshProxyAddressIfChanged(String).
         SCMProxyInfo scmProxyInfo = new SCMProxyInfo(scmServiceId, scmNodeId,
             protocolAddr, protocolAddress);
-        scmProxyInfoMap.put(scmNodeId, scmProxyInfo);
+        newScmProxyInfoMap.put(scmNodeId, scmProxyInfo);
+      }
+    }
+
+    return new ScmProxyConfig(newScmNodeIds, newScmProxyInfoMap);
+  }

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Reverted to else as requested, though it's technically unnecessary since the throw exits the method. Kept the newScmProxyInfoMap/newScmNodeIds locals because buildConfigs() must not touch shared state — it returns a fresh holder so changeConfig() can resolve addresses off-lock, and a mid-loop throw leaves the live config untouched.


/**
* Reloads SCM nodes and proxies from updated config without a restart.
* Stops proxies for removed/changed nodes; fails atomically if config is invalid.
* In-flight calls on removed nodes persist until the next failover.
*/
public void changeConfig() {
// Resolve DNS before locking to prevent slow lookup from blocking callers.
ScmProxyConfig newConfig = buildConfigs();

Map<String, ProxyInfo<T>> staleProxies = new HashMap<>();
synchronized (this) {
Map<String, SCMProxyInfo> oldProxyInfoMap = scmProxyInfoMap;
scmNodeIds = newConfig.nodeIds;
scmProxyInfoMap = newConfig.proxyInfoMap;

// Re-sync proxy index to the new list, or fall back to first node if removed.
int newProxyIndex = scmNodeIds.indexOf(currentProxySCMNodeId);
if (newProxyIndex < 0) {
newProxyIndex = 0;
currentProxySCMNodeId = scmNodeIds.get(newProxyIndex);
}
currentProxyIndex = newProxyIndex;

// Drop removed failover target to prevent NPE on next failover.
if (updatedLeaderNodeID != null
&& !scmProxyInfoMap.containsKey(updatedLeaderNodeID)) {
updatedLeaderNodeID = null;
}

// Evict stale proxies under lock, but defer stopProxy until unlocked to avoid blocking.
for (Map.Entry<String, SCMProxyInfo> entry : oldProxyInfoMap.entrySet()) {
String nodeId = entry.getKey();
SCMProxyInfo newInfo = scmProxyInfoMap.get(nodeId);
if (newInfo == null
|| !newInfo.getAddress().equals(entry.getValue().getAddress())) {
ProxyInfo<T> staleProxy = scmProxies.remove(nodeId);
if (staleProxy != null && staleProxy.proxy != null) {
staleProxies.put(nodeId, staleProxy);
}
}
}

getLogger().info("Reloaded SCM proxy configuration for protocol {} with {} nodes: {}",
protocolClass.getSimpleName(), scmNodeIds.size(), scmProxyInfoMap.values());
}

for (Map.Entry<String, ProxyInfo<T>> entry : staleProxies.entrySet()) {
try {
RPC.stopProxy(entry.getValue().proxy);
} catch (RuntimeException stopEx) {
getLogger().warn("Failed to stop stale proxy for SCM node {}",
entry.getKey(), stopEx);
}
Comment on lines +260 to 266

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.

Why not just stop the stale proxies in the loop above?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

RPC.stopProxy() can block while the IPC client tears down its connections, and doing that inside the synchronized block would stall every concurrent getProxy() / shouldRetry() caller on the provider monitor. Same reason we resolve DNS before locking. refreshProxyAddressIfChanged() already uses the same evict-under-lock / stop-after-unlock shape.

}
}

/** Parsed node list and resolved addresses, built without touching shared state. */
private static final class ScmProxyConfig {
private final List<String> nodeIds;
private final Map<String, SCMProxyInfo> proxyInfoMap;

ScmProxyConfig(List<String> nodeIds,
Map<String, SCMProxyInfo> proxyInfoMap) {
this.nodeIds = nodeIds;
this.proxyInfoMap = proxyInfoMap;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -131,18 +131,37 @@ public static boolean addSCM(OzoneConfiguration conf, AddSCMRequest request,
*/
public static ScmBlockLocationProtocol getScmBlockClient(
OzoneConfiguration conf) {
return getScmBlockClient(conf,
new SCMBlockLocationFailoverProxyProvider(conf));
}

/**
* Creates an SCM block client using the provided proxy provider.
* Retains the provider reference to support dynamic SCM node updates.
*/
public static ScmBlockLocationProtocol getScmBlockClient(
OzoneConfiguration conf,
SCMBlockLocationFailoverProxyProvider proxyProvider) {
ScmBlockLocationProtocolClientSideTranslatorPB scmBlockLocationClient =
new ScmBlockLocationProtocolClientSideTranslatorPB(
new SCMBlockLocationFailoverProxyProvider(conf), conf);
new ScmBlockLocationProtocolClientSideTranslatorPB(proxyProvider, conf);
return TracingUtil
.createProxy(scmBlockLocationClient, ScmBlockLocationProtocol.class,
conf);
}

public static StorageContainerLocationProtocol getScmContainerClient(
ConfigurationSource conf) {
SCMContainerLocationFailoverProxyProvider proxyProvider =
new SCMContainerLocationFailoverProxyProvider(conf, null);
return getScmContainerClient(conf,
new SCMContainerLocationFailoverProxyProvider(conf, null));
}

/**
* Creates an SCM container client using the provided proxy provider.
* Retains the provider reference to support dynamic SCM node updates.
*/
public static StorageContainerLocationProtocol getScmContainerClient(
ConfigurationSource conf,
SCMContainerLocationFailoverProxyProvider proxyProvider) {
StorageContainerLocationProtocol scmContainerClient =
TracingUtil.createProxy(
new StorageContainerLocationProtocolClientSideTranslatorPB(
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
/*
* 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.
*/

package org.apache.hadoop.hdds.scm.proxy;

import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_ADDRESS_KEY;
import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_NODES_KEY;
import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_SERVICE_IDS_KEY;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotSame;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;

import java.util.List;
import org.apache.hadoop.hdds.conf.ConfigurationException;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.ozone.ha.ConfUtils;
import org.junit.jupiter.api.Test;

/**
* Verifies that {@link SCMFailoverProxyProviderBase#changeConfig()} reloads the
* SCM node list from an updated configuration: adding and removing nodes, keeping
* the current proxy pointer valid, and leaving prior state intact if the new
* configuration is incomplete.
*/
public class TestSCMFailoverProxyProviderChangeConfig {

private static final String SERVICE_ID = "scmservice";

private static OzoneConfiguration haConf(String nodes) {
OzoneConfiguration conf = new OzoneConfiguration();
conf.set(OZONE_SCM_SERVICE_IDS_KEY, SERVICE_ID);
conf.set(ConfUtils.addSuffix(OZONE_SCM_NODES_KEY, SERVICE_ID), nodes);
return conf;
}

private static void setAddress(OzoneConfiguration conf, String nodeId,
String host) {
conf.set(ConfUtils.addKeySuffixes(OZONE_SCM_ADDRESS_KEY, SERVICE_ID, nodeId),
host);
}

@Test
public void testChangeConfigAddsNode() {
OzoneConfiguration conf = haConf("scm1,scm2");
setAddress(conf, "scm1", "host1");
setAddress(conf, "scm2", "host2");

SCMBlockLocationFailoverProxyProvider provider =
new SCMBlockLocationFailoverProxyProvider(conf);
assertEquals(2, provider.getSCMNodeIds().size());

// Operator adds a third SCM: its address key first, then the node list.
setAddress(conf, "scm3", "host3");
conf.set(ConfUtils.addSuffix(OZONE_SCM_NODES_KEY, SERVICE_ID),
"scm1,scm2,scm3");
provider.changeConfig();

List<String> nodeIds = provider.getSCMNodeIds();
assertEquals(3, nodeIds.size());
assertTrue(nodeIds.contains("scm3"));
}

@Test
public void testChangeConfigRemovesNode() {
OzoneConfiguration conf = haConf("scm1,scm2,scm3");
setAddress(conf, "scm1", "host1");
setAddress(conf, "scm2", "host2");
setAddress(conf, "scm3", "host3");

SCMBlockLocationFailoverProxyProvider provider =
new SCMBlockLocationFailoverProxyProvider(conf);

// Pass the preceding node so changeCurrentProxy advances to scm3 (to be removed).
List<String> before = provider.getSCMNodeIds();
int size = before.size();
int scm3Index = before.indexOf("scm3");
provider.changeCurrentProxy(before.get((scm3Index - 1 + size) % size));
assertEquals("scm3", provider.getCurrentProxySCMNodeId());

conf.set(ConfUtils.addSuffix(OZONE_SCM_NODES_KEY, SERVICE_ID), "scm1,scm2");
provider.changeConfig();

List<String> nodeIds = provider.getSCMNodeIds();
assertEquals(2, nodeIds.size());
assertTrue(nodeIds.contains("scm1"));
assertTrue(nodeIds.contains("scm2"));
// The current proxy pointer must fall back to a still-configured node.
assertTrue(nodeIds.contains(provider.getCurrentProxySCMNodeId()));
}

@Test
public void testAddressChangeEvictsCachedProxy() {
OzoneConfiguration conf = haConf("scm1,scm2");
setAddress(conf, "scm1", "127.0.0.1");
setAddress(conf, "scm2", "127.0.0.2");

SCMBlockLocationFailoverProxyProvider provider =
new SCMBlockLocationFailoverProxyProvider(conf);

String current = provider.getCurrentProxySCMNodeId();
Object first = provider.getProxy();
// With no configuration change the cached proxy is reused.
assertSame(first, provider.getProxy());

// Change the current node's address: changeConfig must evict the stale proxy
// so the next getProxy() rebuilds against the new endpoint.
setAddress(conf, current, "127.0.0.3");
provider.changeConfig();

assertNotSame(first, provider.getProxy());
}

@Test
public void testChangeConfigFailsWhenAddressMissing() {
OzoneConfiguration conf = haConf("scm1,scm2");
setAddress(conf, "scm1", "host1");
setAddress(conf, "scm2", "host2");

SCMBlockLocationFailoverProxyProvider provider =
new SCMBlockLocationFailoverProxyProvider(conf);

// Node list references scm3 but its address is not set yet, so the reload
// must fail and leave the previous node set intact for a retry.
conf.set(ConfUtils.addSuffix(OZONE_SCM_NODES_KEY, SERVICE_ID),
"scm1,scm2,scm3");
assertThrows(ConfigurationException.class, provider::changeConfig);

List<String> nodeIds = provider.getSCMNodeIds();
assertEquals(2, nodeIds.size());
assertTrue(nodeIds.contains("scm1"));
assertTrue(nodeIds.contains("scm2"));
}
}
Loading
Loading