This is an automated email from the ASF dual-hosted git repository. lhotari pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/pulsar.git
commit fe64a9ec1d68e17f29066760c1ace746b1554d67 Author: sinan liu <[email protected]> AuthorDate: Tue Jun 30 00:58:18 2026 +0800 [fix][broker] Forward topic policy updates after init failures (#26110) (cherry picked from commit b14524eb672181124f8834a2048b8ed6a278385d) --- .../broker/service/persistent/PersistentTopic.java | 22 ++++++++--- .../service/persistent/PersistentTopicTest.java | 46 ++++++++++++++++++++++ 2 files changed, 62 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 26fe918a144..f015e7240b5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -4720,12 +4720,22 @@ public class PersistentTopic extends AbstractTopic implements Topic, AddEntryCal CompletableFuture<Optional<TopicPolicies>> localPoliciesFuture = topicPoliciesService.getTopicPoliciesAsync(partitionedTopicName, TopicPoliciesService.GetType.LOCAL_ONLY); - return globalPoliciesFuture.thenCombine(localPoliciesFuture, (global, local) -> { - // finally update the topic policies with the latest value or loaded value - return CompletableFuture.runAsync(() -> - topicPolicyListener.completeInitialization(global.orElse(null), local.orElse(null)), - getPoliciesNotifyThread()); - }).thenCompose(Function.identity()); + CompletableFuture<Void> initialPoliciesFuture = + globalPoliciesFuture.thenCombine(localPoliciesFuture, (global, local) -> { + // finally update the topic policies with the latest value or loaded value + return CompletableFuture.runAsync(() -> + topicPolicyListener.completeInitialization(global.orElse(null), + local.orElse(null)), + getPoliciesNotifyThread()); + }).thenCompose(Function.identity()); + return initialPoliciesFuture.exceptionallyCompose(ex -> + // The topic load path logs and continues when initial policy loading fails. Make sure the + // already-registered wrapper is not left buffering future live updates forever. + CompletableFuture.runAsync( + () -> topicPolicyListener.completeInitialization(null, null), + getPoliciesNotifyThread()) + .thenCompose(__ -> FutureUtil.failedFuture( + FutureUtil.unwrapCompletionException(ex)))); }); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java index 5b3c705688c..6b7a29ba3c0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java @@ -64,6 +64,7 @@ import org.apache.bookkeeper.client.LedgerHandle; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedLedger; +import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionFactory; import org.apache.bookkeeper.mledger.impl.ManagedCursorContainer; @@ -74,6 +75,7 @@ import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.BrokerTestBase; import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.service.TopicPoliciesService; +import org.apache.pulsar.broker.service.TopicPolicyListener; import org.apache.pulsar.broker.stats.prometheus.PrometheusMetricsClient.Metric; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.Consumer; @@ -97,6 +99,7 @@ import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.common.policies.data.TopicPolicies; import org.apache.pulsar.common.policies.data.TopicStats; import org.awaitility.Awaitility; +import org.mockito.ArgumentCaptor; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -733,6 +736,49 @@ public class PersistentTopicTest extends BrokerTestBase { TimeUnit.MINUTES.toMillis(1)); } + @Test + public void testTopicPolicyListenerForwardsLiveUpdatesAfterInitialLoadFailure() throws Exception { + class RecordingPersistentTopic extends PersistentTopic { + final List<TopicPolicies> receivedUpdates = new ArrayList<>(); + + RecordingPersistentTopic(String topic, ManagedLedger ledger, BrokerService brokerService) { + super(topic, ledger, brokerService); + } + + @Override + public void onUpdate(TopicPolicies policies) { + receivedUpdates.add(policies); + } + } + + final String topic = "persistent://prop/ns-abc/testTopicPolicyInitFailure-" + UUID.randomUUID(); + ManagedLedger ledger = mock(ManagedLedger.class); + doReturn(new ManagedLedgerConfig()).when(ledger).getConfig(); + doReturn(Collections.emptyMap()).when(ledger).getProperties(); + + TopicPoliciesService policiesService = mock(TopicPoliciesService.class); + doReturn(policiesService).when(pulsar).getTopicPoliciesService(); + doReturn(CompletableFuture.completedFuture(true)).when(policiesService) + .registerListenerAsync(any(TopicName.class), any(TopicPolicyListener.class)); + doReturn(CompletableFuture.failedFuture(new RuntimeException("initial topic policy load failed"))) + .when(policiesService).getTopicPoliciesAsync(any(TopicName.class), + any(TopicPoliciesService.GetType.class)); + + RecordingPersistentTopic persistentTopic = + new RecordingPersistentTopic(topic, ledger, pulsar.getBrokerService()); + persistentTopic.initTopicPolicy().handle((ignored, ex) -> null).get(3, TimeUnit.SECONDS); + + ArgumentCaptor<TopicPolicyListener> listenerCaptor = ArgumentCaptor.forClass(TopicPolicyListener.class); + verify(policiesService).registerListenerAsync(any(TopicName.class), listenerCaptor.capture()); + + TopicPolicies livePolicies = new TopicPolicies(); + livePolicies.setIsGlobal(false); + livePolicies.setMaxConsumerPerTopic(10); + listenerCaptor.getValue().onUpdate(livePolicies); + + assertEquals(persistentTopic.receivedUpdates, Collections.singletonList(livePolicies)); + } + @Test public void testDynamicConfigurationAutoSkipNonRecoverableData() throws Exception { pulsar.getConfiguration().setAutoSkipNonRecoverableData(false);
