diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java index 05b888a785479..7dd152827ad4d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageTest.java @@ -721,7 +721,7 @@ public void testConcurrentBatchMessageAck(BatcherBuilder builder) throws Excepti PersistentDispatcherMultipleConsumers dispatcher = (PersistentDispatcherMultipleConsumers) topic .getSubscription(subscriptionName).getDispatcher(); // check strategically to let ack-message receive by broker - retryStrategically((test) -> dispatcher.getConsumers().get(0).getUnackedMessages() == 0, 5, 150); + retryStrategically((test) -> dispatcher.getConsumers().get(0).getUnackedMessages() == 0, 50, 150); assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 0); executor.shutdown(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorGlobalNSTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorGlobalNSTest.java index da48dd9474aa1..aa5a843afcd0c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorGlobalNSTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorGlobalNSTest.java @@ -87,7 +87,7 @@ public void testRemoveLocalClusterOnGlobalNamespace() throws Exception { admin1.namespaces().setNamespaceReplicationClusters(namespace, Sets.newHashSet("r2", "r3")); MockedPulsarServiceBaseTest - .retryStrategically((test) -> !pulsar1.getBrokerService().getTopics().containsKey(topicName), 5, 150); + .retryStrategically((test) -> !pulsar1.getBrokerService().getTopics().containsKey(topicName), 50, 150); Assert.assertFalse(pulsar1.getBrokerService().getTopics().containsKey(topicName)); Assert.assertFalse(producer1.isConnected()); @@ -118,7 +118,7 @@ public void testForcefullyTopicDeletion() throws Exception { admin1.topics().delete(topicName, true); MockedPulsarServiceBaseTest - .retryStrategically((test) -> !pulsar1.getBrokerService().getTopics().containsKey(topicName), 5, 150); + .retryStrategically((test) -> !pulsar1.getBrokerService().getTopics().containsKey(topicName), 50, 150); Assert.assertFalse(pulsar1.getBrokerService().getTopics().containsKey(topicName)); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java index 6dfe06ae3b331..ebe79b13c0531 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java @@ -332,7 +332,7 @@ public void testAuthorizationWithAnonymousUser() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150)); + }, 50, 150)); // validate pulsar sink consumer has started on the topic assertEquals(admin1.topics().getStats(sourceTopic).subscriptions.size(), 1); @@ -352,7 +352,7 @@ public void testAuthorizationWithAnonymousUser() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); Message msg = consumer.receive(5, TimeUnit.SECONDS); String receivedPropertyValue = msg.getProperty(propertyKey); @@ -384,7 +384,7 @@ public void testAuthorizationWithAnonymousUser() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150)); + }, 50, 150)); // test getFunctionInfo try { @@ -603,7 +603,7 @@ public void testAuthorization() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150)); + }, 50, 150)); // validate pulsar sink consumer has started on the topic assertEquals(admin1.topics().getStats(sourceTopic).subscriptions.size(), 1); @@ -623,7 +623,7 @@ public void testAuthorization() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); Message msg = consumer.receive(5, TimeUnit.SECONDS); String receivedPropertyValue = msg.getProperty(propertyKey); @@ -655,7 +655,7 @@ public void testAuthorization() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150)); + }, 50, 150)); // test getFunctionInfo try { @@ -796,7 +796,7 @@ public void testAuthorization() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150)); + }, 50, 150)); } } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java index 72f158f677286..48bb83fb26aed 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java @@ -412,7 +412,7 @@ private void testE2EPulsarFunctionLocalRun(String jarFilePathUrl) throws Excepti } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate pulsar sink consumer has started on the topic TopicStats stats = admin.topics().getStats(sourceTopic); assertTrue(stats.subscriptions.get(subscriptionName) != null @@ -430,7 +430,7 @@ private void testE2EPulsarFunctionLocalRun(String jarFilePathUrl) throws Excepti } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); for (int i = 0; i < totalMsgs; i++) { Message msg = consumer.receive(5, TimeUnit.SECONDS); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java index eab8be6e89331..8a4928f73bdfc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java @@ -304,7 +304,7 @@ public void testPulsarFunctionState() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate pulsar sink consumer has started on the topic assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 1); @@ -320,7 +320,7 @@ public void testPulsarFunctionState() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); retryStrategically((test) -> { try { @@ -329,7 +329,7 @@ public void testPulsarFunctionState() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); for (int i = 0; i < 5; i++) { Message msg = consumer.receive(5, TimeUnit.SECONDS); @@ -353,7 +353,7 @@ public void testPulsarFunctionState() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // make sure subscriptions are cleanup assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 0); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarWorkerAssignmentTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarWorkerAssignmentTest.java index d5a6f08bcebb6..ca9ce3b028bb3 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarWorkerAssignmentTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarWorkerAssignmentTest.java @@ -193,7 +193,7 @@ public void testFunctionAssignments() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate 2 instances have been started assertEquals(admin.topics().getStats(sinkTopic).subscriptions.size(), 1); assertEquals(admin.topics().getStats(sinkTopic).subscriptions.values().iterator().next().consumers.size(), 2); @@ -210,7 +210,7 @@ public void testFunctionAssignments() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate pulsar sink consumer has started on the topic log.info("admin.topics().getStats(sinkTopic): {}", new Gson().toJson(admin.topics().getStats(sinkTopic))); assertEquals(admin.topics().getStats(sinkTopic).subscriptions.values().iterator().next().consumers.size(), 1); @@ -250,7 +250,7 @@ public void testFunctionAssignmentsWithRestart() throws Exception { } catch (Exception e) { return false; } - }, 5, 150); + }, 50, 150); // Validate registered assignments Map assignments = runtimeManager.getCurrentAssignments().values().iterator().next(); @@ -279,7 +279,7 @@ public void testFunctionAssignmentsWithRestart() throws Exception { } catch (Exception e) { return false; } - }, 5, 150); + }, 50, 150); // Validate registered assignments assignments = runtimeManager.getCurrentAssignments().values().iterator().next(); @@ -298,7 +298,7 @@ public void testFunctionAssignmentsWithRestart() throws Exception { } catch (Exception e) { return false; } - }, 5, 150); + }, 50, 150); // Validate registered assignments assignments = runtimeManager2.getCurrentAssignments().values().iterator().next(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java index 1d963e9b9ce44..c5b605a766a23 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java @@ -441,7 +441,7 @@ private void testE2EPulsarFunction(String jarFilePathUrl) throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate pulsar sink consumer has started on the topic assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 1); @@ -457,7 +457,7 @@ private void testE2EPulsarFunction(String jarFilePathUrl) throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); Message msg = consumer.receive(5, TimeUnit.SECONDS); String receivedPropertyValue = msg.getProperty(propertyKey); @@ -477,7 +477,7 @@ private void testE2EPulsarFunction(String jarFilePathUrl) throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // make sure subscriptions are cleanup assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 0); @@ -708,7 +708,7 @@ private void testPulsarSinkStats(String jarFilePathUrl) throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // make sure subscriptions are cleanup assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 0); @@ -905,7 +905,7 @@ public void testPulsarFunctionStats() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate pulsar sink consumer has started on the topic assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 1); @@ -1197,7 +1197,7 @@ public void testPulsarFunctionStats() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // make sure subscriptions are cleanup assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 0); @@ -1241,7 +1241,7 @@ public void testPulsarFunctionStatus() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate pulsar sink consumer has started on the topic assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 1); @@ -1284,7 +1284,7 @@ public void testPulsarFunctionStatus() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // make sure subscriptions are cleanup assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 0); @@ -1348,7 +1348,7 @@ public void testFunctionStopAndRestartApi() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); SubscriptionStats subStats = admin.topics().getStats(sourceTopic).subscriptions.get(subscriptionName); assertEquals(subStats.consumers.size(), 1); @@ -1363,7 +1363,7 @@ public void testFunctionStopAndRestartApi() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); subStats = admin.topics().getStats(sourceTopic).subscriptions.get(subscriptionName); assertEquals(subStats.consumers.size(), 0); @@ -1378,7 +1378,7 @@ public void testFunctionStopAndRestartApi() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); subStats = admin.topics().getStats(sourceTopic).subscriptions.get(subscriptionName); assertEquals(subStats.consumers.size(), 1); @@ -1421,7 +1421,7 @@ public void testFunctionAutomaticSubCleanup() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); assertFalse(admin.functions().getFunction(tenant, namespacePortion, functionName).getCleanupSubscription()); retryStrategically((test) -> { @@ -1430,7 +1430,7 @@ public void testFunctionAutomaticSubCleanup() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate pulsar source consumer has started on the topic assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 1); @@ -1444,7 +1444,7 @@ public void testFunctionAutomaticSubCleanup() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); assertTrue(admin.functions().getFunction(tenant, namespacePortion, functionName).getCleanupSubscription()); int totalMsgs = 10; @@ -1487,7 +1487,7 @@ public void testFunctionAutomaticSubCleanup() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // make sure subscriptions are cleanup assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 0); @@ -1503,7 +1503,7 @@ public void testFunctionAutomaticSubCleanup() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // validate pulsar source consumer has started on the topic assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 1); @@ -1514,7 +1514,7 @@ public void testFunctionAutomaticSubCleanup() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); assertFalse(admin.functions().getFunction(tenant, namespacePortion, functionName).getCleanupSubscription()); // test update another config and making sure that subscription cleanup remains unchanged @@ -1528,7 +1528,7 @@ public void testFunctionAutomaticSubCleanup() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); assertFalse(admin.functions().getFunction(tenant, namespacePortion, functionName).getCleanupSubscription()); // delete functions @@ -1540,7 +1540,7 @@ public void testFunctionAutomaticSubCleanup() throws Exception { } catch (PulsarAdminException e) { return false; } - }, 5, 150); + }, 50, 150); // make sure subscriptions are cleanup assertEquals(admin.topics().getStats(sourceTopic).subscriptions.size(), 1);