From d2ff7206410c27649a9159c967f869bdbe0431da Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Mon, 14 Sep 2026 15:04:35 +0800 Subject: [PATCH 01/13] HDDS-16382. Dedicated SCM client RPC timeout and retry for OM request critical path --- .../apache/hadoop/ozone/OzoneConfigKeys.java | 3 +- .../src/main/resources/ozone-default.xml | 54 ++++++ ...SCMBlockLocationFailoverProxyProvider.java | 9 +- ...ontainerLocationFailoverProxyProvider.java | 14 +- .../proxy/SCMFailoverProxyProviderBase.java | 34 +++- .../org/apache/hadoop/hdds/utils/HAUtils.java | 25 ++- .../ozone/om/IdentitySocketFactory.java | 72 +++++++ .../apache/hadoop/ozone/om/OMConfigKeys.java | 28 +++ .../ozone/om/OMScmLocationClientConfig.java | 123 ++++++++++++ .../ozone/om/TestIdentitySocketFactory.java | 52 +++++ .../om/TestOMScmLocationClientConfig.java | 181 ++++++++++++++++++ .../apache/hadoop/ozone/om/OzoneManager.java | 28 ++- 12 files changed, 595 insertions(+), 28 deletions(-) create mode 100644 hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/IdentitySocketFactory.java create mode 100644 hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMScmLocationClientConfig.java create mode 100644 hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestIdentitySocketFactory.java create mode 100644 hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneConfigKeys.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneConfigKeys.java index ab735c23940c..3a591b67201c 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneConfigKeys.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneConfigKeys.java @@ -710,7 +710,8 @@ public final class OzoneConfigKeys { "hdds.scmclient.max.retry.timeout"; public static final String HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY = "hdds.scmclient.failover.max.retry"; - + public static final String HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL = + "hdds.scmclient.failover.retry.interval"; public static final String OZONE_XCEIVER_CLIENT_METRICS_PERCENTILES_INTERVALS_SECONDS_KEY = "ozone.xceiver.client.metrics.percentiles.intervals.seconds"; diff --git a/hadoop-hdds/common/src/main/resources/ozone-default.xml b/hadoop-hdds/common/src/main/resources/ozone-default.xml index 2de7451e1d22..6a0f5ca7e6e6 100644 --- a/hadoop-hdds/common/src/main/resources/ozone-default.xml +++ b/hadoop-hdds/common/src/main/resources/ozone-default.xml @@ -602,6 +602,60 @@ This config overrides Hadoop configuration "ipc.server.read.threadpool.size" for Ozone Manager. + + ozone.om.scmclient.location.rpc.timeout + 30s + + RPC response timeout for OM's SCM block and container location + clients. + + + + ozone.om.scmclient.location.ipc.connect.timeout + 5s + + TCP connect timeout for OM's SCM block and container location + clients. The value accepts a duration suffix such as ms, s, or m and is + converted to milliseconds for Hadoop IPC. + + + + ozone.om.scmclient.location.ipc.connect.max.retries.on.timeouts + 0 + + Maximum TCP connect retries after a timeout for OM's SCM block and + container location clients. A value of 0 immediately fails over to + another SCM peer. + + + + ozone.om.scmclient.location.failover.max.retry + 3 + + Maximum failover retries for OM's SCM block and container location + clients. The effective retry count is the larger of this value and the + quotient of ozone.om.scmclient.location.max.retry.timeout divided by + ozone.om.scmclient.location.failover.retry.interval. + + + + ozone.om.scmclient.location.max.retry.timeout + 6s + + Duration used with ozone.om.scmclient.location.failover.retry.interval + to derive the failover retry count for OM's SCM block and container + location clients. This is not a limit on the total duration of an RPC + call. + + + + ozone.om.scmclient.location.failover.retry.interval + 2s + + Delay between failover retries for OM's SCM block and container location + clients. + + ozone.om.http-address 0.0.0.0:9874 diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMBlockLocationFailoverProxyProvider.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMBlockLocationFailoverProxyProvider.java index 256eb30d80a9..22ae1799599d 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMBlockLocationFailoverProxyProvider.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMBlockLocationFailoverProxyProvider.java @@ -17,6 +17,7 @@ package org.apache.hadoop.hdds.scm.proxy; +import javax.net.SocketFactory; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.scm.ha.SCMNodeInfo; import org.apache.hadoop.hdds.scm.protocolPB.ScmBlockLocationProtocolPB; @@ -32,7 +33,12 @@ public class SCMBlockLocationFailoverProxyProvider extends LoggerFactory.getLogger(SCMBlockLocationFailoverProxyProvider.class); public SCMBlockLocationFailoverProxyProvider(ConfigurationSource conf) { - super(ScmBlockLocationProtocolPB.class, conf, null); + this(conf, null); + } + + public SCMBlockLocationFailoverProxyProvider(ConfigurationSource conf, + SocketFactory socketFactory) { + super(ScmBlockLocationProtocolPB.class, conf, null, socketFactory); } @Override @@ -45,4 +51,3 @@ protected String getProtocolAddress(SCMNodeInfo scmNodeInfo) { return scmNodeInfo.getBlockClientAddress(); } } - diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMContainerLocationFailoverProxyProvider.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMContainerLocationFailoverProxyProvider.java index 242123bd0054..02fa5a290f6e 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMContainerLocationFailoverProxyProvider.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMContainerLocationFailoverProxyProvider.java @@ -17,6 +17,8 @@ package org.apache.hadoop.hdds.scm.proxy; +import jakarta.annotation.Nullable; +import javax.net.SocketFactory; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.scm.ha.SCMNodeInfo; import org.apache.hadoop.hdds.scm.protocolPB.StorageContainerLocationProtocolPB; @@ -32,14 +34,14 @@ public class SCMContainerLocationFailoverProxyProvider extends private static final Logger LOG = LoggerFactory.getLogger(SCMContainerLocationFailoverProxyProvider.class); - /** - * Construct SCMContainerLocationFailoverProxyProvider. - * If userGroupInformation is not null, use the passed ugi, else obtain - * from {@link UserGroupInformation#getCurrentUser()} - */ public SCMContainerLocationFailoverProxyProvider(ConfigurationSource conf, UserGroupInformation userGroupInformation) { - super(StorageContainerLocationProtocolPB.class, conf, userGroupInformation); + this(conf, userGroupInformation, null); + } + + public SCMContainerLocationFailoverProxyProvider(ConfigurationSource conf, + UserGroupInformation userGroupInformation, @Nullable SocketFactory socketFactory) { + super(StorageContainerLocationProtocolPB.class, conf, userGroupInformation, socketFactory); } @Override diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java index 05bcbf6bffa2..04d700bb515f 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java @@ -19,6 +19,7 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.net.InetAddresses; +import jakarta.annotation.Nullable; import java.io.IOException; import java.net.InetAddress; import java.net.InetSocketAddress; @@ -32,6 +33,7 @@ import java.util.Optional; import java.util.OptionalInt; import java.util.stream.Collectors; +import javax.net.SocketFactory; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hdds.HddsUtils; import org.apache.hadoop.hdds.conf.ConfigurationException; @@ -48,6 +50,7 @@ import org.apache.hadoop.ipc_.ProtobufRpcEngine; import org.apache.hadoop.ipc_.RPC; import org.apache.hadoop.net.NetUtils; +import org.apache.hadoop.net.StandardSocketFactory; import org.apache.hadoop.ozone.OzoneConfigKeys; import org.apache.hadoop.security.UserGroupInformation; import org.slf4j.Logger; @@ -88,9 +91,16 @@ public abstract class SCMFailoverProxyProviderBase implements FailoverProxyPr private final long retryInterval; private final UserGroupInformation ugi; + @Nullable + private final SocketFactory socketFactory; private String updatedLeaderNodeID = null; + public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, + UserGroupInformation userGroupInformation) { + this(protocol, conf, userGroupInformation, null); + } + /** * When true, on each connection-class failure the provider re-resolves * the cached SCM hostname and rebuilds the proxy if the IP has changed @@ -101,13 +111,31 @@ public abstract class SCMFailoverProxyProviderBase implements FailoverProxyPr /** * Construct SCMFailoverProxyProviderBase. + * Construct SCMContainerLocationFailoverProxyProvider. + *

* If userGroupInformation is not null, use the passed ugi, else obtain * from {@link UserGroupInformation#getCurrentUser()} + *

+ * Additionally, we can optionally specify {@link SocketFactory} which can be used for use + * case when two clients require different timeout configuration. This is because Hadoop + * {@link org.apache.hadoop.ipc.ClientCache} uses {@link SocketFactory} as the cache key. + * The default {@link StandardSocketFactory#hashCode()} hashes to the class name, which + * means that all clients with the same socket factory share the same {@link org.apache.hadoop.ipc.Client}. + * However, sometimes we want to have separate client with different timeout. In that case, + * we need to have a separate {@link SocketFactory}. + * @param protocol SCM RPC protocol + * @param conf Configuration + * @param userGroupInformation custom ugi for this client. If not specified the implementation defaults to + * {@link UserGroupInformation#getCurrentUser()} + * @param socketFactory custom socket factory for the client. If not specified, the implement defaults to + * {@link NetUtils#getDefaultSocketFactory(Configuration)} */ public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, - UserGroupInformation userGroupInformation) { + UserGroupInformation userGroupInformation, + @Nullable SocketFactory socketFactory) { this.protocolClass = protocol; this.conf = conf; + this.socketFactory = socketFactory; if (userGroupInformation == null) { try { @@ -496,10 +524,12 @@ private T createSCMProxy(InetSocketAddress scmAddress) throws IOException { // retries on the same SCM in case of connection exception. This retry // policy essentially results in TRY_ONCE_THEN_FAIL. RetryPolicy connectionRetryPolicy = RetryPolicies.failoverOnNetworkException(0); + SocketFactory sockFactory = socketFactory == null ? + NetUtils.getDefaultSocketFactory(hadoopConf) : socketFactory; return RPC.getProtocolProxy( protocolClass, scmVersion, scmAddress, ugi, - hadoopConf, NetUtils.getDefaultSocketFactory(hadoopConf), + hadoopConf, sockFactory, (int)scmClientConfig.getRpcTimeOut(), connectionRetryPolicy).getProxy(); } diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/HAUtils.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/HAUtils.java index 3ae1b451e1fc..74fb81c28b22 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/HAUtils.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/HAUtils.java @@ -40,6 +40,7 @@ import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import java.util.stream.Stream; +import javax.net.SocketFactory; import org.apache.hadoop.fs.FileUtil; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.HddsUtils; @@ -131,9 +132,15 @@ public static boolean addSCM(OzoneConfiguration conf, AddSCMRequest request, */ public static ScmBlockLocationProtocol getScmBlockClient( OzoneConfiguration conf) { + return getScmBlockClient(conf, null); + } + + public static ScmBlockLocationProtocol getScmBlockClient( + OzoneConfiguration conf, SocketFactory socketFactory) { ScmBlockLocationProtocolClientSideTranslatorPB scmBlockLocationClient = new ScmBlockLocationProtocolClientSideTranslatorPB( - new SCMBlockLocationFailoverProxyProvider(conf), conf); + new SCMBlockLocationFailoverProxyProvider(conf, socketFactory), + conf); return TracingUtil .createProxy(scmBlockLocationClient, ScmBlockLocationProtocol.class, conf); @@ -141,21 +148,21 @@ public static ScmBlockLocationProtocol getScmBlockClient( public static StorageContainerLocationProtocol getScmContainerClient( ConfigurationSource conf) { - SCMContainerLocationFailoverProxyProvider proxyProvider = - new SCMContainerLocationFailoverProxyProvider(conf, null); - StorageContainerLocationProtocol scmContainerClient = - TracingUtil.createProxy( - new StorageContainerLocationProtocolClientSideTranslatorPB( - proxyProvider), StorageContainerLocationProtocol.class, conf); - return scmContainerClient; + return getScmContainerClient(conf, null, null); } @VisibleForTesting public static StorageContainerLocationProtocol getScmContainerClient( ConfigurationSource conf, UserGroupInformation userGroupInformation) { + return getScmContainerClient(conf, userGroupInformation, null); + } + + public static StorageContainerLocationProtocol getScmContainerClient( + ConfigurationSource conf, UserGroupInformation userGroupInformation, + SocketFactory socketFactory) { SCMContainerLocationFailoverProxyProvider proxyProvider = new SCMContainerLocationFailoverProxyProvider(conf, - userGroupInformation); + userGroupInformation, socketFactory); StorageContainerLocationProtocol scmContainerClient = TracingUtil.createProxy( new StorageContainerLocationProtocolClientSideTranslatorPB( diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/IdentitySocketFactory.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/IdentitySocketFactory.java new file mode 100644 index 000000000000..cd86df329d0e --- /dev/null +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/IdentitySocketFactory.java @@ -0,0 +1,72 @@ +/* + * 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.ozone.om; + +import java.io.IOException; +import java.net.InetAddress; +import java.net.Socket; +import org.apache.hadoop.ipc.ClientCache; +import org.apache.hadoop.net.StandardSocketFactory; + +/** + * Delegates socket creation while retaining identity-based equality. + * + *

Hadoop's {@link ClientCache} uses a {@link SocketFactory} as the cache key + * and deduplicates IPC clients when socket factories compare equal. Wrapping a + * factory in this class provides a distinct cache key so a client can use + * configuration that must not be shared with other Hadoop IPC clients. + * + *

Unlike {@link StandardSocketFactory}, this class intentionally does not + * implement {@link Object#equals(Object)} or {@link Object#hashCode()}. + * Retaining {@link Object}'s identity-based implementations ensures that each + * instance remains a distinct Hadoop IPC client cache key. + */ +final class IdentitySocketFactory extends SocketFactory { + private final SocketFactory delegate; + + IdentitySocketFactory(SocketFactory delegate) { + this.delegate = delegate; + } + + @Override + public Socket createSocket() throws IOException { + return delegate.createSocket(); + } + + @Override + public Socket createSocket(String host, int port) throws IOException { + return delegate.createSocket(host, port); + } + + @Override + public Socket createSocket(String host, int port, InetAddress localHost, + int localPort) throws IOException { + return delegate.createSocket(host, port, localHost, localPort); + } + + @Override + public Socket createSocket(InetAddress host, int port) throws IOException { + return delegate.createSocket(host, port); + } + + @Override + public Socket createSocket(InetAddress address, int port, + InetAddress localAddress, int localPort) throws IOException { + return delegate.createSocket(address, port, localAddress, localPort); + } +} diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMConfigKeys.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMConfigKeys.java index 5852e06ea754..ab2262598f2f 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMConfigKeys.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMConfigKeys.java @@ -55,6 +55,34 @@ public final class OMConfigKeys { public static final int OZONE_OM_DB_MAX_OPEN_FILES_DEFAULT = -1; + public static final String OZONE_OM_SCM_LOCATION_CLIENT_RPC_TIMEOUT = + "ozone.om.scmclient.location.rpc.timeout"; + public static final String OZONE_OM_SCM_LOCATION_CLIENT_RPC_TIMEOUT_DEFAULT = + "30s"; + public static final String OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT = + "ozone.om.scmclient.location.ipc.connect.timeout"; + public static final String + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_DEFAULT = + "5s"; + public static final String + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_RETRIES = + "ozone.om.scmclient.location.ipc.connect.max.retries.on.timeouts"; + public static final String + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_RETRIES_DEFAULT = "0"; + public static final String OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY = + "ozone.om.scmclient.location.failover.max.retry"; + public static final String + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY_DEFAULT = "3"; + public static final String OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT = + "ozone.om.scmclient.location.max.retry.timeout"; + public static final String + OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT_DEFAULT = "6s"; + public static final String + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL = + "ozone.om.scmclient.location.failover.retry.interval"; + public static final String + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL_DEFAULT = "2s"; + public static final String OZONE_OM_INTERNAL_SERVICE_ID = "ozone.om.internal.service.id"; diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMScmLocationClientConfig.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMScmLocationClientConfig.java new file mode 100644 index 000000000000..fec79304b7f4 --- /dev/null +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMScmLocationClientConfig.java @@ -0,0 +1,123 @@ +/* + * 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.ozone.om; + +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY_DEFAULT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL_DEFAULT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_DEFAULT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_RETRIES; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_RETRIES_DEFAULT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT_DEFAULT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_RPC_TIMEOUT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_RPC_TIMEOUT_DEFAULT; + +import java.util.concurrent.TimeUnit; +import javax.net.SocketFactory; +import org.apache.hadoop.fs.CommonConfigurationKeysPublic; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.net.NetUtils; +import org.apache.hadoop.ozone.OzoneConfigKeys; + +/** + * Builds the configuration used by OM's SCM location clients. + * + *

The default effective retry limit is the larger of the configured retry + * count and retry timeout divided by retry interval: + *

+ * max(3, 6s / 2s) = 3
+ * 
+ * For failures that always trigger failover, this allows four attempts. The + * configured timeout envelopes are: + *
+ * established connection: 4 * 30s + 3 * 2s = 126s
+ * including TCP connect:   4 * (5s + 30s) + 3 * 2s = 146s
+ * 
+ * + *

Hadoop tracks retries and failovers separately. A sequence of three + * retry-without-failover responses followed by three failover failures can + * therefore make seven attempts. Its configured timeout envelopes are: + *

+ * established connection: 7 * 30s + 6 * 2s = 222s
+ * including TCP connect:   7 * (5s + 30s) + 6 * 2s = 257s
+ * 
+ * Actual calls can complete sooner. These calculations are not caller-side + * deadlines and exclude time outside the configured connect and RPC waits. + */ +public final class OMScmLocationClientConfig { + private OMScmLocationClientConfig() { + } + + public static OzoneConfiguration createScmClientConfiguration( + OzoneConfiguration configuration) { + OzoneConfiguration scmClientConfiguration = + new OzoneConfiguration(configuration); + copy(configuration, scmClientConfiguration, + OZONE_OM_SCM_LOCATION_CLIENT_RPC_TIMEOUT, + OZONE_OM_SCM_LOCATION_CLIENT_RPC_TIMEOUT_DEFAULT, + OzoneConfigKeys.HDDS_SCM_CLIENT_RPC_TIME_OUT); + copyTimeDurationInMillis(configuration, scmClientConfiguration, + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT, + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_DEFAULT, + CommonConfigurationKeysPublic.IPC_CLIENT_CONNECT_TIMEOUT_KEY); + copy(configuration, scmClientConfiguration, + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_RETRIES, + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_RETRIES_DEFAULT, + CommonConfigurationKeysPublic + .IPC_CLIENT_CONNECT_MAX_RETRIES_ON_SOCKET_TIMEOUTS_KEY); + copy(configuration, scmClientConfiguration, + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY, + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY_DEFAULT, + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY); + copy(configuration, scmClientConfiguration, + OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT, + OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT_DEFAULT, + OzoneConfigKeys.HDDS_SCM_CLIENT_MAX_RETRY_TIMEOUT); + copy(configuration, scmClientConfiguration, + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL, + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL_DEFAULT, + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL); + return scmClientConfiguration; + } + + /** + * Creates a socket factory with a separate Hadoop IPC client cache entry so + * the SCM location connection timeout is not shared with other clients. + */ + static SocketFactory createSocketFactory(OzoneConfiguration configuration) { + return new IdentitySocketFactory( + NetUtils.getDefaultSocketFactory(configuration)); + } + + private static void copy(OzoneConfiguration source, + OzoneConfiguration target, String sourceKey, String defaultValue, + String targetKey) { + target.set(targetKey, source.get(sourceKey, defaultValue)); + } + + private static void copyTimeDurationInMillis(OzoneConfiguration source, + OzoneConfiguration target, String sourceKey, String defaultValue, + String targetKey) { + long value = source.getTimeDuration(sourceKey, defaultValue, + TimeUnit.MILLISECONDS); + target.setInt(targetKey, Math.toIntExact(value)); + } +} diff --git a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestIdentitySocketFactory.java b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestIdentitySocketFactory.java new file mode 100644 index 000000000000..e0139d45d8a9 --- /dev/null +++ b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestIdentitySocketFactory.java @@ -0,0 +1,52 @@ +/* + * 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.ozone.om; + +import static org.junit.jupiter.api.Assertions.assertNotSame; + +import javax.net.SocketFactory; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.io.ObjectWritable; +import org.apache.hadoop.ipc.Client; +import org.apache.hadoop.ipc.ClientCache; +import org.apache.hadoop.net.NetUtils; +import org.junit.jupiter.api.Test; + +/** Tests {@link IdentitySocketFactory}. */ +class TestIdentitySocketFactory { + + @Test + void usesIdentityForHadoopClientCacheKeys() { + OzoneConfiguration configuration = new OzoneConfiguration(); + SocketFactory delegate = NetUtils.getDefaultSocketFactory(configuration); + SocketFactory firstFactory = new IdentitySocketFactory(delegate); + SocketFactory secondFactory = new IdentitySocketFactory(delegate); + ClientCache clientCache = new ClientCache(); + Client firstClient = clientCache.getClient(configuration, firstFactory, + ObjectWritable.class); + Client secondClient = clientCache.getClient(configuration, secondFactory, + ObjectWritable.class); + + try { + assertNotSame(firstClient, secondClient); + } finally { + clientCache.stopClient(firstClient); + clientCache.stopClient(secondClient); + } + } +} diff --git a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java new file mode 100644 index 000000000000..1a5c3627e66b --- /dev/null +++ b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java @@ -0,0 +1,181 @@ +/* + * 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.ozone.om; + +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_RETRIES; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT; +import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_RPC_TIMEOUT; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotSame; + +import java.util.concurrent.TimeUnit; +import javax.net.SocketFactory; +import org.apache.hadoop.fs.CommonConfigurationKeysPublic; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.scm.proxy.SCMClientConfig; +import org.apache.hadoop.io.ObjectWritable; +import org.apache.hadoop.ipc.Client; +import org.apache.hadoop.ipc.ClientCache; +import org.apache.hadoop.net.NetUtils; +import org.apache.hadoop.ozone.OzoneConfigKeys; +import org.junit.jupiter.api.Test; + +/** Tests OM-specific SCM client configuration. */ +class TestOMScmLocationClientConfig { + @Test + void overridesScmClientAndIpcTimeoutsWithoutMutatingSourceConfiguration() { + OzoneConfiguration configuration = new OzoneConfiguration(); + configuration.setTimeDuration(OzoneConfigKeys.HDDS_SCM_CLIENT_RPC_TIME_OUT, + 15, TimeUnit.MINUTES); + configuration.setTimeDuration( + CommonConfigurationKeysPublic.IPC_CLIENT_CONNECT_TIMEOUT_KEY, + 20, TimeUnit.SECONDS); + configuration.setInt( + CommonConfigurationKeysPublic + .IPC_CLIENT_CONNECT_MAX_RETRIES_ON_SOCKET_TIMEOUTS_KEY, + 45); + configuration.setInt(OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, + 15); + configuration.setTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_MAX_RETRY_TIMEOUT, + 10, TimeUnit.MINUTES); + configuration.setTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL, + 9, TimeUnit.SECONDS); + configuration.setTimeDuration(OZONE_OM_SCM_LOCATION_CLIENT_RPC_TIMEOUT, + 60, TimeUnit.SECONDS); + configuration.setTimeDuration( + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT, + 7, TimeUnit.SECONDS); + configuration.setInt( + OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT_RETRIES, 1); + configuration.setInt(OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY, 4); + configuration.setTimeDuration( + OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT, + 8, TimeUnit.SECONDS); + configuration.setTimeDuration( + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL, + 1, TimeUnit.SECONDS); + + OzoneConfiguration scmClientConfiguration = + OMScmLocationClientConfig.createScmClientConfiguration(configuration); + + assertEquals(60, scmClientConfiguration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_RPC_TIME_OUT, 0, TimeUnit.SECONDS)); + assertEquals(7_000, scmClientConfiguration.getInt( + CommonConfigurationKeysPublic.IPC_CLIENT_CONNECT_TIMEOUT_KEY, 0)); + assertEquals(1, scmClientConfiguration.getInt( + CommonConfigurationKeysPublic + .IPC_CLIENT_CONNECT_MAX_RETRIES_ON_SOCKET_TIMEOUTS_KEY, + 0)); + assertEquals(4, scmClientConfiguration.getInt( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, 0)); + assertEquals(8, scmClientConfiguration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_MAX_RETRY_TIMEOUT, + 0, TimeUnit.SECONDS)); + assertEquals(1, scmClientConfiguration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL, + 0, TimeUnit.SECONDS)); + assertEquals(15, configuration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_RPC_TIME_OUT, 0, TimeUnit.MINUTES)); + assertEquals(20, configuration.getTimeDuration( + CommonConfigurationKeysPublic.IPC_CLIENT_CONNECT_TIMEOUT_KEY, + 0, TimeUnit.SECONDS)); + assertEquals(45, configuration.getInt( + CommonConfigurationKeysPublic + .IPC_CLIENT_CONNECT_MAX_RETRIES_ON_SOCKET_TIMEOUTS_KEY, + 0)); + assertEquals(15, configuration.getInt( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, 0)); + assertEquals(10, configuration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_MAX_RETRY_TIMEOUT, + 0, TimeUnit.MINUTES)); + assertEquals(9, configuration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL, + 0, TimeUnit.SECONDS)); + } + + @Test + void appliesOmDefaultsWhenOverridesAreNotSet() { + OzoneConfiguration configuration = new OzoneConfiguration(); + configuration.setTimeDuration(OzoneConfigKeys.HDDS_SCM_CLIENT_RPC_TIME_OUT, + 2, TimeUnit.MINUTES); + configuration.setTimeDuration( + CommonConfigurationKeysPublic.IPC_CLIENT_CONNECT_TIMEOUT_KEY, + 10, TimeUnit.SECONDS); + configuration.setInt( + CommonConfigurationKeysPublic + .IPC_CLIENT_CONNECT_MAX_RETRIES_ON_SOCKET_TIMEOUTS_KEY, + 3); + configuration.setInt(OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, + 15); + configuration.setTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_MAX_RETRY_TIMEOUT, + 10, TimeUnit.MINUTES); + configuration.setTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL, + 9, TimeUnit.SECONDS); + + OzoneConfiguration scmClientConfiguration = + OMScmLocationClientConfig.createScmClientConfiguration(configuration); + + assertEquals(30, scmClientConfiguration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_RPC_TIME_OUT, 0, TimeUnit.SECONDS)); + assertEquals(5_000, scmClientConfiguration.getInt( + CommonConfigurationKeysPublic.IPC_CLIENT_CONNECT_TIMEOUT_KEY, 0)); + assertEquals(0, scmClientConfiguration.getInt( + CommonConfigurationKeysPublic + .IPC_CLIENT_CONNECT_MAX_RETRIES_ON_SOCKET_TIMEOUTS_KEY, + 0)); + assertEquals(3, scmClientConfiguration.getInt( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, 0)); + assertEquals(6, scmClientConfiguration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_MAX_RETRY_TIMEOUT, + 0, TimeUnit.SECONDS)); + assertEquals(2, scmClientConfiguration.getTimeDuration( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL, + 0, TimeUnit.SECONDS)); + assertEquals(3, + scmClientConfiguration.getObject(SCMClientConfig.class) + .getRetryCount()); + } + + @Test + void usesAnIpcClientSeparateFromTheDefaultSocketFactory() { + OzoneConfiguration configuration = new OzoneConfiguration(); + SocketFactory defaultSocketFactory = + NetUtils.getDefaultSocketFactory(configuration); + SocketFactory scmLocationSocketFactory = + OMScmLocationClientConfig.createSocketFactory(configuration); + ClientCache clientCache = new ClientCache(); + Client defaultClient = clientCache.getClient(configuration, + defaultSocketFactory, ObjectWritable.class); + Client criticalScmClient = clientCache.getClient(configuration, + scmLocationSocketFactory, ObjectWritable.class); + + try { + assertNotSame(defaultClient, criticalScmClient); + } finally { + clientCache.stopClient(defaultClient); + clientCache.stopClient(criticalScmClient); + } + } +} diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java index 1960e1fc3bdc..9c4bd961bc2d 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java @@ -165,6 +165,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import javax.management.ObjectName; +import javax.net.SocketFactory; import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.hadoop.conf.Configuration; @@ -660,9 +661,20 @@ private OzoneManager(OzoneConfiguration conf, StartupOption startupOption) // Honor property 'hadoop.security.token.service.use_ip' omRpcAddressTxt = new Text(SecurityUtil.buildTokenService(omNodeRpcAddr)); - final StorageContainerLocationProtocol scmContainerClient = getScmContainerClient(configuration); + metrics = OMMetrics.create(); + perfMetrics = OMPerformanceMetrics.register(conf); + + OzoneConfiguration scmLocationClientConfiguration = + OMScmLocationClientConfig.createScmClientConfiguration(configuration); + SocketFactory scmLocationSocketFactory = + OMScmLocationClientConfig.createSocketFactory(configuration); + StorageContainerLocationProtocol scmContainerClient = + getScmContainerClient(scmLocationClientConfiguration, + scmLocationSocketFactory); // verifies that the SCM info in the OM Version file is correct. - final ScmBlockLocationProtocol scmBlockClient = getScmBlockClient(configuration); + ScmBlockLocationProtocol scmBlockClient = + getScmBlockClient(scmLocationClientConfiguration, + scmLocationSocketFactory); scmTopologyClient = new ScmTopologyClient(scmBlockClient); this.scmClient = new ScmClient(scmBlockClient, scmContainerClient, configuration); @@ -1512,23 +1524,23 @@ private static void loginOMUser(OzoneConfiguration conf) } /** - * Create a scm block client, used by putKey() and getKey(). + * Creates a scm block client. * * @return {@link ScmBlockLocationProtocol} */ private static ScmBlockLocationProtocol getScmBlockClient( - OzoneConfiguration conf) { - return HAUtils.getScmBlockClient(conf); + OzoneConfiguration conf, SocketFactory socketFactory) { + return HAUtils.getScmBlockClient(conf, socketFactory); } /** - * Returns a scm container client. + * Returns am scm container client. * * @return {@link StorageContainerLocationProtocol} */ private static StorageContainerLocationProtocol getScmContainerClient( - OzoneConfiguration conf) { - return HAUtils.getScmContainerClient(conf); + OzoneConfiguration conf, SocketFactory socketFactory) { + return HAUtils.getScmContainerClient(conf, null, socketFactory); } /** From 20ae7a9b297c0d71d0b715bf5b8363b6a49ba835 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Mon, 14 Sep 2026 15:15:47 +0800 Subject: [PATCH 02/13] HDDS-16382. Separate key deletion SCM client configuration --- .../org/apache/hadoop/ozone/om/KeyManagerImpl.java | 2 +- .../org/apache/hadoop/ozone/om/OzoneManager.java | 4 +++- .../java/org/apache/hadoop/ozone/om/ScmClient.java | 13 +++++++++++++ .../org/apache/hadoop/ozone/om/TestScmClient.java | 12 ++++++++++++ 4 files changed, 29 insertions(+), 2 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManagerImpl.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManagerImpl.java index d74614579e24..dd12a45056bd 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManagerImpl.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/KeyManagerImpl.java @@ -280,7 +280,7 @@ public void start(OzoneConfiguration configuration) { keyDeletingServiceCorePoolSize = 1; } keyDeletingService = new KeyDeletingService(ozoneManager, - scmClient.getBlockClient(), blockDeleteInterval, + scmClient.getBlockClientForKeyDeletion(), blockDeleteInterval, serviceTimeout, configuration, keyDeletingServiceCorePoolSize, isSnapshotDeepCleaningEnabled); keyDeletingService.start(); } diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java index 9c4bd961bc2d..a7c9653e9fd8 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java @@ -675,9 +675,11 @@ private OzoneManager(OzoneConfiguration conf, StartupOption startupOption) ScmBlockLocationProtocol scmBlockClient = getScmBlockClient(scmLocationClientConfiguration, scmLocationSocketFactory); + ScmBlockLocationProtocol keyDeletionScmBlockClient = + getScmBlockClient(configuration, null); scmTopologyClient = new ScmTopologyClient(scmBlockClient); this.scmClient = new ScmClient(scmBlockClient, scmContainerClient, - configuration); + configuration, keyDeletionScmBlockClient); this.ozoneLockProvider = new OzoneLockProvider(getKeyPathLockEnabled(), getEnableFileSystemPaths()); diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java index a61595faf069..5b86ad965458 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java @@ -54,6 +54,7 @@ public class ScmClient { private final ScmBlockLocationProtocol blockClient; + private final ScmBlockLocationProtocol blockClientForKeyDeletion; private final StorageContainerLocationProtocol containerClient; private final LoadingCache containerLocationCache; private final CacheMetrics containerCacheMetrics; @@ -62,8 +63,16 @@ public class ScmClient { ScmClient(ScmBlockLocationProtocol blockClient, StorageContainerLocationProtocol containerClient, OzoneConfiguration configuration) { + this(blockClient, containerClient, configuration, blockClient); + } + + ScmClient(ScmBlockLocationProtocol blockClient, + StorageContainerLocationProtocol containerClient, + OzoneConfiguration configuration, + ScmBlockLocationProtocol blockClientForKeyDeletion) { this.containerClient = containerClient; this.blockClient = blockClient; + this.blockClientForKeyDeletion = blockClientForKeyDeletion; Cache datanodeDetailsCache = createDatanodeDetailsCache(configuration); this.containerLocationCache = @@ -148,6 +157,10 @@ public ScmBlockLocationProtocol getBlockClient() { return this.blockClient; } + public ScmBlockLocationProtocol getBlockClientForKeyDeletion() { + return this.blockClientForKeyDeletion; + } + public StorageContainerLocationProtocol getContainerClient() { return this.containerClient; } diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestScmClient.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestScmClient.java index 89d8c2162c38..a9611fccd22f 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestScmClient.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestScmClient.java @@ -73,6 +73,18 @@ public void setUp() { containerLocationProtocol, conf); } + @Test + void usesDedicatedBlockClientForKeyDeletion() { + ScmBlockLocationProtocol foregroundClient = mock(ScmBlockLocationProtocol.class); + ScmBlockLocationProtocol keyDeletionClient = mock(ScmBlockLocationProtocol.class); + OzoneConfiguration conf = new OzoneConfiguration(); + ScmClient client = new ScmClient(foregroundClient, + containerLocationProtocol, conf, keyDeletionClient); + + assertSame(foregroundClient, client.getBlockClient()); + assertSame(keyDeletionClient, client.getBlockClientForKeyDeletion()); + } + private static Stream getContainerLocationsTestCases() { return Stream.of( Arguments.of("Existing keys", From 2e1d367e5b23143d0a68eb977bf610d10f490928 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Mon, 14 Sep 2026 15:37:58 +0800 Subject: [PATCH 03/13] HDDS-16382. Fix CI issues for dedicated SCM clients --- .../scm/proxy/SCMFailoverProxyProviderBase.java | 14 +++++++------- .../hadoop/ozone/om/IdentitySocketFactory.java | 1 + 2 files changed, 8 insertions(+), 7 deletions(-) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java index 04d700bb515f..524ca8d1c785 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java @@ -94,13 +94,6 @@ public abstract class SCMFailoverProxyProviderBase implements FailoverProxyPr @Nullable private final SocketFactory socketFactory; - private String updatedLeaderNodeID = null; - - public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, - UserGroupInformation userGroupInformation) { - this(protocol, conf, userGroupInformation, null); - } - /** * When true, on each connection-class failure the provider re-resolves * the cached SCM hostname and rebuilds the proxy if the IP has changed @@ -109,6 +102,13 @@ public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, */ private final boolean resolveOnFailureEnabled; + private String updatedLeaderNodeID = null; + + public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, + UserGroupInformation userGroupInformation) { + this(protocol, conf, userGroupInformation, null); + } + /** * Construct SCMFailoverProxyProviderBase. * Construct SCMContainerLocationFailoverProxyProvider. diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/IdentitySocketFactory.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/IdentitySocketFactory.java index cd86df329d0e..bf32abcba065 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/IdentitySocketFactory.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/IdentitySocketFactory.java @@ -20,6 +20,7 @@ import java.io.IOException; import java.net.InetAddress; import java.net.Socket; +import javax.net.SocketFactory; import org.apache.hadoop.ipc.ClientCache; import org.apache.hadoop.net.StandardSocketFactory; From e085b4f05024f7fab60a17a549251127295af57d Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Mon, 14 Sep 2026 15:55:41 +0800 Subject: [PATCH 04/13] HDDS-16382. Fix common module test dependency --- .../apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java | 4 ---- 1 file changed, 4 deletions(-) diff --git a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java index 1a5c3627e66b..327605f09041 100644 --- a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java +++ b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java @@ -30,7 +30,6 @@ import javax.net.SocketFactory; import org.apache.hadoop.fs.CommonConfigurationKeysPublic; import org.apache.hadoop.hdds.conf.OzoneConfiguration; -import org.apache.hadoop.hdds.scm.proxy.SCMClientConfig; import org.apache.hadoop.io.ObjectWritable; import org.apache.hadoop.ipc.Client; import org.apache.hadoop.ipc.ClientCache; @@ -153,9 +152,6 @@ void appliesOmDefaultsWhenOverridesAreNotSet() { assertEquals(2, scmClientConfiguration.getTimeDuration( OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL, 0, TimeUnit.SECONDS)); - assertEquals(3, - scmClientConfiguration.getObject(SCMClientConfig.class) - .getRetryCount()); } @Test From fac99aa430c1b69aa0a10ff6affa1946dec16aa3 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Mon, 14 Sep 2026 16:24:36 +0800 Subject: [PATCH 05/13] HDDS-16382. Fix Ozone Manager metrics initialization --- .../src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java index a7c9653e9fd8..fd598a2892c5 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java @@ -661,9 +661,6 @@ private OzoneManager(OzoneConfiguration conf, StartupOption startupOption) // Honor property 'hadoop.security.token.service.use_ip' omRpcAddressTxt = new Text(SecurityUtil.buildTokenService(omNodeRpcAddr)); - metrics = OMMetrics.create(); - perfMetrics = OMPerformanceMetrics.register(conf); - OzoneConfiguration scmLocationClientConfiguration = OMScmLocationClientConfig.createScmClientConfiguration(configuration); SocketFactory scmLocationSocketFactory = From b95f3dae859e53f71e332297758fca3a7d7a2c17 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Tue, 15 Sep 2026 10:08:06 +0800 Subject: [PATCH 06/13] HDDS-16382. Document SCM failover retry interval --- hadoop-hdds/common/src/main/resources/ozone-default.xml | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/hadoop-hdds/common/src/main/resources/ozone-default.xml b/hadoop-hdds/common/src/main/resources/ozone-default.xml index 6a0f5ca7e6e6..b999ca31734f 100644 --- a/hadoop-hdds/common/src/main/resources/ozone-default.xml +++ b/hadoop-hdds/common/src/main/resources/ozone-default.xml @@ -656,6 +656,14 @@ clients.
+ + hdds.scmclient.failover.retry.interval + 2s + OZONE, SCM, CLIENT + + SCM Client timeout on waiting for the next connection retry to other SCM IP. + + ozone.om.http-address 0.0.0.0:9874 From 651acbf1644eae8818c4dfe27c401a7f7857c9ea Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Tue, 15 Sep 2026 10:08:42 +0800 Subject: [PATCH 07/13] HDDS-16382. Clean up SCM proxy provider Javadoc --- .../hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java | 1 - 1 file changed, 1 deletion(-) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java index 524ca8d1c785..223dc2b9a31c 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java @@ -111,7 +111,6 @@ public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, /** * Construct SCMFailoverProxyProviderBase. - * Construct SCMContainerLocationFailoverProxyProvider. *

* If userGroupInformation is not null, use the passed ugi, else obtain * from {@link UserGroupInformation#getCurrentUser()} From 7308a61f8a3fa05da94022cc35e112f652c12130 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Tue, 15 Sep 2026 10:13:27 +0800 Subject: [PATCH 08/13] HDDS-16382. Fix SCM container client Javadoc --- .../src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java index fd598a2892c5..1a68c7396d30 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java @@ -1533,7 +1533,7 @@ private static ScmBlockLocationProtocol getScmBlockClient( } /** - * Returns am scm container client. + * Returns a scm container client. * * @return {@link StorageContainerLocationProtocol} */ From 04d5eec30ba0800265f5c1e2a735de41fc751692 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Tue, 15 Sep 2026 10:16:48 +0800 Subject: [PATCH 09/13] HDDS-16382. Preserve SCM proxy field ordering --- .../hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java index 223dc2b9a31c..e7fb779bc960 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java @@ -94,6 +94,8 @@ public abstract class SCMFailoverProxyProviderBase implements FailoverProxyPr @Nullable private final SocketFactory socketFactory; + private String updatedLeaderNodeID = null; + /** * When true, on each connection-class failure the provider re-resolves * the cached SCM hostname and rebuilds the proxy if the IP has changed @@ -102,8 +104,6 @@ public abstract class SCMFailoverProxyProviderBase implements FailoverProxyPr */ private final boolean resolveOnFailureEnabled; - private String updatedLeaderNodeID = null; - public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, UserGroupInformation userGroupInformation) { this(protocol, conf, userGroupInformation, null); From 5155d3409d6c28e88eb77880f0a94b5a928de1a8 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Tue, 15 Sep 2026 10:34:13 +0800 Subject: [PATCH 10/13] HDDS-16382. Make SCM failover retries topology-aware --- .../proxy/SCMFailoverProxyProviderBase.java | 2 +- ...tSCMFailoverProxyProviderRefreshWired.java | 52 +++++++++++++++++++ 2 files changed, 53 insertions(+), 1 deletion(-) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java index e7fb779bc960..32a037edb99e 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java @@ -156,7 +156,7 @@ public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, currentProxySCMNodeId = scmNodeIds.get(currentProxyIndex); scmClientConfig = conf.getObject(SCMClientConfig.class); - this.maxRetryCount = scmClientConfig.getRetryCount(); + this.maxRetryCount = Math.max(scmClientConfig.getRetryCount(), scmNodeIds.size()); this.retryInterval = scmClientConfig.getRetryInterval(); this.resolveOnFailureEnabled = conf.getBoolean( OzoneConfigKeys.OZONE_CLIENT_FAILOVER_RESOLVE_NEEDED_KEY, diff --git a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderRefreshWired.java b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderRefreshWired.java index 83a08a7d8df7..1667a6c82779 100644 --- a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderRefreshWired.java +++ b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderRefreshWired.java @@ -31,9 +31,11 @@ import java.net.InetSocketAddress; import java.net.SocketTimeoutException; import java.util.ArrayList; +import java.util.Arrays; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.ratis.ServerNotLeaderException; import org.apache.hadoop.io.retry.RetryPolicy; +import org.apache.hadoop.io.retry.RetryPolicy.RetryAction.RetryDecision; import org.apache.hadoop.ozone.ha.ConfUtils; import org.apache.ozone.test.GenericTestUtils.LogCapturer; import org.apache.ratis.protocol.RaftPeerId; @@ -56,6 +58,9 @@ public class TestSCMFailoverProxyProviderRefreshWired { private static final String SCM_NODE_1 = "scm1"; private static final String SCM_NODE_2 = "scm2"; private static final String SCM_NODE_3 = "scm3"; + private static final String SCM_NODE_4 = "scm4"; + private static final String SCM_NODE_5 = "scm5"; + private static final String SCM_NODE_6 = "scm6"; private OzoneConfiguration conf; @@ -129,6 +134,34 @@ public void testApplicationLevelErrorDoesNotTriggerRefresh() throws Exception { "ServerNotLeaderException is application-level; refresh must NOT fire"); } + @Test + public void testRetryCountCoversAllConfiguredScmNodes() throws Exception { + RetryPolicy policy = newProvider(6, 3).getRetryPolicy(); + + assertEquals(RetryDecision.FAILOVER_AND_RETRY, + policy.shouldRetry(new IOException("failure"), 0, 5, false).action); + assertEquals(RetryDecision.FAIL, + policy.shouldRetry(new IOException("failure"), 0, 6, false).action); + } + + @Test + public void testRetryCountIsUnchangedWhenItCoversScmNodes() throws Exception { + RetryPolicy policy = newProvider(3, 3).getRetryPolicy(); + + assertEquals(RetryDecision.FAIL, + policy.shouldRetry(new IOException("failure"), 0, 3, false).action); + } + + @Test + public void testConfiguredRetryCountRemainsMinimum() throws Exception { + RetryPolicy policy = newProvider(3, 5).getRetryPolicy(); + + assertEquals(RetryDecision.FAILOVER_AND_RETRY, + policy.shouldRetry(new IOException("failure"), 0, 4, false).action); + assertEquals(RetryDecision.FAIL, + policy.shouldRetry(new IOException("failure"), 0, 5, false).action); + } + @Test public void testFailoverToIpv6SuggestedLeader() { SCMBlockLocationFailoverProxyProvider provider = newThreeNodeProvider(); @@ -268,6 +301,25 @@ private SCMBlockLocationFailoverProxyProvider newThreeNodeProvider() { return new SCMBlockLocationFailoverProxyProvider(conf); } + private SCMBlockLocationFailoverProxyProvider newProvider(int nodeCount, + int retryCount) { + String[] nodeIds = { + SCM_NODE_1, SCM_NODE_2, SCM_NODE_3, SCM_NODE_4, SCM_NODE_5, SCM_NODE_6 + }; + conf.set(OZONE_SCM_NODES_KEY + "." + SCM_SERVICE_ID, + String.join(",", Arrays.copyOf(nodeIds, nodeCount))); + for (int i = 0; i < nodeCount; i++) { + conf.set(ConfUtils.addKeySuffixes(OZONE_SCM_ADDRESS_KEY, + SCM_SERVICE_ID, nodeIds[i]), "localhost"); + } + SCMClientConfig scmClientConfig = new SCMClientConfig(); + scmClientConfig.setRetryCount(retryCount); + scmClientConfig.setMaxRetryTimeout(retryCount * 1000L); + scmClientConfig.setRetryInterval(1000L); + conf.setFromObject(scmClientConfig); + return new SCMBlockLocationFailoverProxyProvider(conf); + } + private static SCMProxyInfo proxyInfoOf( SCMBlockLocationFailoverProxyProvider provider, String nodeId) { return provider.getSCMProxyInfoList().stream() From 07931405cf7982e4d6e579e9b4d0545c08aa9b42 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Tue, 15 Sep 2026 10:50:07 +0800 Subject: [PATCH 11/13] HDDS-16382. Scope topology-aware retries to OM client --- .../proxy/SCMFailoverProxyProviderBase.java | 2 +- ...tSCMFailoverProxyProviderRefreshWired.java | 52 ------------------- .../ozone/om/OMScmLocationClientConfig.java | 35 ++++++++----- .../om/TestOMScmLocationClientConfig.java | 26 ++++++++++ 4 files changed, 50 insertions(+), 65 deletions(-) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java index 32a037edb99e..e7fb779bc960 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java @@ -156,7 +156,7 @@ public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, currentProxySCMNodeId = scmNodeIds.get(currentProxyIndex); scmClientConfig = conf.getObject(SCMClientConfig.class); - this.maxRetryCount = Math.max(scmClientConfig.getRetryCount(), scmNodeIds.size()); + this.maxRetryCount = scmClientConfig.getRetryCount(); this.retryInterval = scmClientConfig.getRetryInterval(); this.resolveOnFailureEnabled = conf.getBoolean( OzoneConfigKeys.OZONE_CLIENT_FAILOVER_RESOLVE_NEEDED_KEY, diff --git a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderRefreshWired.java b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderRefreshWired.java index 1667a6c82779..83a08a7d8df7 100644 --- a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderRefreshWired.java +++ b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderRefreshWired.java @@ -31,11 +31,9 @@ import java.net.InetSocketAddress; import java.net.SocketTimeoutException; import java.util.ArrayList; -import java.util.Arrays; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.ratis.ServerNotLeaderException; import org.apache.hadoop.io.retry.RetryPolicy; -import org.apache.hadoop.io.retry.RetryPolicy.RetryAction.RetryDecision; import org.apache.hadoop.ozone.ha.ConfUtils; import org.apache.ozone.test.GenericTestUtils.LogCapturer; import org.apache.ratis.protocol.RaftPeerId; @@ -58,9 +56,6 @@ public class TestSCMFailoverProxyProviderRefreshWired { private static final String SCM_NODE_1 = "scm1"; private static final String SCM_NODE_2 = "scm2"; private static final String SCM_NODE_3 = "scm3"; - private static final String SCM_NODE_4 = "scm4"; - private static final String SCM_NODE_5 = "scm5"; - private static final String SCM_NODE_6 = "scm6"; private OzoneConfiguration conf; @@ -134,34 +129,6 @@ public void testApplicationLevelErrorDoesNotTriggerRefresh() throws Exception { "ServerNotLeaderException is application-level; refresh must NOT fire"); } - @Test - public void testRetryCountCoversAllConfiguredScmNodes() throws Exception { - RetryPolicy policy = newProvider(6, 3).getRetryPolicy(); - - assertEquals(RetryDecision.FAILOVER_AND_RETRY, - policy.shouldRetry(new IOException("failure"), 0, 5, false).action); - assertEquals(RetryDecision.FAIL, - policy.shouldRetry(new IOException("failure"), 0, 6, false).action); - } - - @Test - public void testRetryCountIsUnchangedWhenItCoversScmNodes() throws Exception { - RetryPolicy policy = newProvider(3, 3).getRetryPolicy(); - - assertEquals(RetryDecision.FAIL, - policy.shouldRetry(new IOException("failure"), 0, 3, false).action); - } - - @Test - public void testConfiguredRetryCountRemainsMinimum() throws Exception { - RetryPolicy policy = newProvider(3, 5).getRetryPolicy(); - - assertEquals(RetryDecision.FAILOVER_AND_RETRY, - policy.shouldRetry(new IOException("failure"), 0, 4, false).action); - assertEquals(RetryDecision.FAIL, - policy.shouldRetry(new IOException("failure"), 0, 5, false).action); - } - @Test public void testFailoverToIpv6SuggestedLeader() { SCMBlockLocationFailoverProxyProvider provider = newThreeNodeProvider(); @@ -301,25 +268,6 @@ private SCMBlockLocationFailoverProxyProvider newThreeNodeProvider() { return new SCMBlockLocationFailoverProxyProvider(conf); } - private SCMBlockLocationFailoverProxyProvider newProvider(int nodeCount, - int retryCount) { - String[] nodeIds = { - SCM_NODE_1, SCM_NODE_2, SCM_NODE_3, SCM_NODE_4, SCM_NODE_5, SCM_NODE_6 - }; - conf.set(OZONE_SCM_NODES_KEY + "." + SCM_SERVICE_ID, - String.join(",", Arrays.copyOf(nodeIds, nodeCount))); - for (int i = 0; i < nodeCount; i++) { - conf.set(ConfUtils.addKeySuffixes(OZONE_SCM_ADDRESS_KEY, - SCM_SERVICE_ID, nodeIds[i]), "localhost"); - } - SCMClientConfig scmClientConfig = new SCMClientConfig(); - scmClientConfig.setRetryCount(retryCount); - scmClientConfig.setMaxRetryTimeout(retryCount * 1000L); - scmClientConfig.setRetryInterval(1000L); - conf.setFromObject(scmClientConfig); - return new SCMBlockLocationFailoverProxyProvider(conf); - } - private static SCMProxyInfo proxyInfoOf( SCMBlockLocationFailoverProxyProvider provider, String nodeId) { return provider.getSCMProxyInfoList().stream() diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMScmLocationClientConfig.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMScmLocationClientConfig.java index fec79304b7f4..4fff70789ec2 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMScmLocationClientConfig.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/OMScmLocationClientConfig.java @@ -33,6 +33,7 @@ import java.util.concurrent.TimeUnit; import javax.net.SocketFactory; import org.apache.hadoop.fs.CommonConfigurationKeysPublic; +import org.apache.hadoop.hdds.HddsUtils; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.net.NetUtils; import org.apache.hadoop.ozone.OzoneConfigKeys; @@ -40,24 +41,26 @@ /** * Builds the configuration used by OM's SCM location clients. * - *

The default effective retry limit is the larger of the configured retry - * count and retry timeout divided by retry interval: + *

The default effective retry limit is the largest of the configured retry + * count, retry timeout divided by retry interval, and configured SCM node + * count: *

- * max(3, 6s / 2s) = 3
+ * R = max(3, 6s / 2s, SCM node count)
  * 
- * For failures that always trigger failover, this allows four attempts. The - * configured timeout envelopes are: + * For failures that always trigger failover, this allows {@code R + 1} + * attempts. The configured timeout envelopes are: *
- * established connection: 4 * 30s + 3 * 2s = 126s
- * including TCP connect:   4 * (5s + 30s) + 3 * 2s = 146s
+ * established connection: (R + 1) * 30s + R * 2s
+ * including TCP connect:   (R + 1) * (5s + 30s) + R * 2s
  * 
* - *

Hadoop tracks retries and failovers separately. A sequence of three - * retry-without-failover responses followed by three failover failures can - * therefore make seven attempts. Its configured timeout envelopes are: + *

Hadoop tracks retries and failovers separately. A sequence of {@code R} + * retry-without-failover responses followed by {@code R} failover failures + * can therefore make {@code 2R + 1} attempts. Its configured timeout + * envelopes are: *

- * established connection: 7 * 30s + 6 * 2s = 222s
- * including TCP connect:   7 * (5s + 30s) + 6 * 2s = 257s
+ * established connection: (2R + 1) * 30s + 2R * 2s
+ * including TCP connect:   (2R + 1) * (5s + 30s) + 2R * 2s
  * 
* Actual calls can complete sooner. These calculations are not caller-side * deadlines and exclude time outside the configured connect and RPC waits. @@ -87,6 +90,14 @@ public static OzoneConfiguration createScmClientConfiguration( OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY, OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY_DEFAULT, OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY); + String scmServiceId = HddsUtils.getScmServiceId(configuration); + int scmNodeCount = scmServiceId == null ? 1 : + HddsUtils.getSCMNodeIds(configuration, scmServiceId).size(); + int retryCount = scmClientConfiguration.getInt( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, 0); + scmClientConfiguration.setInt( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, + Math.max(retryCount, scmNodeCount)); copy(configuration, scmClientConfiguration, OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT, OZONE_OM_SCM_LOCATION_CLIENT_MAX_RETRY_TIMEOUT_DEFAULT, diff --git a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java index 327605f09041..d41b586c33b2 100644 --- a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java +++ b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/TestOMScmLocationClientConfig.java @@ -17,6 +17,8 @@ package org.apache.hadoop.ozone.om; +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.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY; import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_RETRY_INTERVAL; import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_SCM_LOCATION_CLIENT_IPC_CONNECT_TIMEOUT; @@ -154,6 +156,30 @@ void appliesOmDefaultsWhenOverridesAreNotSet() { 0, TimeUnit.SECONDS)); } + @Test + void boundsRetryCountByConfiguredScmNodes() { + OzoneConfiguration configuration = new OzoneConfiguration(); + String scmServiceId = "scmservice"; + configuration.set(OZONE_SCM_SERVICE_IDS_KEY, scmServiceId); + configuration.set(OZONE_SCM_NODES_KEY + "." + scmServiceId, + "scm1,scm2,scm3,scm4,scm5,scm6"); + configuration.setInt(OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY, 3); + + OzoneConfiguration scmClientConfiguration = + OMScmLocationClientConfig.createScmClientConfiguration(configuration); + + assertEquals(6, scmClientConfiguration.getInt( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, 0)); + assertEquals(3, configuration.getInt( + OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY, 0)); + + configuration.setInt(OZONE_OM_SCM_LOCATION_CLIENT_FAILOVER_MAX_RETRY, 8); + scmClientConfiguration = + OMScmLocationClientConfig.createScmClientConfiguration(configuration); + assertEquals(8, scmClientConfiguration.getInt( + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, 0)); + } + @Test void usesAnIpcClientSeparateFromTheDefaultSocketFactory() { OzoneConfiguration configuration = new OzoneConfiguration(); From 9c2f81cf879e9e95c16c52ba9cd6a1110cc4767a Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Tue, 15 Sep 2026 13:34:03 +0800 Subject: [PATCH 12/13] HDDS-16382. Remove duplicate SCM retry configuration --- hadoop-hdds/common/src/main/resources/ozone-default.xml | 8 -------- 1 file changed, 8 deletions(-) diff --git a/hadoop-hdds/common/src/main/resources/ozone-default.xml b/hadoop-hdds/common/src/main/resources/ozone-default.xml index b999ca31734f..6a0f5ca7e6e6 100644 --- a/hadoop-hdds/common/src/main/resources/ozone-default.xml +++ b/hadoop-hdds/common/src/main/resources/ozone-default.xml @@ -656,14 +656,6 @@ clients.
- - hdds.scmclient.failover.retry.interval - 2s - OZONE, SCM, CLIENT - - SCM Client timeout on waiting for the next connection retry to other SCM IP. - - ozone.om.http-address 0.0.0.0:9874 From 13485e7de47ff6341cba846b0607cb5d6a34f678 Mon Sep 17 00:00:00 2001 From: Ivan Andika Date: Tue, 15 Sep 2026 15:21:46 +0800 Subject: [PATCH 13/13] HDDS-16382. Exclude generated SCM retry configuration --- .../org/apache/hadoop/ozone/TestOzoneConfigurationFields.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestOzoneConfigurationFields.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestOzoneConfigurationFields.java index 418c62a1d12e..79965a993960 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestOzoneConfigurationFields.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/TestOzoneConfigurationFields.java @@ -187,8 +187,8 @@ private void addPropertiesNotInXml() { DatanodeConfiguration.HDDS_DATANODE_VOLUME_MIN_FREE_SPACE_PERCENT, OzoneConfigKeys.HDDS_SCM_CLIENT_RPC_TIME_OUT, OzoneConfigKeys.HDDS_SCM_CLIENT_MAX_RETRY_TIMEOUT, - OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_MAX_RETRY, + OzoneConfigKeys.HDDS_SCM_CLIENT_FAILOVER_RETRY_INTERVAL )); } } -