From b709b238ec0309543556d90eeef41e698bc4575e Mon Sep 17 00:00:00 2001 From: penghui Date: Tue, 2 Jun 2020 17:55:08 +0800 Subject: [PATCH 1/4] Only close active consumer for Failover subscription when seek(). --- ...bstractDispatcherSingleActiveConsumer.java | 9 ++ .../pulsar/broker/service/Dispatcher.java | 5 ++ ...PersistentDispatcherMultipleConsumers.java | 5 ++ ...PersistentDispatcherMultipleConsumers.java | 5 ++ .../persistent/PersistentSubscription.java | 2 +- .../PrecisTopicPublishRateThrottleTest.java | 6 +- .../broker/service/SubscriptionSeekTest.java | 83 +++++++++++++++++++ 7 files changed, 111 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java index a093d07557ea0..6c5f8a723638b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java @@ -245,6 +245,15 @@ public synchronized CompletableFuture disconnectAllConsumers(boolean isRes return closeFuture; } + public synchronized CompletableFuture disconnectActiveConsumers(boolean isResetCursor) { + closeFuture = new CompletableFuture<>(); + if (activeConsumer != null) { + activeConsumer.disconnect(isResetCursor); + } + closeFuture.complete(null); + return closeFuture; + } + @Override public synchronized void resetCloseFuture() { closeFuture = null; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Dispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Dispatcher.java index cffa024360c45..7b789e6da04bc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Dispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Dispatcher.java @@ -55,6 +55,11 @@ public interface Dispatcher { boolean isClosed(); + /** + * Disconnect active consumers + */ + CompletableFuture disconnectActiveConsumers(boolean isResetCursor); + /** * disconnect all consumers * diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java index 2053f5c51c2d7..9648e2c9be94b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java @@ -174,6 +174,11 @@ public synchronized CompletableFuture disconnectAllConsumers(boolean isRes return closeFuture; } + @Override + public CompletableFuture disconnectActiveConsumers(boolean isResetCursor) { + return disconnectAllConsumers(isResetCursor); + } + @Override public synchronized void resetCloseFuture() { closeFuture = null; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 1402b40e7fbea..c7af33817a233 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -419,6 +419,11 @@ public synchronized CompletableFuture disconnectAllConsumers(boolean isRes return closeFuture; } + @Override + public CompletableFuture disconnectActiveConsumers(boolean isResetCursor) { + return disconnectAllConsumers(isResetCursor); + } + @Override public synchronized void resetCloseFuture() { closeFuture = null; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index c2078326f5875..1ba4ce73a1881 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -698,7 +698,7 @@ private void resetCursor(Position finalPosition, CompletableFuture future) // Lock the Subscription object before locking the Dispatcher object to avoid deadlocks synchronized (this) { if (dispatcher != null && dispatcher.isConsumerConnected()) { - disconnectFuture = dispatcher.disconnectAllConsumers(true); + disconnectFuture = dispatcher.disconnectActiveConsumers(true); } else { disconnectFuture = CompletableFuture.completedFuture(null); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PrecisTopicPublishRateThrottleTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PrecisTopicPublishRateThrottleTest.java index c7a02aa812615..31130fafe9fa1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PrecisTopicPublishRateThrottleTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PrecisTopicPublishRateThrottleTest.java @@ -45,7 +45,7 @@ public void testPrecisTopicPublishRateLimitingDisabled() throws Exception { PublishRate publishRate = new PublishRate(1,10); // disable precis topic publish rate limiting conf.setPreciseTopicPublishRateLimiterEnable(false); - conf.setMaxPendingPublishdRequestsPerConnection(0); + conf.setMaxPendingPublishRequestsPerConnection(0); super.baseSetup(); final String topic = "persistent://prop/ns-abc/testPrecisTopicPublishRateLimiting"; org.apache.pulsar.client.api.Producer producer = pulsarClient.newProducer() @@ -84,7 +84,7 @@ public void testPrecisTopicPublishRateLimitingDisabled() throws Exception { public void testProducerBlockedByPrecisTopicPublishRateLimiting() throws Exception { PublishRate publishRate = new PublishRate(1,10); conf.setPreciseTopicPublishRateLimiterEnable(true); - conf.setMaxPendingPublishdRequestsPerConnection(0); + conf.setMaxPendingPublishRequestsPerConnection(0); super.baseSetup(); final String topic = "persistent://prop/ns-abc/testPrecisTopicPublishRateLimiting"; org.apache.pulsar.client.api.Producer producer = pulsarClient.newProducer() @@ -116,7 +116,7 @@ public void testProducerBlockedByPrecisTopicPublishRateLimiting() throws Excepti public void testPrecisTopicPublishRateLimitingProduceRefresh() throws Exception { PublishRate publishRate = new PublishRate(1,10); conf.setPreciseTopicPublishRateLimiterEnable(true); - conf.setMaxPendingPublishdRequestsPerConnection(0); + conf.setMaxPendingPublishRequestsPerConnection(0); super.baseSetup(); final String topic = "persistent://prop/ns-abc/testPrecisTopicPublishRateLimiting"; org.apache.pulsar.client.api.Producer producer = pulsarClient.newProducer() diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SubscriptionSeekTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SubscriptionSeekTest.java index f1b58ff441b50..54afc0e911102 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SubscriptionSeekTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SubscriptionSeekTest.java @@ -19,11 +19,15 @@ package org.apache.pulsar.broker.service; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; +import java.util.Set; import java.util.concurrent.TimeUnit; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; @@ -31,6 +35,7 @@ import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.RelativeTimeUtil; import org.testng.annotations.AfterClass; @@ -199,4 +204,82 @@ public void testSeekTimeOnPartitionedTopic() throws Exception { assertEquals(backlogs, 10); } + @Test + public void testShouldCloseAllConsumersForMultipleConsumerDispatcherWhenSeek() throws Exception { + final String topicName = "persistent://prop/use/ns-abc/testShouldCloseAllConsumersForMultipleConsumerDispatcherWhenSeek"; + // Disable pre-fetch in consumer to track the messages received + org.apache.pulsar.client.api.Consumer consumer1 = pulsarClient.newConsumer() + .topic(topicName) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("my-subscription") + .subscribe(); + + pulsarClient.newConsumer() + .topic(topicName) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("my-subscription") + .subscribe(); + + PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); + assertNotNull(topicRef); + assertEquals(topicRef.getSubscriptions().size(), 1); + List consumers = topicRef.getSubscriptions().get("my-subscription").getConsumers(); + assertEquals(consumers.size(), 2); + Set connectedSinceSet = new HashSet<>(); + for (Consumer consumer : consumers) { + connectedSinceSet.add(consumer.getStats().getConnectedSince()); + } + assertEquals(connectedSinceSet.size(), 2); + consumer1.seek(MessageId.earliest); + // Wait for consumer to reconnect + Thread.sleep(1000); + + consumers = topicRef.getSubscriptions().get("my-subscription").getConsumers(); + assertEquals(consumers.size(), 2); + for (Consumer consumer : consumers) { + assertFalse(connectedSinceSet.contains(consumer.getStats().getConnectedSince())); + } + } + + @Test + public void testOnlyCloseActiveConsumerForSingleActiveConsumerDispatcherWhenSeek() throws Exception { + final String topicName = "persistent://prop/use/ns-abc/testOnlyCloseActiveConsumerForSingleActiveConsumerDispatcherWhenSeek"; + // Disable pre-fetch in consumer to track the messages received + org.apache.pulsar.client.api.Consumer consumer1 = pulsarClient.newConsumer() + .topic(topicName) + .subscriptionType(SubscriptionType.Failover) + .subscriptionName("my-subscription") + .subscribe(); + + pulsarClient.newConsumer() + .topic(topicName) + .subscriptionType(SubscriptionType.Failover) + .subscriptionName("my-subscription") + .subscribe(); + + PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); + assertNotNull(topicRef); + assertEquals(topicRef.getSubscriptions().size(), 1); + List consumers = topicRef.getSubscriptions().get("my-subscription").getConsumers(); + assertEquals(consumers.size(), 2); + Set connectedSinceSet = new HashSet<>(); + for (Consumer consumer : consumers) { + connectedSinceSet.add(consumer.getStats().getConnectedSince()); + } + assertEquals(connectedSinceSet.size(), 2); + consumer1.seek(MessageId.earliest); + // Wait for consumer to reconnect + Thread.sleep(1000); + + consumers = topicRef.getSubscriptions().get("my-subscription").getConsumers(); + assertEquals(consumers.size(), 2); + + boolean hasConsumerNotDisconnected = false; + for (Consumer consumer : consumers) { + if (connectedSinceSet.contains(consumer.getStats().getConnectedSince())) { + hasConsumerNotDisconnected = true; + } + } + assertTrue(hasConsumerNotDisconnected); + } } From e391fa2e8c6e4bd0de3f2c740ad3fc9dfc783b8b Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 3 Jun 2020 11:40:47 +0800 Subject: [PATCH 2/4] Fix tests. --- .../java/org/apache/pulsar/io/batch/BatchSourceExecutor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-io/batch/src/main/java/org/apache/pulsar/io/batch/BatchSourceExecutor.java b/pulsar-io/batch/src/main/java/org/apache/pulsar/io/batch/BatchSourceExecutor.java index 88a011e324bc6..e624db9788141 100644 --- a/pulsar-io/batch/src/main/java/org/apache/pulsar/io/batch/BatchSourceExecutor.java +++ b/pulsar-io/batch/src/main/java/org/apache/pulsar/io/batch/BatchSourceExecutor.java @@ -25,7 +25,7 @@ import org.apache.pulsar.functions.api.Record; import org.apache.pulsar.functions.utils.Actions; import org.apache.pulsar.functions.utils.FunctionCommon; -import org.apache.pulsar.functions.utils.Reflections; +import org.apache.pulsar.common.util.Reflections; import org.apache.pulsar.functions.utils.SourceConfigUtils; import org.apache.pulsar.io.core.*; From 168aa674eea7cef82e28f42869d3028899a18aa9 Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 3 Jun 2020 13:35:59 +0800 Subject: [PATCH 3/4] Clean disk --- .github/workflows/ci-integration-backwards-compatibility.yaml | 2 +- .github/workflows/ci-integration-cli.yaml | 2 +- .github/workflows/ci-integration-function-state.yaml | 2 +- .github/workflows/ci-integration-messaging.yaml | 2 +- .github/workflows/ci-integration-schema.yaml | 2 +- .github/workflows/ci-integration-sql.yaml | 2 +- .github/workflows/ci-integration-standalone.yaml | 2 +- .github/workflows/ci-integration-tiered-filesystem.yaml | 2 +- .github/workflows/ci-integration-tiered-jcloud.yaml | 2 +- 9 files changed, 9 insertions(+), 9 deletions(-) diff --git a/.github/workflows/ci-integration-backwards-compatibility.yaml b/.github/workflows/ci-integration-backwards-compatibility.yaml index c9a5ced27e6d9..58fba8cd6b81e 100644 --- a/.github/workflows/ci-integration-backwards-compatibility.yaml +++ b/.github/workflows/ci-integration-backwards-compatibility.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h diff --git a/.github/workflows/ci-integration-cli.yaml b/.github/workflows/ci-integration-cli.yaml index f74510b5e556a..4ec3991e0530b 100644 --- a/.github/workflows/ci-integration-cli.yaml +++ b/.github/workflows/ci-integration-cli.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h diff --git a/.github/workflows/ci-integration-function-state.yaml b/.github/workflows/ci-integration-function-state.yaml index 1574f7ec96f27..beb4b7ba0165f 100644 --- a/.github/workflows/ci-integration-function-state.yaml +++ b/.github/workflows/ci-integration-function-state.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h diff --git a/.github/workflows/ci-integration-messaging.yaml b/.github/workflows/ci-integration-messaging.yaml index 7db4a1b87ed74..907ba2d927170 100644 --- a/.github/workflows/ci-integration-messaging.yaml +++ b/.github/workflows/ci-integration-messaging.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h diff --git a/.github/workflows/ci-integration-schema.yaml b/.github/workflows/ci-integration-schema.yaml index c21dfe7fd481c..f563b7a6e6c6d 100644 --- a/.github/workflows/ci-integration-schema.yaml +++ b/.github/workflows/ci-integration-schema.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h diff --git a/.github/workflows/ci-integration-sql.yaml b/.github/workflows/ci-integration-sql.yaml index daefeda288178..5b4b0713467c6 100644 --- a/.github/workflows/ci-integration-sql.yaml +++ b/.github/workflows/ci-integration-sql.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h diff --git a/.github/workflows/ci-integration-standalone.yaml b/.github/workflows/ci-integration-standalone.yaml index 2b6a231750406..e6c6a21fab2dd 100644 --- a/.github/workflows/ci-integration-standalone.yaml +++ b/.github/workflows/ci-integration-standalone.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h diff --git a/.github/workflows/ci-integration-tiered-filesystem.yaml b/.github/workflows/ci-integration-tiered-filesystem.yaml index 472cbced4d066..95e3c4eceb720 100644 --- a/.github/workflows/ci-integration-tiered-filesystem.yaml +++ b/.github/workflows/ci-integration-tiered-filesystem.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h diff --git a/.github/workflows/ci-integration-tiered-jcloud.yaml b/.github/workflows/ci-integration-tiered-jcloud.yaml index 15944525b050b..36042e3895a92 100644 --- a/.github/workflows/ci-integration-tiered-jcloud.yaml +++ b/.github/workflows/ci-integration-tiered-jcloud.yaml @@ -53,7 +53,7 @@ jobs: if: steps.docs.outputs.changed_only == 'no' run: | sudo swapoff -a - sudo rm -f /swapfile + sudo rm -rf /swapfile /usr/share/dotnet /usr/local/lib/android /opt/ghc sudo apt clean docker rmi $(docker images -q) -f df -h From 77d1d7d175445fa0a411e5e8a049ed3198de7217 Mon Sep 17 00:00:00 2001 From: penghui Date: Wed, 3 Jun 2020 15:16:28 +0800 Subject: [PATCH 4/4] Fix flaky test in topic policies service test. --- .../SystemTopicBasedTopicPoliciesServiceTest.java | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java index 132d8b17557be..69b6a5ef91cd7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java @@ -73,16 +73,14 @@ public void testGetPolicy() throws ExecutionException, InterruptedException, Top TopicPolicies initPolicy = TopicPolicies.builder() .maxConsumerPerTopic(10) .build(); - systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, initPolicy); - - Assert.assertNull(systemTopicBasedTopicPoliciesService.getPoliciesCacheInit(TOPIC1.getNamespaceObject())); - - Thread.sleep(1000); - + systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, initPolicy).get(); Assert.assertTrue(systemTopicBasedTopicPoliciesService.getPoliciesCacheInit(TOPIC1.getNamespaceObject())); + // Wait for all topic policies updated. + Thread.sleep(3000); + // Assert broker is cache all topic policies - Assert.assertEquals(10, systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1).getMaxConsumerPerTopic().intValue()); + Assert.assertEquals(systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1).getMaxConsumerPerTopic().intValue(), 10); // Update policy for TOPIC1 TopicPolicies policies1 = TopicPolicies.builder()