[ 
https://issues.apache.org/jira/browse/HDFS-17909?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18080287#comment-18080287
 ] 

ASF GitHub Bot commented on HDFS-17909:
---------------------------------------

KeeProMise commented on code in PR #8448:
URL: https://github.com/apache/hadoop/pull/8448#discussion_r3226213942


##########
hadoop-hdfs-project/hadoop-hdfs-rbf/src/test/java/org/apache/hadoop/hdfs/server/federation/router/async/TestRouterAsyncHandlerQueueOverflow.java:
##########
@@ -0,0 +1,179 @@
+/**
+ * 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
+ * <p>
+ * http://www.apache.org/licenses/LICENSE-2.0
+ * <p>
+ * 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.hdfs.server.federation.router.async;
+
+import java.util.EnumSet;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ThreadPoolExecutor;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+import org.mockito.Mockito;
+
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.hdfs.protocol.OpenFilesIterator;
+import org.apache.hadoop.hdfs.server.federation.MiniRouterDFSCluster;
+import org.apache.hadoop.hdfs.server.federation.RouterConfigBuilder;
+import 
org.apache.hadoop.hdfs.server.federation.fairness.RouterAsyncRpcFairnessPolicyController;
+import 
org.apache.hadoop.hdfs.server.federation.fairness.RouterRpcFairnessPolicyController;
+import 
org.apache.hadoop.hdfs.server.federation.resolver.FederationNamenodeContext;
+import org.apache.hadoop.hdfs.server.federation.router.RemoteMethod;
+import org.apache.hadoop.hdfs.server.federation.router.RemoteParam;
+import org.apache.hadoop.hdfs.server.federation.router.RouterRpcServer;
+import org.apache.hadoop.hdfs.server.namenode.FSNamesystem;
+import org.apache.hadoop.hdfs.server.namenode.NameNodeAdapterMockitoUtil;
+import org.apache.hadoop.ipc.StandbyException;
+import org.apache.hadoop.security.UserGroupInformation;
+import org.apache.hadoop.test.LambdaTestUtils;
+
+import static 
org.apache.hadoop.hdfs.server.federation.FederationTestUtils.NAMENODES;
+import static 
org.apache.hadoop.hdfs.server.federation.MiniRouterDFSCluster.DEFAULT_HEARTBEAT_INTERVAL_MS;
+import static 
org.apache.hadoop.hdfs.server.federation.router.RBFConfigKeys.DFS_ROUTER_ASYNC_RPC_HANDLER_COUNT_KEY;
+import static 
org.apache.hadoop.hdfs.server.federation.router.RBFConfigKeys.DFS_ROUTER_ASYNC_RPC_MAX_ASYNCCALL_PERMIT_KEY;
+import static 
org.apache.hadoop.hdfs.server.federation.router.RBFConfigKeys.DFS_ROUTER_ASYNC_RPC_QUEUE_SIZE;
+import static 
org.apache.hadoop.hdfs.server.federation.router.RBFConfigKeys.DFS_ROUTER_ASYNC_RPC_RESPONDER_COUNT_KEY;
+import static 
org.apache.hadoop.hdfs.server.federation.router.RBFConfigKeys.DFS_ROUTER_FAIRNESS_ACQUIRE_TIMEOUT;
+import static 
org.apache.hadoop.hdfs.server.federation.router.RBFConfigKeys.DFS_ROUTER_FAIRNESS_POLICY_CONTROLLER_CLASS;
+import static 
org.apache.hadoop.hdfs.server.federation.router.RBFConfigKeys.DFS_ROUTER_MONITOR_NAMENODE;
+import static 
org.apache.hadoop.hdfs.server.federation.router.async.utils.AsyncUtil.syncReturn;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+
+public class TestRouterAsyncHandlerQueueOverflow {
+  /**
+   * Federated HDFS cluster.
+   */
+  private static MiniRouterDFSCluster cluster;
+  private static String ns0;
+  private static RouterRpcServer routerRpcServer;
+  private static RouterAsyncRpcClient asyncRpcClient;
+
+  private final static int QUEUE_CAP = 2;
+  private static CountDownLatch testLatch;
+
+  @BeforeAll
+  public static void setUpCluster() throws Exception {
+    cluster = new MiniRouterDFSCluster(true, 2, 3, 
DEFAULT_HEARTBEAT_INTERVAL_MS, 1000);
+    cluster.startCluster();
+
+    // Making one Namenode active per nameservice
+    if (cluster.isHighAvailability()) {
+      for (String ns : cluster.getNameservices()) {
+        cluster.switchToActive(ns, NAMENODES[0]);
+        cluster.switchToStandby(ns, NAMENODES[1]);
+        cluster.switchToObserver(ns, NAMENODES[2]);
+      }
+    }
+    // Start routers with only an RPC service
+    Configuration routerConf = new 
RouterConfigBuilder().metrics().rpc().build();
+
+    routerConf.setInt(DFS_ROUTER_ASYNC_RPC_QUEUE_SIZE, QUEUE_CAP);
+    routerConf.setInt(DFS_ROUTER_ASYNC_RPC_HANDLER_COUNT_KEY, 1);
+    routerConf.setInt(DFS_ROUTER_ASYNC_RPC_RESPONDER_COUNT_KEY, 1);
+    routerConf.set(DFS_ROUTER_MONITOR_NAMENODE,
+        cluster.getNameservices().get(0) + "," + 
cluster.getNameservices().get(1));
+    routerConf.setClass(DFS_ROUTER_FAIRNESS_POLICY_CONTROLLER_CLASS,
+        RouterAsyncRpcFairnessPolicyController.class, 
RouterRpcFairnessPolicyController.class);
+    routerConf.setInt(DFS_ROUTER_FAIRNESS_ACQUIRE_TIMEOUT, 60000);
+    routerConf.setInt(DFS_ROUTER_ASYNC_RPC_MAX_ASYNCCALL_PERMIT_KEY, 1);
+    cluster.addRouterOverrides(routerConf);
+    cluster.startRouters();
+
+    cluster.registerNamenodes();
+    cluster.waitNamenodeRegistration();
+    cluster.waitActiveNamespaces();
+
+    testLatch = new CountDownLatch(1);
+    ns0 = cluster.getNameservices().get(0);
+    MiniRouterDFSCluster.NamenodeContext nn0 = cluster.getNamenode(ns0, null);
+    FSNamesystem spyNamesystem = 
NameNodeAdapterMockitoUtil.spyOnNamesystem(nn0.getNamenode());
+    // Mock one slow operation. Any public interface from FSNamesystem will do.
+    Mockito.doAnswer(invocationOnMock -> {
+      String invokePath = invocationOnMock.getArgument(1);
+      if (invokePath.startsWith("/veryBigOperation")) {
+        testLatch.await();
+      } else {
+        return invocationOnMock.callRealMethod();
+      }
+      return null;
+    }).when(spyNamesystem).getFilesBlockingDecom(anyLong(), anyString());
+
+    MiniRouterDFSCluster.RouterContext router = cluster.getRandomRouter();
+    routerRpcServer = router.getRouterRpcServer();
+    routerRpcServer.initAsyncThreadPools(routerConf);
+    asyncRpcClient = new RouterAsyncRpcClient(routerConf, router.getRouter(),
+        routerRpcServer.getNamenodeResolver(), routerRpcServer.getRPCMonitor(),
+        routerRpcServer.getRouterStateIdContext());
+  }
+
+  @AfterAll
+  public static void shutdownCluster() {
+    if (cluster != null) {
+      cluster.shutdown();
+    }
+  }
+
+  @Test
+  @Timeout(value = 10)
+  public void testInvokeMethodQueueOverflow() throws Exception {
+    RemoteMethod method =
+        new RemoteMethod("listOpenFiles", new Class<?>[] {long.class, 
EnumSet.class, String.class},
+            0, 
EnumSet.of(OpenFilesIterator.OpenFilesType.BLOCKING_DECOMMISSION),
+            new RemoteParam());
+    UserGroupInformation ugi = RouterRpcServer.getRemoteUser();
+    Class<?> protocol = method.getProtocol();
+    String bigPath = "/veryBigOperation";
+    Object[] params =
+        new Object[] {0, 
EnumSet.of(OpenFilesIterator.OpenFilesType.BLOCKING_DECOMMISSION),
+            bigPath};
+    List<? extends FederationNamenodeContext> namenodes =
+        asyncRpcClient.getOrderedNamenodes(ns0, true);
+    // Downstream namespace processing this huge request
+    asyncRpcClient.invokeMethod(ugi, namenodes, true, protocol, 
method.getMethod(), params);
+    ThreadPoolExecutor nsExecutor = 
routerRpcServer.getAsyncExecutorForNamespace(ns0);
+    Thread.sleep(500);
+    assertEquals(0, nsExecutor.getQueue().size());
+    // Successfully sent this request downstream, but all subsequent ones will 
get stuck
+    assertEquals(1, nsExecutor.getCompletedTaskCount());
+
+    // Async handler handling, blocking at acquirePermit
+    asyncRpcClient.invokeMethod(ugi, namenodes, true, protocol, 
method.getMethod(), params);
+    Thread.sleep(500);

Review Comment:
   Could we use GenericTestUtils.waitFor instead of sleep here?





> [ARR] AsyncRouterHandlerExecutors should use bounded queue
> ----------------------------------------------------------
>
>                 Key: HDFS-17909
>                 URL: https://issues.apache.org/jira/browse/HDFS-17909
>             Project: Hadoop HDFS
>          Issue Type: Improvement
>          Components: rbf
>            Reporter: Felix N
>            Assignee: Felix N
>            Priority: Major
>              Labels: pull-request-available
>
> AsyncRouterHandlerExecutors should use bounded queues to prevent queues from 
> growing without a limit.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to