From c3ecf1312f41f3d91beaaeacd5eeae7f3b615ca1 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Tue, 5 May 2026 08:39:37 -0700 Subject: [PATCH] [fix][test] Reduce admin client churn in ExtensibleLoadManagerTest.startBroker The @BeforeMethod startBroker() polled readiness by creating a new PulsarAdmin instance per broker on every poll iteration (3 brokers x N polls of build/close). When earlier tests stopped brokers (testStopBroker, testIsolationPolicy), the brokers in the next method's startBroker() are still warming up, and the per-tick connection setup contends with that warmup. Build the per-broker PulsarAdmin clients once before the await loop and reuse them across all poll iterations, closing them in a finally block. The readiness checks themselves (getActiveBrokers / createPartitioned- Topic / lookupPartitionedTopic) are unchanged. --- .../ExtensibleLoadManagerTest.java | 64 +++++++++++-------- 1 file changed, 39 insertions(+), 25 deletions(-) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/loadbalance/ExtensibleLoadManagerTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/loadbalance/ExtensibleLoadManagerTest.java index 38cda507fac97..fa0ba2ef9b123 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/loadbalance/ExtensibleLoadManagerTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/loadbalance/ExtensibleLoadManagerTest.java @@ -134,7 +134,7 @@ public void cleanup() { } @BeforeMethod(alwaysRun = true) - public void startBroker() { + public void startBroker() throws Exception { if (pulsarCluster == null) { return; } @@ -144,35 +144,49 @@ public void startBroker() { } }); String topicName = "persistent://" + DEFAULT_NAMESPACE + "/startBrokerCheck"; - Awaitility.await().atMost(180, TimeUnit.SECONDS).until( - () -> { - for (BrokerContainer brokerContainer : pulsarCluster.getBrokers()) { - try (PulsarAdmin brokerAdmin = PulsarAdmin.builder().serviceHttpUrl( - brokerContainer.getHttpServiceUrl()).build()) { - if (brokerAdmin.brokers().getActiveBrokers(clusterName).size() != NUM_BROKERS) { - log.info() + // Build admin clients once and reuse across poll iterations to avoid the per-tick + // connection churn (3 brokers x N polls of admin builder/close). Connection setup + // contends with brokers that are still warming up after a stop/restart. + List brokers = new ArrayList<>(pulsarCluster.getBrokers()); + List brokerAdmins = new ArrayList<>(brokers.size()); + try { + for (BrokerContainer brokerContainer : brokers) { + brokerAdmins.add(PulsarAdmin.builder() + .serviceHttpUrl(brokerContainer.getHttpServiceUrl()).build()); + } + Awaitility.await().atMost(180, TimeUnit.SECONDS).until( + () -> { + for (int i = 0; i < brokers.size(); i++) { + BrokerContainer brokerContainer = brokers.get(i); + PulsarAdmin brokerAdmin = brokerAdmins.get(i); + try { + if (brokerAdmin.brokers().getActiveBrokers(clusterName).size() != NUM_BROKERS) { + log.info() + .attr("broker", brokerContainer.getHostName()) + .attr("see", NUM_BROKERS) + .log("Broker does not see active brokers yet"); + return false; + } + try { + brokerAdmin.topics().createPartitionedTopic(topicName, 10); + } catch (PulsarAdminException.ConflictException e) { + // expected - topic already exists + } + brokerAdmin.lookups().lookupPartitionedTopic(topicName); + } catch (Exception e) { + log.warn() .attr("broker", brokerContainer.getHostName()) - .attr("see", NUM_BROKERS) - .log("Broker does not see active brokers yet"); + .attr("yet", e.getMessage()) + .log("Broker is not ready yet"); return false; } - try { - brokerAdmin.topics().createPartitionedTopic(topicName, 10); - } catch (PulsarAdminException.ConflictException e) { - // expected - topic already exists - } - brokerAdmin.lookups().lookupPartitionedTopic(topicName); - } catch (Exception e) { - log.warn() - .attr("broker", brokerContainer.getHostName()) - .attr("yet", e.getMessage()) - .log("Broker is not ready yet"); - return false; } + return true; } - return true; - } - ); + ); + } finally { + brokerAdmins.forEach(PulsarAdmin::close); + } } @Test(timeOut = 40 * 1000)