This is an automated email from the ASF dual-hosted git repository.
lhotari pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new b14524eb672 [fix][broker] Forward topic policy updates after init
failures (#26110)
b14524eb672 is described below
commit b14524eb672181124f8834a2048b8ed6a278385d
Author: sinan liu <[email protected]>
AuthorDate: Tue Jun 30 00:58:18 2026 +0800
[fix][broker] Forward topic policy updates after init failures (#26110)
---
.../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 89c34e417a1..c2f656c776d 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
@@ -4903,12 +4903,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 1b258ed53ad..535b0ae75ad 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
@@ -65,6 +65,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.PositionBound;
import org.apache.bookkeeper.mledger.PositionFactory;
@@ -76,6 +77,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;
@@ -100,6 +102,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;
@@ -860,6 +863,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);