From 0547a730f8a6fb5b4f9668a941104c0352e1335c Mon Sep 17 00:00:00 2001 From: Eugen Dueck Date: Tue, 4 Feb 2020 12:16:53 +0900 Subject: [PATCH 1/4] add feature BrokerDeduplicationAcrossProducers --- .../pulsar/broker/ServiceConfiguration.java | 9 + .../persistent/MessageDeduplication.java | 10 + .../persistent/MessageDuplicationTest.java | 315 ++++++++++++++++++ 3 files changed, 334 insertions(+) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 2771692b82f30..ef2efe3e19efb 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -338,6 +338,15 @@ public class ServiceConfiguration implements PulsarConfiguration { + " relative to a disconnected producer. Default is 6 hours.") private int brokerDeduplicationProducerInactivityTimeoutMinutes = 360; + @FieldContext( + category = CATEGORY_POLICIES, + doc = "Enable message deduplication across all producers.\n\n" + + "This can be overridden per-namespace. If enabled, brokers will reject" + + " messages with sequence ids that were already stored in the topic," + + " regardless of which producer sent the message. Default is false" + ) + private boolean brokerDeduplicationAcrossProducersEnabled = false; + @FieldContext( category = CATEGORY_POLICIES, doc = "When a namespace is created without specifying the number of bundle, this" diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index dc4604e8a401c..bdd21e9e2846b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -91,6 +91,7 @@ public MessageDupUnknownException() { } } + public static final String ALL_PRODUCERS = "__all"; private volatile Status status; @@ -118,6 +119,8 @@ public MessageDupUnknownException() { private final String replicatorPrefix; + private final boolean deduplicationAcrossProducersEnabled; + public MessageDeduplication(PulsarService pulsar, PersistentTopic topic, ManagedLedger managedLedger) { this.pulsar = pulsar; this.topic = topic; @@ -127,6 +130,7 @@ public MessageDeduplication(PulsarService pulsar, PersistentTopic topic, Managed this.maxNumberOfProducers = pulsar.getConfiguration().getBrokerDeduplicationMaxNumberOfProducers(); this.snapshotCounter = 0; this.replicatorPrefix = pulsar.getConfiguration().getReplicatorPrefix(); + this.deduplicationAcrossProducersEnabled = pulsar.getConfiguration().isBrokerDeduplicationAcrossProducersEnabled(); } private CompletableFuture recoverSequenceIdsMap() { @@ -298,6 +302,9 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade headersAndPayload.readerIndex(readerIndex); md.recycle(); } + if (deduplicationAcrossProducersEnabled) { + producerName = ALL_PRODUCERS; + } // Synchronize the get() and subsequent put() on the map. This would only be relevant if the producer // disconnects and re-connects very quickly. At that point the call can be coming from a different thread @@ -342,6 +349,9 @@ public void recordMessagePersisted(PublishContext publishContext, PositionImpl p sequenceId = publishContext.getOriginalSequenceId(); highestSequenceId = publishContext.getOriginalHighestSequenceId(); } + if (deduplicationAcrossProducersEnabled) { + producerName = ALL_PRODUCERS; + } highestSequencedPersisted.put(producerName, Math.max(highestSequenceId, sequenceId)); if (++snapshotCounter >= snapshotInterval) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java index 5cfdef839db5a..1fc217713cbc0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java @@ -47,6 +47,7 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; +import static org.testng.AssertJUnit.assertNull; @Slf4j public class MessageDuplicationTest { @@ -146,6 +147,153 @@ public void testIsDuplicate() { assertEquals(lastSequenceIdPushed.longValue(), 5); } + @Test + public void testIsDuplicateAcrossProducers() { + PulsarService pulsarService = mock(PulsarService.class); + ServiceConfiguration serviceConfiguration = new ServiceConfiguration(); + serviceConfiguration.setBrokerDeduplicationEntriesInterval(BROKER_DEDUPLICATION_ENTRIES_INTERVAL); + serviceConfiguration.setBrokerDeduplicationMaxNumberOfProducers(BROKER_DEDUPLICATION_MAX_NUMBER_PRODUCERS); + serviceConfiguration.setReplicatorPrefix(REPLICATOR_PREFIX); + serviceConfiguration.setBrokerDeduplicationAcrossProducersEnabled(true); + + doReturn(serviceConfiguration).when(pulsarService).getConfiguration(); + PersistentTopic persistentTopic = mock(PersistentTopic.class); + ManagedLedger managedLedger = mock(ManagedLedger.class); + MessageDeduplication messageDeduplication = spy(new MessageDeduplication(pulsarService, persistentTopic, managedLedger)); + doReturn(true).when(messageDeduplication).isEnabled(); + + String producerName1 = "producer1"; + ByteBuf byteBuf1 = getMessage(producerName1, 0); + Topic.PublishContext publishContext1 = getPublishContext(producerName1, 0); + + String producerName2 = "producer2"; + ByteBuf byteBuf2 = getMessage(producerName2, 1); + Topic.PublishContext publishContext2 = getPublishContext(producerName2, 1); + + MessageDeduplication.MessageDupStatus status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); + + Long lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName1); + assertNull(lastSequenceIdPushed); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 0); + + status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); + assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName2); + assertNull(lastSequenceIdPushed); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 1); + + byteBuf1 = getMessage(producerName1, 1); + publishContext1 = getPublishContext(producerName1, 1); + status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + // should expect unknown because highestSequencePersisted is empty + assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 1); + + byteBuf2 = getMessage(producerName2, 1); + publishContext2 = getPublishContext(producerName2, 1); + status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); + // should expect unknown because highestSequencePersisted is empty + assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 1); + + byteBuf1 = getMessage(producerName1, 5); + publishContext1 = getPublishContext(producerName1, 5); + status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 5); + + byteBuf2 = getMessage(producerName2, 6); + publishContext2 = getPublishContext(producerName2, 6); + status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); + assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + byteBuf1 = getMessage(producerName1, 0); + publishContext1 = getPublishContext(producerName1, 0); + status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + // should expect unknown because highestSequencePersisted is empty + assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + byteBuf2 = getMessage(producerName2, 0); + publishContext2 = getPublishContext(producerName2, 0); + status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); + // should expect unknown because highestSequencePersisted is empty + assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + // update highest sequence persisted + messageDeduplication.highestSequencedPersisted.put(MessageDeduplication.ALL_PRODUCERS, 6L); + + byteBuf1 = getMessage(producerName1, 0); + publishContext1 = getPublishContext(producerName1, 0); + status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + // now that highestSequencedPersisted, message with seqId of zero can be classified as a dup + assertEquals(status, MessageDeduplication.MessageDupStatus.Dup); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + byteBuf2 = getMessage(producerName2, 0); + publishContext2 = getPublishContext(producerName2, 0); + status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); + // now that highestSequencedPersisted, message with seqId of zero can be classified as a dup + assertEquals(status, MessageDeduplication.MessageDupStatus.Dup); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + // update highest sequence persisted + messageDeduplication.highestSequencedPushed.put(MessageDeduplication.ALL_PRODUCERS, 0L); + messageDeduplication.highestSequencedPersisted.put(MessageDeduplication.ALL_PRODUCERS, 0L); + byteBuf1 = getMessage(producerName1, 0); + publishContext1 = getPublishContext(producerName1, 1, 6); + status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + byteBuf2 = getMessage(producerName2, 0); + publishContext2 = getPublishContext(producerName2, 2, 6); + status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); + assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + publishContext1 = getPublishContext(producerName1, 4, 8); + status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); + assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + publishContext2 = getPublishContext(producerName2, 4, 8); + status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); + assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 6); + } + @Test public void testIsDuplicateWithFailure() { @@ -310,6 +458,173 @@ public Object answer(InvocationOnMock invocationOnMock) throws Throwable { } + @Test + public void testIsDuplicateAcrossProducersWithFailure() { + + PulsarService pulsarService = mock(PulsarService.class); + ServiceConfiguration serviceConfiguration = new ServiceConfiguration(); + serviceConfiguration.setBrokerDeduplicationEntriesInterval(BROKER_DEDUPLICATION_ENTRIES_INTERVAL); + serviceConfiguration.setBrokerDeduplicationMaxNumberOfProducers(BROKER_DEDUPLICATION_MAX_NUMBER_PRODUCERS); + serviceConfiguration.setReplicatorPrefix(REPLICATOR_PREFIX); + serviceConfiguration.setBrokerDeduplicationAcrossProducersEnabled(true); + + doReturn(serviceConfiguration).when(pulsarService).getConfiguration(); + + ManagedLedger managedLedger = mock(ManagedLedger.class); + MessageDeduplication messageDeduplication = spy(new MessageDeduplication(pulsarService, mock(PersistentTopic.class), managedLedger)); + doReturn(true).when(messageDeduplication).isEnabled(); + + + ScheduledExecutorService scheduledExecutorService = mock(ScheduledExecutorService.class); + + doAnswer(new Answer() { + @Override + public Object answer(InvocationOnMock invocationOnMock) throws Throwable { + Object[] args = invocationOnMock.getArguments(); + Runnable test = (Runnable) args[0]; + test.run(); + return null; + } + }).when(scheduledExecutorService).submit(any(Runnable.class)); + + BrokerService brokerService = mock(BrokerService.class); + doReturn(scheduledExecutorService).when(brokerService).executor(); + doReturn(pulsarService).when(brokerService).pulsar(); + + PersistentTopic persistentTopic = spy(new PersistentTopic("topic-1", brokerService, managedLedger, messageDeduplication)); + + String producerName1 = "producer1"; + ByteBuf byteBuf1 = getMessage(producerName1, 0); + Topic.PublishContext publishContext1 = getPublishContext(producerName1, 0); + + String producerName2 = "producer2"; + ByteBuf byteBuf2 = getMessage(producerName2, 1); + Topic.PublishContext publishContext2 = getPublishContext(producerName2, 1); + + persistentTopic.publishMessage(byteBuf1, publishContext1); + persistentTopic.addComplete(new PositionImpl(0, 1), publishContext1); + verify(managedLedger, times(1)).asyncAddEntry(any(ByteBuf.class), any(), any()); + // just to make sure nothing is stored per producer + Long lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName1); + assertNull(lastSequenceIdPushed); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 0); + // just to make sure nothing is stored per producer + lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(producerName1); + assertNull(lastSequenceIdPushed); + lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); + assertNotNull(lastSequenceIdPushed); + assertEquals(lastSequenceIdPushed.longValue(), 0); + + persistentTopic.publishMessage(byteBuf2, publishContext2); + persistentTopic.addComplete(new PositionImpl(0, 2), publishContext2); + verify(managedLedger, times(2)).asyncAddEntry(any(ByteBuf.class), any(), any()); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 1); + lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 1); + + byteBuf1 = getMessage(producerName1, 3); + publishContext1 = getPublishContext(producerName1, 3); + persistentTopic.publishMessage(byteBuf1, publishContext1); + persistentTopic.addComplete(new PositionImpl(0, 3), publishContext1); + verify(managedLedger, times(3)).asyncAddEntry(any(ByteBuf.class), any(), any()); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 3); + lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 3); + + byteBuf1 = getMessage(producerName1, 5); + publishContext1 = getPublishContext(producerName1, 5); + persistentTopic.publishMessage(byteBuf1, publishContext1); + persistentTopic.addComplete(new PositionImpl(0, 4), publishContext1); + verify(managedLedger, times(4)).asyncAddEntry(any(ByteBuf.class), any(), any()); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 5); + lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 5); + + // publish dup + byteBuf1 = getMessage(producerName1, 0); + publishContext1 = getPublishContext(producerName1, 0); + persistentTopic.publishMessage(byteBuf1, publishContext1); + verify(managedLedger, times(4)).asyncAddEntry(any(ByteBuf.class), any(), any()); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 5); + verify(publishContext1, times(1)).completed(eq(null), eq(-1L), eq(-1L)); + + // publish message unknown dup status + byteBuf1 = getMessage(producerName1, 6); + publishContext1 = getPublishContext(producerName1, 6); + // don't complete message + persistentTopic.publishMessage(byteBuf1, publishContext1); + verify(managedLedger, times(5)).asyncAddEntry(any(ByteBuf.class), any(), any()); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 6); + lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 5); + + // publish same message again + byteBuf1 = getMessage(producerName1, 6); + publishContext1 = getPublishContext(producerName1, 6); + persistentTopic.publishMessage(byteBuf1, publishContext1); + verify(managedLedger, times(5)).asyncAddEntry(any(ByteBuf.class), any(), any()); + verify(publishContext1, times(1)).completed(any(MessageDeduplication.MessageDupUnknownException.class), eq(-1L), eq(-1L)); + + // complete seq 6 message eventually + persistentTopic.addComplete(new PositionImpl(0, 5), publishContext1); + + // simulate failure + byteBuf1 = getMessage(producerName1, 7); + publishContext1 = getPublishContext(producerName1, 7); + persistentTopic.publishMessage(byteBuf1, publishContext1); + verify(managedLedger, times(6)).asyncAddEntry(any(ByteBuf.class), any(), any()); + + persistentTopic.addFailed(new ManagedLedgerException("test"), publishContext1); + // check highestSequencedPushed is reset + assertEquals(messageDeduplication.highestSequencedPushed.size(), 1); + assertEquals(messageDeduplication.highestSequencedPersisted.size(), 1); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertEquals(lastSequenceIdPushed.longValue(), 6); + lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); + assertEquals(lastSequenceIdPushed.longValue(), 6); + verify(messageDeduplication, times(1)).resetHighestSequenceIdPushed(); + + // try dup + byteBuf1 = getMessage(producerName1, 6); + publishContext1 = getPublishContext(producerName1, 6); + persistentTopic.publishMessage(byteBuf1, publishContext1); + verify(managedLedger, times(6)).asyncAddEntry(any(ByteBuf.class), any(), any()); + verify(publishContext1, times(1)).completed(eq(null), eq(-1L), eq(-1L)); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 6); + + // try new message + byteBuf1 = getMessage(producerName1, 8); + publishContext1 = getPublishContext(producerName1, 8); + persistentTopic.publishMessage(byteBuf1, publishContext1); + verify(managedLedger, times(7)).asyncAddEntry(any(ByteBuf.class), any(), any()); + persistentTopic.addComplete(new PositionImpl(0, 5), publishContext1); + lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 8); + lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); + assertTrue(lastSequenceIdPushed != null); + assertEquals(lastSequenceIdPushed.longValue(), 8); + + } + public ByteBuf getMessage(String producerName, long seqId) { PulsarApi.MessageMetadata messageMetadata = PulsarApi.MessageMetadata.newBuilder() .setProducerName(producerName).setSequenceId(seqId) From 88a71697b4cb01166365570f87afed5d0986a2d9 Mon Sep 17 00:00:00 2001 From: Eugen Dueck Date: Tue, 4 Feb 2020 15:58:44 +0900 Subject: [PATCH 2/4] Revert "add feature BrokerDeduplicationAcrossProducers" This reverts commit 0547a730f8a6fb5b4f9668a941104c0352e1335c. --- .../pulsar/broker/ServiceConfiguration.java | 9 - .../persistent/MessageDeduplication.java | 10 - .../persistent/MessageDuplicationTest.java | 315 ------------------ 3 files changed, 334 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index ef2efe3e19efb..2771692b82f30 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -338,15 +338,6 @@ public class ServiceConfiguration implements PulsarConfiguration { + " relative to a disconnected producer. Default is 6 hours.") private int brokerDeduplicationProducerInactivityTimeoutMinutes = 360; - @FieldContext( - category = CATEGORY_POLICIES, - doc = "Enable message deduplication across all producers.\n\n" - + "This can be overridden per-namespace. If enabled, brokers will reject" - + " messages with sequence ids that were already stored in the topic," - + " regardless of which producer sent the message. Default is false" - ) - private boolean brokerDeduplicationAcrossProducersEnabled = false; - @FieldContext( category = CATEGORY_POLICIES, doc = "When a namespace is created without specifying the number of bundle, this" diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index bdd21e9e2846b..dc4604e8a401c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -91,7 +91,6 @@ public MessageDupUnknownException() { } } - public static final String ALL_PRODUCERS = "__all"; private volatile Status status; @@ -119,8 +118,6 @@ public MessageDupUnknownException() { private final String replicatorPrefix; - private final boolean deduplicationAcrossProducersEnabled; - public MessageDeduplication(PulsarService pulsar, PersistentTopic topic, ManagedLedger managedLedger) { this.pulsar = pulsar; this.topic = topic; @@ -130,7 +127,6 @@ public MessageDeduplication(PulsarService pulsar, PersistentTopic topic, Managed this.maxNumberOfProducers = pulsar.getConfiguration().getBrokerDeduplicationMaxNumberOfProducers(); this.snapshotCounter = 0; this.replicatorPrefix = pulsar.getConfiguration().getReplicatorPrefix(); - this.deduplicationAcrossProducersEnabled = pulsar.getConfiguration().isBrokerDeduplicationAcrossProducersEnabled(); } private CompletableFuture recoverSequenceIdsMap() { @@ -302,9 +298,6 @@ public MessageDupStatus isDuplicate(PublishContext publishContext, ByteBuf heade headersAndPayload.readerIndex(readerIndex); md.recycle(); } - if (deduplicationAcrossProducersEnabled) { - producerName = ALL_PRODUCERS; - } // Synchronize the get() and subsequent put() on the map. This would only be relevant if the producer // disconnects and re-connects very quickly. At that point the call can be coming from a different thread @@ -349,9 +342,6 @@ public void recordMessagePersisted(PublishContext publishContext, PositionImpl p sequenceId = publishContext.getOriginalSequenceId(); highestSequenceId = publishContext.getOriginalHighestSequenceId(); } - if (deduplicationAcrossProducersEnabled) { - producerName = ALL_PRODUCERS; - } highestSequencedPersisted.put(producerName, Math.max(highestSequenceId, sequenceId)); if (++snapshotCounter >= snapshotInterval) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java index 1fc217713cbc0..5cfdef839db5a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/MessageDuplicationTest.java @@ -47,7 +47,6 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; -import static org.testng.AssertJUnit.assertNull; @Slf4j public class MessageDuplicationTest { @@ -147,153 +146,6 @@ public void testIsDuplicate() { assertEquals(lastSequenceIdPushed.longValue(), 5); } - @Test - public void testIsDuplicateAcrossProducers() { - PulsarService pulsarService = mock(PulsarService.class); - ServiceConfiguration serviceConfiguration = new ServiceConfiguration(); - serviceConfiguration.setBrokerDeduplicationEntriesInterval(BROKER_DEDUPLICATION_ENTRIES_INTERVAL); - serviceConfiguration.setBrokerDeduplicationMaxNumberOfProducers(BROKER_DEDUPLICATION_MAX_NUMBER_PRODUCERS); - serviceConfiguration.setReplicatorPrefix(REPLICATOR_PREFIX); - serviceConfiguration.setBrokerDeduplicationAcrossProducersEnabled(true); - - doReturn(serviceConfiguration).when(pulsarService).getConfiguration(); - PersistentTopic persistentTopic = mock(PersistentTopic.class); - ManagedLedger managedLedger = mock(ManagedLedger.class); - MessageDeduplication messageDeduplication = spy(new MessageDeduplication(pulsarService, persistentTopic, managedLedger)); - doReturn(true).when(messageDeduplication).isEnabled(); - - String producerName1 = "producer1"; - ByteBuf byteBuf1 = getMessage(producerName1, 0); - Topic.PublishContext publishContext1 = getPublishContext(producerName1, 0); - - String producerName2 = "producer2"; - ByteBuf byteBuf2 = getMessage(producerName2, 1); - Topic.PublishContext publishContext2 = getPublishContext(producerName2, 1); - - MessageDeduplication.MessageDupStatus status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); - assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); - - Long lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName1); - assertNull(lastSequenceIdPushed); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 0); - - status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); - assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName2); - assertNull(lastSequenceIdPushed); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 1); - - byteBuf1 = getMessage(producerName1, 1); - publishContext1 = getPublishContext(producerName1, 1); - status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); - // should expect unknown because highestSequencePersisted is empty - assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 1); - - byteBuf2 = getMessage(producerName2, 1); - publishContext2 = getPublishContext(producerName2, 1); - status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); - // should expect unknown because highestSequencePersisted is empty - assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 1); - - byteBuf1 = getMessage(producerName1, 5); - publishContext1 = getPublishContext(producerName1, 5); - status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); - assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 5); - - byteBuf2 = getMessage(producerName2, 6); - publishContext2 = getPublishContext(producerName2, 6); - status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); - assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - byteBuf1 = getMessage(producerName1, 0); - publishContext1 = getPublishContext(producerName1, 0); - status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); - // should expect unknown because highestSequencePersisted is empty - assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - byteBuf2 = getMessage(producerName2, 0); - publishContext2 = getPublishContext(producerName2, 0); - status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); - // should expect unknown because highestSequencePersisted is empty - assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - // update highest sequence persisted - messageDeduplication.highestSequencedPersisted.put(MessageDeduplication.ALL_PRODUCERS, 6L); - - byteBuf1 = getMessage(producerName1, 0); - publishContext1 = getPublishContext(producerName1, 0); - status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); - // now that highestSequencedPersisted, message with seqId of zero can be classified as a dup - assertEquals(status, MessageDeduplication.MessageDupStatus.Dup); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - byteBuf2 = getMessage(producerName2, 0); - publishContext2 = getPublishContext(producerName2, 0); - status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); - // now that highestSequencedPersisted, message with seqId of zero can be classified as a dup - assertEquals(status, MessageDeduplication.MessageDupStatus.Dup); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - // update highest sequence persisted - messageDeduplication.highestSequencedPushed.put(MessageDeduplication.ALL_PRODUCERS, 0L); - messageDeduplication.highestSequencedPersisted.put(MessageDeduplication.ALL_PRODUCERS, 0L); - byteBuf1 = getMessage(producerName1, 0); - publishContext1 = getPublishContext(producerName1, 1, 6); - status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); - assertEquals(status, MessageDeduplication.MessageDupStatus.NotDup); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - byteBuf2 = getMessage(producerName2, 0); - publishContext2 = getPublishContext(producerName2, 2, 6); - status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); - assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - publishContext1 = getPublishContext(producerName1, 4, 8); - status = messageDeduplication.isDuplicate(publishContext1, byteBuf1); - assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - publishContext2 = getPublishContext(producerName2, 4, 8); - status = messageDeduplication.isDuplicate(publishContext2, byteBuf2); - assertEquals(status, MessageDeduplication.MessageDupStatus.Unknown); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 6); - } - @Test public void testIsDuplicateWithFailure() { @@ -458,173 +310,6 @@ public Object answer(InvocationOnMock invocationOnMock) throws Throwable { } - @Test - public void testIsDuplicateAcrossProducersWithFailure() { - - PulsarService pulsarService = mock(PulsarService.class); - ServiceConfiguration serviceConfiguration = new ServiceConfiguration(); - serviceConfiguration.setBrokerDeduplicationEntriesInterval(BROKER_DEDUPLICATION_ENTRIES_INTERVAL); - serviceConfiguration.setBrokerDeduplicationMaxNumberOfProducers(BROKER_DEDUPLICATION_MAX_NUMBER_PRODUCERS); - serviceConfiguration.setReplicatorPrefix(REPLICATOR_PREFIX); - serviceConfiguration.setBrokerDeduplicationAcrossProducersEnabled(true); - - doReturn(serviceConfiguration).when(pulsarService).getConfiguration(); - - ManagedLedger managedLedger = mock(ManagedLedger.class); - MessageDeduplication messageDeduplication = spy(new MessageDeduplication(pulsarService, mock(PersistentTopic.class), managedLedger)); - doReturn(true).when(messageDeduplication).isEnabled(); - - - ScheduledExecutorService scheduledExecutorService = mock(ScheduledExecutorService.class); - - doAnswer(new Answer() { - @Override - public Object answer(InvocationOnMock invocationOnMock) throws Throwable { - Object[] args = invocationOnMock.getArguments(); - Runnable test = (Runnable) args[0]; - test.run(); - return null; - } - }).when(scheduledExecutorService).submit(any(Runnable.class)); - - BrokerService brokerService = mock(BrokerService.class); - doReturn(scheduledExecutorService).when(brokerService).executor(); - doReturn(pulsarService).when(brokerService).pulsar(); - - PersistentTopic persistentTopic = spy(new PersistentTopic("topic-1", brokerService, managedLedger, messageDeduplication)); - - String producerName1 = "producer1"; - ByteBuf byteBuf1 = getMessage(producerName1, 0); - Topic.PublishContext publishContext1 = getPublishContext(producerName1, 0); - - String producerName2 = "producer2"; - ByteBuf byteBuf2 = getMessage(producerName2, 1); - Topic.PublishContext publishContext2 = getPublishContext(producerName2, 1); - - persistentTopic.publishMessage(byteBuf1, publishContext1); - persistentTopic.addComplete(new PositionImpl(0, 1), publishContext1); - verify(managedLedger, times(1)).asyncAddEntry(any(ByteBuf.class), any(), any()); - // just to make sure nothing is stored per producer - Long lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(producerName1); - assertNull(lastSequenceIdPushed); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 0); - // just to make sure nothing is stored per producer - lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(producerName1); - assertNull(lastSequenceIdPushed); - lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); - assertNotNull(lastSequenceIdPushed); - assertEquals(lastSequenceIdPushed.longValue(), 0); - - persistentTopic.publishMessage(byteBuf2, publishContext2); - persistentTopic.addComplete(new PositionImpl(0, 2), publishContext2); - verify(managedLedger, times(2)).asyncAddEntry(any(ByteBuf.class), any(), any()); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 1); - lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 1); - - byteBuf1 = getMessage(producerName1, 3); - publishContext1 = getPublishContext(producerName1, 3); - persistentTopic.publishMessage(byteBuf1, publishContext1); - persistentTopic.addComplete(new PositionImpl(0, 3), publishContext1); - verify(managedLedger, times(3)).asyncAddEntry(any(ByteBuf.class), any(), any()); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 3); - lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 3); - - byteBuf1 = getMessage(producerName1, 5); - publishContext1 = getPublishContext(producerName1, 5); - persistentTopic.publishMessage(byteBuf1, publishContext1); - persistentTopic.addComplete(new PositionImpl(0, 4), publishContext1); - verify(managedLedger, times(4)).asyncAddEntry(any(ByteBuf.class), any(), any()); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 5); - lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 5); - - // publish dup - byteBuf1 = getMessage(producerName1, 0); - publishContext1 = getPublishContext(producerName1, 0); - persistentTopic.publishMessage(byteBuf1, publishContext1); - verify(managedLedger, times(4)).asyncAddEntry(any(ByteBuf.class), any(), any()); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 5); - verify(publishContext1, times(1)).completed(eq(null), eq(-1L), eq(-1L)); - - // publish message unknown dup status - byteBuf1 = getMessage(producerName1, 6); - publishContext1 = getPublishContext(producerName1, 6); - // don't complete message - persistentTopic.publishMessage(byteBuf1, publishContext1); - verify(managedLedger, times(5)).asyncAddEntry(any(ByteBuf.class), any(), any()); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 6); - lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 5); - - // publish same message again - byteBuf1 = getMessage(producerName1, 6); - publishContext1 = getPublishContext(producerName1, 6); - persistentTopic.publishMessage(byteBuf1, publishContext1); - verify(managedLedger, times(5)).asyncAddEntry(any(ByteBuf.class), any(), any()); - verify(publishContext1, times(1)).completed(any(MessageDeduplication.MessageDupUnknownException.class), eq(-1L), eq(-1L)); - - // complete seq 6 message eventually - persistentTopic.addComplete(new PositionImpl(0, 5), publishContext1); - - // simulate failure - byteBuf1 = getMessage(producerName1, 7); - publishContext1 = getPublishContext(producerName1, 7); - persistentTopic.publishMessage(byteBuf1, publishContext1); - verify(managedLedger, times(6)).asyncAddEntry(any(ByteBuf.class), any(), any()); - - persistentTopic.addFailed(new ManagedLedgerException("test"), publishContext1); - // check highestSequencedPushed is reset - assertEquals(messageDeduplication.highestSequencedPushed.size(), 1); - assertEquals(messageDeduplication.highestSequencedPersisted.size(), 1); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertEquals(lastSequenceIdPushed.longValue(), 6); - lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); - assertEquals(lastSequenceIdPushed.longValue(), 6); - verify(messageDeduplication, times(1)).resetHighestSequenceIdPushed(); - - // try dup - byteBuf1 = getMessage(producerName1, 6); - publishContext1 = getPublishContext(producerName1, 6); - persistentTopic.publishMessage(byteBuf1, publishContext1); - verify(managedLedger, times(6)).asyncAddEntry(any(ByteBuf.class), any(), any()); - verify(publishContext1, times(1)).completed(eq(null), eq(-1L), eq(-1L)); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 6); - - // try new message - byteBuf1 = getMessage(producerName1, 8); - publishContext1 = getPublishContext(producerName1, 8); - persistentTopic.publishMessage(byteBuf1, publishContext1); - verify(managedLedger, times(7)).asyncAddEntry(any(ByteBuf.class), any(), any()); - persistentTopic.addComplete(new PositionImpl(0, 5), publishContext1); - lastSequenceIdPushed = messageDeduplication.highestSequencedPushed.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 8); - lastSequenceIdPushed = messageDeduplication.highestSequencedPersisted.get(MessageDeduplication.ALL_PRODUCERS); - assertTrue(lastSequenceIdPushed != null); - assertEquals(lastSequenceIdPushed.longValue(), 8); - - } - public ByteBuf getMessage(String producerName, long seqId) { PulsarApi.MessageMetadata messageMetadata = PulsarApi.MessageMetadata.newBuilder() .setProducerName(producerName).setSequenceId(seqId) From a0bfcbf6d34d7a023d674b1aa35fde7d2efaafca Mon Sep 17 00:00:00 2001 From: Eugen Dueck Date: Fri, 7 Feb 2020 19:03:16 +0900 Subject: [PATCH 3/4] add producer group mode --- .../pulsar/broker/service/AbstractTopic.java | 134 ++++++++++++------ .../broker/service/BacklogQuotaManager.java | 3 +- .../pulsar/broker/service/Producer.java | 14 +- .../pulsar/broker/service/ServerCnx.java | 8 +- .../apache/pulsar/broker/service/Topic.java | 3 +- .../nonpersistent/NonPersistentTopic.java | 24 ++-- .../service/persistent/PersistentTopic.java | 31 ++-- .../prometheus/NamespaceStatsAggregator.java | 2 +- .../broker/service/BatchMessageTest.java | 18 +-- .../broker/service/BrokerServiceTest.java | 16 +-- .../service/PersistentQueueE2ETest.java | 4 +- .../service/PersistentTopicE2ETest.java | 13 +- .../broker/service/PersistentTopicTest.java | 61 ++++---- .../broker/service/ResendRequestTest.java | 10 +- .../pulsar/broker/service/ServerCnxTest.java | 14 +- .../broker/service/SubscriptionSeekTest.java | 6 +- .../client/api/ProducerCreationTest.java | 89 ++++++++++++ .../impl/MessagePublishThrottlingTest.java | 12 +- .../impl/TopicPublishThrottlingInitTest.java | 2 +- .../pulsar/client/api/ProducerBuilder.java | 10 ++ .../pulsar/client/api/ProducerGroupMode.java | 44 ++++++ .../pulsar/client/impl/ProducerBase.java | 23 ++- .../client/impl/ProducerBuilderImpl.java | 7 + .../pulsar/client/impl/ProducerImpl.java | 5 +- .../impl/conf/ProducerConfigurationData.java | 3 + .../pulsar/common/api/proto/PulsarApi.java | 105 ++++++++++++++ .../pulsar/common/protocol/Commands.java | 6 +- pulsar-common/src/main/proto/PulsarApi.proto | 8 ++ 28 files changed, 497 insertions(+), 178 deletions(-) create mode 100644 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerGroupMode.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index 638c8c7937f1e..14d333e7c680f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -18,23 +18,27 @@ */ package org.apache.pulsar.broker.service; +import static com.google.common.base.Preconditions.checkArgument; import static org.apache.bookkeeper.mledger.impl.ManagedLedgerMBeanImpl.ENTRY_LATENCY_BUCKETS_USEC; import static org.apache.pulsar.broker.cache.ConfigurationCacheService.POLICIES; import com.google.common.base.MoreObjects; -import java.util.Map; -import java.util.Objects; + +import java.util.*; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.LongAdder; import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.stream.Stream; + +import com.google.common.collect.Sets; import org.apache.bookkeeper.mledger.util.StatsBuckets; import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.service.schema.SchemaRegistryService; import org.apache.pulsar.broker.service.schema.exceptions.IncompatibleSchemaException; import org.apache.pulsar.broker.stats.prometheus.metrics.Summary; +import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.PublishRate; @@ -52,7 +56,7 @@ public abstract class AbstractTopic implements Topic { protected final String topic; // Producers currently connected to this topic - protected final ConcurrentHashMap producers; + protected final ConcurrentHashMap> producerGroups; protected final BrokerService brokerService; @@ -88,7 +92,7 @@ public abstract class AbstractTopic implements Topic { public AbstractTopic(String topic, BrokerService brokerService) { this.topic = topic; this.brokerService = brokerService; - this.producers = new ConcurrentHashMap<>(); + this.producerGroups = new ConcurrentHashMap<>(); this.isFenced = false; this.replicatorPrefix = brokerService.pulsar().getConfiguration().getReplicatorPrefix(); this.lastActive = System.nanoTime(); @@ -117,21 +121,38 @@ protected boolean isProducersExceeded() { final int maxProducers = policies.max_producers_per_topic > 0 ? policies.max_producers_per_topic : brokerService.pulsar().getConfiguration().getMaxProducersPerTopic(); - if (maxProducers > 0 && maxProducers <= producers.size()) { - return true; - } - return false; + return maxProducers > 0 && maxProducers <= getProducers().count(); } - protected boolean hasLocalProducers() { - AtomicBoolean foundLocal = new AtomicBoolean(false); - producers.values().forEach(producer -> { - if (!producer.isRemote()) { - foundLocal.set(true); + @Override + public void removeProducer(Producer producer) { + checkArgument(producer.getTopic() == this); + + boolean[] removed = { false }; + producerGroups.computeIfPresent(producer.getProducerName(), (s, producerSet) -> { + if (producerSet instanceof HashSet) { // a non-exclusive producer group + if (producerSet.remove(producer)) { + removed[0] = true; + if (producerSet.size() == 0) + return null; + } + return producerSet; + } else { // an exclusive producer "group" + if (producerSet.contains(producer)) { + removed[0] = true; + return null; + } else { + return producerSet; + } } }); + if (removed[0]) { + handleProducerRemoved(producer); + } + } - return foundLocal.get(); + protected boolean hasLocalProducers() { + return getProducers().anyMatch(producer -> !producer.isRemote()); } @Override @@ -140,8 +161,8 @@ public String toString() { } @Override - public Map getProducers() { - return producers; + public Stream getProducers() { + return producerGroups.values().stream().flatMap(Collection::stream); } @@ -291,8 +312,8 @@ public void resetBrokerPublishCountAndEnableReadIfRequired(boolean doneBrokerRes * it sets cnx auto-readable if producer's cnx is disabled due to publish-throttling */ protected void enableProducerRead() { - if (producers != null) { - producers.values().forEach(producer -> producer.getCnx().enableCnxAutoRead()); + if (producerGroups != null) { + getProducers().forEach(producer -> producer.getCnx().enableCnxAutoRead()); } } @@ -313,32 +334,61 @@ protected void internalAddProducer(Producer producer) throws BrokerServiceExcept log.debug("[{}] {} Got request to create producer ", topic, producer.getProducerName()); } - Producer existProducer = producers.putIfAbsent(producer.getProducerName(), producer); - if (existProducer != null) { - tryOverwriteOldProducer(existProducer, producer); - } - } - - private void tryOverwriteOldProducer(Producer oldProducer, Producer newProducer) - throws BrokerServiceException { - boolean canOverwrite = false; - if (oldProducer.equals(newProducer) && !isUserProvidedProducerName(oldProducer) - && !isUserProvidedProducerName(newProducer) && newProducer.getEpoch() > oldProducer.getEpoch()) { - oldProducer.close(false); - canOverwrite = true; - } - if (canOverwrite) { - if(!producers.replace(newProducer.getProducerName(), oldProducer, newProducer)) { - // Met concurrent update, throw exception here so that client can try reconnect later. - throw new BrokerServiceException.NamingException("Producer with name '" + newProducer.getProducerName() - + "' replace concurrency error"); + // the following the variables are used to get state out the compute function + // (we want to get out of producers.compute as quickly as possible so as not to block other concurrent actions) + BrokerServiceException[] bse = {null}; + Collection oldProducers = new LinkedList<>(); + boolean parallelGroupMode = producer.getGroupMode() == PulsarApi.CommandProducer.GroupMode.Parallel; + producerGroups.compute(producer.getProducerName(), (s, producerSet) -> { + if (producerSet == null) { // no producer under that name and topic connected yet + if (parallelGroupMode) { + Set identityHashSet = Sets.newIdentityHashSet(); + identityHashSet.add(producer); + return identityHashSet; + } else { + return Collections.singleton(producer); // in reality a "set" that can contain only one element is enough here + } } else { - handleProducerRemoved(oldProducer); + Producer existingProducer = producerSet.iterator().next(); + boolean existingProducerIsExclusive = + existingProducer.getGroupMode() == PulsarApi.CommandProducer.GroupMode.Exclusive; + if (existingProducerIsExclusive) { // an exclusive producer is already connected under that producerName + if (parallelGroupMode) { + bse[0] = new BrokerServiceException.NamingException( + "Exclusive Producer with name '" + producer.getProducerName() + "' is already connected to topic"); + return producerSet; + } else { + if (!isUserProvidedProducerName(existingProducer) && !isUserProvidedProducerName(producer) + && producer.getEpoch() > existingProducer.getEpoch()) { + oldProducers.add(existingProducer); + return Collections.singleton(producer); + } else { + bse[0] = new BrokerServiceException.NamingException( + "Producer with name '" + producer.getProducerName() + "' is already connected to topic"); + return producerSet; + } + } + } else { // a non-exclusive producer is already connected under that producerName + if (parallelGroupMode) { + if (!producerSet.add(producer)) { + bse[0] = new BrokerServiceException.NamingException( + "Non-exclusive producer with name '" + producer.getProducerName() + "' and address '" + + producer.getCnx().clientAddress() + "' is already connected to topic"); + } + return producerSet; + } else { + oldProducers.addAll(producerSet); + return Collections.singleton(producer); + } + } } - } else { - throw new BrokerServiceException.NamingException( - "Producer with name '" + newProducer.getProducerName() + "' is already connected to topic"); + }); + for (Producer oldProducer : oldProducers) { + oldProducer.close(false); + handleProducerRemoved(oldProducer); } + if (bse[0] != null) + throw bse[0]; } private boolean isUserProvidedProducerName(Producer producer){ diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BacklogQuotaManager.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BacklogQuotaManager.java index 1d15b003d8c0f..b63e2fadf3528 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BacklogQuotaManager.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BacklogQuotaManager.java @@ -185,9 +185,8 @@ private void dropBacklog(PersistentTopic persistentTopic, BacklogQuota quota) { */ private void disconnectProducers(PersistentTopic persistentTopic) { List> futures = Lists.newArrayList(); - Map producers = persistentTopic.getProducers(); - producers.values().forEach(producer -> { + persistentTopic.getProducers().forEach(producer -> { log.info("Producer [{}] has exceeded backlog quota on topic [{}]. Disconnecting producer", producer.getProducerName(), persistentTopic.getName()); futures.add(producer.disconnect()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index b6c2b496c522c..e3c1f6d3a6f4d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -42,6 +42,7 @@ import org.apache.pulsar.broker.service.Topic.PublishContext; import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic; import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; import org.apache.pulsar.common.api.proto.PulsarApi.ServerError; @@ -63,6 +64,7 @@ public class Producer { private final String producerName; private final long epoch; private final boolean userProvidedProducerName; + private final PulsarApi.CommandProducer.GroupMode groupMode; private final long producerId; private final String appId; private Rate msgIn; @@ -88,13 +90,17 @@ public class Producer { private final SchemaVersion schemaVersion; public Producer(Topic topic, ServerCnx cnx, long producerId, String producerName, String appId, - boolean isEncrypted, Map metadata, SchemaVersion schemaVersion, long epoch, - boolean userProvidedProducerName) { + boolean isEncrypted, Map metadata, SchemaVersion schemaVersion, long epoch, + boolean userProvidedProducerName, PulsarApi.CommandProducer.GroupMode groupMode) + throws BrokerServiceException { this.topic = topic; this.cnx = cnx; this.producerId = producerId; this.producerName = checkNotNull(producerName); this.userProvidedProducerName = userProvidedProducerName; + this.groupMode = groupMode; + if (groupMode != PulsarApi.CommandProducer.GroupMode.Exclusive && !userProvidedProducerName) + throw new BrokerServiceException.NotAllowedException("producerName must be specified in non-exclusive group modes"); this.epoch = epoch; this.closeFuture = new CompletableFuture<>(); this.appId = appId; @@ -277,6 +283,10 @@ public ServerCnx getCnx() { return this.cnx; } + public PulsarApi.CommandProducer.GroupMode getGroupMode() { + return groupMode; + } + private static final class MessagePublishContext implements PublishContext, Runnable { private Producer producer; private long sequenceId; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 2a78508ed0304..fb38a91e44ed5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -825,6 +825,7 @@ protected void handleProducer(final CommandProducer cmdProducer) { final boolean isEncrypted = cmdProducer.getEncrypted(); final Map metadata = CommandUtils.metadataFromCommand(cmdProducer); final SchemaData schema = cmdProducer.hasSchema() ? getSchema(cmdProducer.getSchema()) : null; + final CommandProducer.GroupMode groupMode = cmdProducer.getGroupMode(); TopicName topicName = validateTopicName(cmdProducer.getTopic(), requestId, cmdProducer); if (topicName == null) { @@ -942,10 +943,11 @@ protected void handleProducer(final CommandProducer cmdProducer) { }); schemaVersionFuture.thenAccept(schemaVersion -> { - Producer producer = new Producer(topic, ServerCnx.this, producerId, producerName, authRole, - isEncrypted, metadata, schemaVersion, epoch, userProvidedProducerName); - try { + Producer producer = new Producer(topic, ServerCnx.this, producerId, producerName, authRole, + isEncrypted, metadata, schemaVersion, epoch, userProvidedProducerName, + groupMode); + topic.addProducer(producer); if (isActive()) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index 26af1c1c5c8bf..e3fa8b06ee44c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -22,6 +22,7 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import java.util.stream.Stream; import org.apache.bookkeeper.mledger.Position; import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter; @@ -116,7 +117,7 @@ CompletableFuture createSubscription(String subscriptionName, Init CompletableFuture delete(); - Map getProducers(); + Stream getProducers(); String getName(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java index 8e3b14fa7f43a..dad9fd24469dc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java @@ -214,14 +214,6 @@ public void checkMessageDeduplicationInfo() { // No-op } - @Override - public void removeProducer(Producer producer) { - checkArgument(producer.getTopic() == this); - if (producers.remove(producer.getProducerName(), producer)) { - handleProducerRemoved(producer); - } - } - @Override public void handleProducerRemoved(Producer producer) { // decrement usage only if this was a valid producer close @@ -354,7 +346,7 @@ private CompletableFuture delete(boolean failIfHasSubscriptions, boolean c if (closeIfClientsConnected) { List> futures = Lists.newArrayList(); replicators.forEach((cluster, replicator) -> futures.add(replicator.disconnect())); - producers.values().forEach(producer -> futures.add(producer.disconnect())); + getProducers().forEach(producer -> futures.add(producer.disconnect())); subscriptions.forEach((s, sub) -> futures.add(sub.disconnect())); FutureUtil.waitForAll(futures).thenRun(() -> { closeClientFuture.complete(null); @@ -445,7 +437,7 @@ public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect List> futures = Lists.newArrayList(); replicators.forEach((cluster, replicator) -> futures.add(replicator.disconnect())); - producers.values().forEach(producer -> futures.add(producer.disconnect())); + getProducers().forEach(producer -> futures.add(producer.disconnect())); subscriptions.forEach((s, sub) -> futures.add(sub.disconnect())); CompletableFuture clientCloseFuture = closeWithoutWaitingClientDisconnect ? CompletableFuture.completedFuture(null) @@ -635,12 +627,12 @@ public void updateRates(NamespaceStats nsStats, NamespaceBundleStats bundleStats replicators.forEach((region, replicator) -> replicator.updateRates()); - nsStats.producerCount += producers.size(); - bundleStats.producerCount += producers.size(); + nsStats.producerCount += getProducers().count(); + bundleStats.producerCount += getProducers().count(); topicStatsStream.startObject(topic); topicStatsStream.startList("publishers"); - producers.values().forEach(producer -> { + getProducers().forEach(producer -> { producer.updateRates(); PublisherStats publisherStats = producer.getStats(); @@ -727,7 +719,7 @@ public void updateRates(NamespaceStats nsStats, NamespaceBundleStats bundleStats // Remaining dest stats. topicStats.averageMsgSize = topicStats.aggMsgRateIn == 0.0 ? 0.0 : (topicStats.aggMsgThroughputIn / topicStats.aggMsgRateIn); - topicStatsStream.writePair("producerCount", producers.size()); + topicStatsStream.writePair("producerCount", getProducers().count()); topicStatsStream.writePair("averageMsgSize", topicStats.averageMsgSize); topicStatsStream.writePair("msgRateIn", topicStats.aggMsgRateIn); topicStatsStream.writePair("msgRateOut", topicStats.aggMsgRateOut); @@ -757,7 +749,7 @@ public NonPersistentTopicStats getStats() { ObjectObjectHashMap remotePublishersStats = new ObjectObjectHashMap(); - producers.values().forEach(producer -> { + getProducers().forEach(producer -> { NonPersistentPublisherStats publisherStats = (NonPersistentPublisherStats) producer.getStats(); stats.msgRateIn += publisherStats.msgRateIn; stats.msgThroughputIn += publisherStats.msgThroughputIn; @@ -875,7 +867,7 @@ public CompletableFuture onPoliciesUpdate(Policies data) { isAllowAutoUpdateSchema = data.is_allow_auto_update_schema; schemaValidationEnforced = data.schema_validation_enforced; - producers.values().forEach(producer -> { + getProducers().forEach(producer -> { producer.checkPermissions(); producer.checkEncryption(); }); 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 b177b116a735b..254e39bfe2b3a 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 @@ -373,7 +373,7 @@ public synchronized void addFailed(ManagedLedgerException exception, Object ctx) isFenced = true; // close all producers List> futures = Lists.newArrayList(); - producers.values().forEach(producer -> futures.add(producer.disconnect())); + getProducers().forEach(producer -> futures.add(producer.disconnect())); FutureUtil.waitForAll(futures).handle((BiFunction) (aVoid, throwable) -> { decrementPendingWriteOpsAndCheck(); return null; @@ -435,7 +435,7 @@ public void addProducer(Producer producer) throws BrokerServiceException { private boolean hasRemoteProducers() { AtomicBoolean foundRemote = new AtomicBoolean(false); - producers.values().forEach(producer -> { + getProducers().forEach(producer -> { if (producer.isRemote()) { foundRemote.set(true); } @@ -478,15 +478,6 @@ private synchronized CompletableFuture closeReplProducersIfNoBacklog() { return FutureUtil.waitForAll(closeFutures); } - @Override - public void removeProducer(Producer producer) { - checkArgument(producer.getTopic() == this); - - if (producers.remove(producer.getProducerName(), producer)) { - handleProducerRemoved(producer); - } - } - @Override protected void handleProducerRemoved(Producer producer) { // decrement usage only if this was a valid producer close @@ -823,7 +814,7 @@ private CompletableFuture delete(boolean failIfHasSubscriptions, if (closeIfClientsConnected) { List> futures = Lists.newArrayList(); replicators.forEach((cluster, replicator) -> futures.add(replicator.disconnect())); - producers.values().forEach(producer -> futures.add(producer.disconnect())); + getProducers().forEach(producer -> futures.add(producer.disconnect())); subscriptions.forEach((s, sub) -> futures.add(sub.disconnect())); FutureUtil.waitForAll(futures).thenRun(() -> { closeClientFuture.complete(null); @@ -946,7 +937,7 @@ public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect List> futures = Lists.newArrayList(); replicators.forEach((cluster, replicator) -> futures.add(replicator.disconnect())); - producers.values().forEach(producer -> futures.add(producer.disconnect())); + getProducers().forEach(producer -> futures.add(producer.disconnect())); subscriptions.forEach((s, sub) -> futures.add(sub.disconnect())); CompletableFuture clientCloseFuture = closeWithoutWaitingClientDisconnect ? CompletableFuture.completedFuture(null) @@ -1285,13 +1276,13 @@ public void updateRates(NamespaceStats nsStats, NamespaceBundleStats bundleStats replicators.forEach((region, replicator) -> replicator.updateRates()); - nsStats.producerCount += producers.size(); - bundleStats.producerCount += producers.size(); + nsStats.producerCount += getProducers().count(); + bundleStats.producerCount += getProducers().count(); topicStatsStream.startObject(topic); // start publisher stats topicStatsStream.startList("publishers"); - producers.values().forEach(producer -> { + getProducers().forEach(producer -> { producer.updateRates(); PublisherStats publisherStats = producer.getStats(); @@ -1449,7 +1440,7 @@ public void updateRates(NamespaceStats nsStats, NamespaceBundleStats bundleStats // Remaining dest stats. topicStatsHelper.averageMsgSize = topicStatsHelper.aggMsgRateIn == 0.0 ? 0.0 : (topicStatsHelper.aggMsgThroughputIn / topicStatsHelper.aggMsgRateIn); - topicStatsStream.writePair("producerCount", producers.size()); + topicStatsStream.writePair("producerCount", getProducers().count()); topicStatsStream.writePair("averageMsgSize", topicStatsHelper.averageMsgSize); topicStatsStream.writePair("msgRateIn", topicStatsHelper.aggMsgRateIn); topicStatsStream.writePair("msgRateOut", topicStatsHelper.aggMsgRateOut); @@ -1494,7 +1485,7 @@ public TopicStats getStats() { ObjectObjectHashMap remotePublishersStats = new ObjectObjectHashMap(); - producers.values().forEach(producer -> { + getProducers().forEach(producer -> { PublisherStats publisherStats = producer.getStats(); stats.msgRateIn += publisherStats.msgRateIn; stats.msgThroughputIn += publisherStats.msgThroughputIn; @@ -1751,7 +1742,7 @@ public CompletableFuture onPoliciesUpdate(Policies data) { this.updateMaxPublishRate(data); - producers.values().forEach(producer -> { + getProducers().forEach(producer -> { producer.checkPermissions(); producer.checkEncryption(); }); @@ -1827,7 +1818,7 @@ public CompletableFuture terminate() { ledger.asyncTerminate(new TerminateCallback() { @Override public void terminateComplete(Position lastCommittedPosition, Object ctx) { - producers.values().forEach(Producer::disconnect); + getProducers().forEach(Producer::disconnect); subscriptions.forEach((name, sub) -> sub.topicTerminated()); PositionImpl lastPosition = (PositionImpl) lastCommittedPosition; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/NamespaceStatsAggregator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/NamespaceStatsAggregator.java index 19aa043ae823a..9ac8ae21a58c9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/NamespaceStatsAggregator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/NamespaceStatsAggregator.java @@ -108,7 +108,7 @@ private static void getTopicStats(Topic topic, TopicStats stats, boolean include stats.bytesInCounter = topic.getStats().bytesInCounter; stats.producersCount = 0; - topic.getProducers().values().forEach(producer -> { + topic.getProducers().forEach(producer -> { if (producer.isRemote()) { AggregatedReplicationStats replStats = stats.replicationStats .computeIfAbsent(producer.getRemoteCluster(), k -> new AggregatedReplicationStats()); 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..70fc5b086f8d6 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 @@ -118,7 +118,7 @@ public void testSimpleBatchProducerWithFixedBatchSize(CompressionType compressio PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); // we expect 2 messages in the backlog since we sent 50 messages with the batch size set to 25. We have set the // batch time high enough for it to not affect the number of messages in the batch assertEquals(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog(), 2); @@ -167,7 +167,7 @@ public void testSimpleBatchProducerWithFixedBatchBytes(CompressionType compressi PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); // we expect 2 messages in the backlog since we sent 50 messages with the batch size set to 25. We have set the // batch time high enough for it to not affect the number of messages in the batch assertEquals(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog(), 2); @@ -213,7 +213,7 @@ public void testSimpleBatchProducerWithFixedBatchTime(CompressionType compressio PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); LOG.info("Sent {} messages, backlog is {} messages", numMsgs, topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog()); assertTrue(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog() < numMsgs); @@ -249,7 +249,7 @@ public void testSimpleBatchProducerWithFixedBatchSizeAndTime(CompressionType com PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); LOG.info("Sent {} messages, backlog is {} messages", numMsgs, topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog()); assertTrue(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog() < numMsgs); @@ -295,7 +295,7 @@ public void testBatchProducerWithLargeMessage(CompressionType compressionType, B PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); // we expect 3 messages in the backlog since the large message in the middle should // close out the batch and be sent in a batch of its own assertEquals(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog(), 3); @@ -353,7 +353,7 @@ public void testSimpleBatchProducerConsumer(CompressionType compressionType, Bat PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); assertEquals(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog(), numMsgs / numMsgsInBatch); consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName).subscribe(); @@ -400,7 +400,7 @@ public void testSimpleBatchSyncProducerWithFixedBatchSize(BatcherBuilder builder PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); // we expect 10 messages in the backlog since we sent 10 messages with the batch size set to 5. // However, we are using synchronous send and so each message will go as an individual message assertEquals(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog(), 10); @@ -579,7 +579,7 @@ public void testNonBatchCumulativeAckAfterBatchPublish(BatcherBuilder builder) t PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); assertEquals(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog(), 2); consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName).subscribe(); @@ -636,7 +636,7 @@ public void testBatchAndNonBatchCumulativeAcks(BatcherBuilder builder) throws Ex PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); rolloverPerIntervalStats(); - assertTrue(topic.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topic.getProducers().findFirst().get().getStats().msgRateIn > 0.0); assertEquals(topic.getSubscription(subscriptionName).getNumberOfEntriesInBacklog(), (numMsgs / 2) / numMsgsInBatch + numMsgs / 2); consumer = pulsarClient.newConsumer() diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java index 505038bede387..e4e290e6be7c9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java @@ -31,13 +31,7 @@ import java.io.IOException; import java.lang.reflect.Field; -import java.net.URI; -import java.util.HashMap; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.Optional; -import java.util.Set; +import java.util.*; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; @@ -50,7 +44,6 @@ import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerException; -import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.pulsar.broker.service.BrokerServiceException.PersistenceException; @@ -69,7 +62,6 @@ import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.LocalPolicies; import org.apache.pulsar.common.policies.data.TopicStats; -import org.apache.pulsar.common.util.collections.ConcurrentOpenHashSet; import org.apache.pulsar.common.policies.data.SubscriptionStats; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; @@ -946,9 +938,9 @@ public void testStuckTopicUnloading() throws Exception { .get(mlFactory); assertNotNull(ledgers.get(topicMlName)); - org.apache.pulsar.broker.service.Producer prod = (org.apache.pulsar.broker.service.Producer) spy(topic.producers.values().toArray()[0]); - topic.producers.clear(); - topic.producers.put(prod.getProducerName(), prod); + org.apache.pulsar.broker.service.Producer prod = spy(topic.getProducers().findFirst().get()); + topic.producerGroups.clear(); + topic.producerGroups.put(prod.getProducerName(), Collections.singleton(prod)); CompletableFuture waitFuture = new CompletableFuture(); doReturn(waitFuture).when(prod).disconnect(); Set bundles = pulsar.getNamespaceService().getOwnedServiceUnits(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentQueueE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentQueueE2ETest.java index 7ffb00565c63c..2fe99a537b5ce 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentQueueE2ETest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentQueueE2ETest.java @@ -318,7 +318,7 @@ public void testSharedSingleAckedNormalTopic() throws Exception { Producer producer = pulsarClient.newProducer().topic(topicName).create(); PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. Create consumer ConsumerBuilder consumerBuilder1 = pulsarClient.newConsumer().topic(topicName) @@ -389,7 +389,7 @@ public void testCancelReadRequestOnLastDisconnect() throws Exception { Producer producer = pulsarClient.newProducer().topic(topicName).create(); PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. Create consumer ConsumerBuilder consumerBuilder = pulsarClient.newConsumer().topic(topicName) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java index 38a080f0035e0..f9ec28f43b937 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java @@ -112,7 +112,7 @@ public void testSimpleProducerEvents() throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. producer publish messages for (int i = 0; i < 10; i++) { @@ -121,13 +121,13 @@ public void testSimpleProducerEvents() throws Exception { } rolloverPerIntervalStats(); - assertTrue(topicRef.getProducers().values().iterator().next().getStats().msgRateIn > 0.0); + assertTrue(topicRef.getProducers().findFirst().get().getStats().msgRateIn > 0.0); // 3. producer disconnect producer.close(); Thread.sleep(ASYNC_EVENT_COMPLETION_WAIT); - assertEquals(topicRef.getProducers().size(), 0); + assertEquals(topicRef.getProducers().count(), 0); } @Test @@ -404,7 +404,7 @@ public void testGracefulClose() throws Exception { // 1. verify there are no pending publish acks once the producer close // is completed on client - assertEquals(topicRef.getProducers().values().iterator().next().getPendingPublishAcks(), 0); + assertEquals(topicRef.getProducers().findFirst().get().getPendingPublishAcks(), 0); // safety latch in case of failure, // wait for the spawned thread to complete @@ -1034,7 +1034,7 @@ public void testProducerReturnedMessageId() throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) topicRef.getManagedLedger(); long ledgerId = managedLedger.getLedgersInfoAsList().get(0).getLedgerId(); @@ -1232,7 +1232,7 @@ public void testCompression(CompressionType compressionType) throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. producer publish messages for (int i = 0; i < 10; i++) { @@ -1500,6 +1500,7 @@ public void testCreateProducerWithSameName() throws Exception { fail("Should have thrown ProducerBusyException"); } catch (ProducerBusyException e) { // Expected + System.out.println("busy " + e); } p1.close(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java index bc92ee492c230..b4e82b25dbd95 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java @@ -93,6 +93,7 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.impl.PulsarClientImpl; import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; +import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.api.proto.PulsarApi.CommandAck.AckType; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosition; @@ -356,9 +357,9 @@ public void testAddRemoveProducer() throws Exception { String role = "appid1"; // 1. simple add producer Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, false); + role, false, null, SchemaVersion.Latest, 0, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer); - assertEquals(topic.getProducers().size(), 1); + assertEquals(topic.getProducers().count(), 1); // 2. duplicate add try { @@ -367,12 +368,12 @@ public void testAddRemoveProducer() throws Exception { } catch (BrokerServiceException e) { assertTrue(e instanceof BrokerServiceException.NamingException); } - assertEquals(topic.getProducers().size(), 1); + assertEquals(topic.getProducers().count(), 1); // 3. add producer for a different topic PersistentTopic failTopic = new PersistentTopic(failTopicName, ledgerMock, brokerService); Producer failProducer = new Producer(failTopic, serverCnx, 2 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest,0, false); + role, false, null, SchemaVersion.Latest,0, false, PulsarApi.CommandProducer.GroupMode.Exclusive); try { topic.addProducer(failProducer); fail("should have failed"); @@ -382,7 +383,7 @@ public void testAddRemoveProducer() throws Exception { // 4. simple remove producer topic.removeProducer(producer); - assertEquals(topic.getProducers().size(), 0); + assertEquals(topic.getProducers().count(), 0); // 5. duplicate remove topic.removeProducer(producer); /* noop */ @@ -393,9 +394,9 @@ public void testProducerOverwrite() throws Exception { PersistentTopic topic = new PersistentTopic(successTopicName, ledgerMock, brokerService); String role = "appid1"; Producer producer1 = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, true); + role, false, null, SchemaVersion.Latest, 0, true, PulsarApi.CommandProducer.GroupMode.Exclusive); Producer producer2 = new Producer(topic, serverCnx, 2 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, true); + role, false, null, SchemaVersion.Latest, 0, true, PulsarApi.CommandProducer.GroupMode.Exclusive); try { topic.addProducer(producer1); topic.addProducer(producer2); @@ -404,10 +405,10 @@ public void testProducerOverwrite() throws Exception { // OK } - Assert.assertEquals(topic.getProducers().size(), 1); + Assert.assertEquals(topic.getProducers().count(), 1); Producer producer3 = new Producer(topic, serverCnx, 2 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 1, false); + role, false, null, SchemaVersion.Latest, 1, false, PulsarApi.CommandProducer.GroupMode.Exclusive); try { topic.addProducer(producer3); @@ -416,44 +417,44 @@ public void testProducerOverwrite() throws Exception { // OK } - Assert.assertEquals(topic.getProducers().size(), 1); + Assert.assertEquals(topic.getProducers().count(), 1); topic.removeProducer(producer1); - Assert.assertEquals(topic.getProducers().size(), 0); + Assert.assertEquals(topic.getProducers().count(), 0); Producer producer4 = new Producer(topic, serverCnx, 2 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 2, false); + role, false, null, SchemaVersion.Latest, 2, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer3); topic.addProducer(producer4); - Assert.assertEquals(topic.getProducers().size(), 1); + Assert.assertEquals(topic.getProducers().count(), 1); - topic.getProducers().values().forEach(producer -> Assert.assertEquals(producer.getEpoch(), 2)); + topic.getProducers().forEach(producer -> Assert.assertEquals(producer.getEpoch(), 2)); topic.removeProducer(producer4); - Assert.assertEquals(topic.getProducers().size(), 0); + Assert.assertEquals(topic.getProducers().count(), 0); Producer producer5 = new Producer(topic, serverCnx, 2 /* producer id */, "pulsar.repl.cluster1", - role, false, null, SchemaVersion.Latest, 1, false); + role, false, null, SchemaVersion.Latest, 1, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer5); - Assert.assertEquals(topic.getProducers().size(), 1); + Assert.assertEquals(topic.getProducers().count(), 1); Producer producer6 = new Producer(topic, serverCnx, 2 /* producer id */, "pulsar.repl.cluster1", - role, false, null, SchemaVersion.Latest, 2, false); + role, false, null, SchemaVersion.Latest, 2, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer6); - Assert.assertEquals(topic.getProducers().size(), 1); + Assert.assertEquals(topic.getProducers().count(), 1); - topic.getProducers().values().forEach(producer -> Assert.assertEquals(producer.getEpoch(), 2)); + topic.getProducers().forEach(producer -> Assert.assertEquals(producer.getEpoch(), 2)); Producer producer7 = new Producer(topic, serverCnx, 2 /* producer id */, "pulsar.repl.cluster1", - role, false, null, SchemaVersion.Latest, 3, true); + role, false, null, SchemaVersion.Latest, 3, true, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer7); - Assert.assertEquals(topic.getProducers().size(), 1); - topic.getProducers().values().forEach(producer -> Assert.assertEquals(producer.getEpoch(), 3)); + Assert.assertEquals(topic.getProducers().count(), 1); + topic.getProducers().forEach(producer -> Assert.assertEquals(producer.getEpoch(), 3)); } public void testMaxProducers() throws Exception { @@ -461,20 +462,20 @@ public void testMaxProducers() throws Exception { String role = "appid1"; // 1. add producer1 Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name1", role, - false, null, SchemaVersion.Latest,0, false); + false, null, SchemaVersion.Latest,0, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer); - assertEquals(topic.getProducers().size(), 1); + assertEquals(topic.getProducers().count(), 1); // 2. add producer2 Producer producer2 = new Producer(topic, serverCnx, 2 /* producer id */, "prod-name2", role, - false, null, SchemaVersion.Latest,0, false); + false, null, SchemaVersion.Latest,0, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer2); - assertEquals(topic.getProducers().size(), 2); + assertEquals(topic.getProducers().count(), 2); // 3. add producer3 but reached maxProducersPerTopic try { Producer producer3 = new Producer(topic, serverCnx, 3 /* producer id */, "prod-name3", role, - false, null, SchemaVersion.Latest,0, false); + false, null, SchemaVersion.Latest,0, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer3); fail("should have failed"); } catch (BrokerServiceException e) { @@ -871,7 +872,7 @@ public void testDeleteTopic() throws Exception { // 2. delete topic with producer topic = (PersistentTopic) brokerService.getOrCreateTopic(successTopicName).get(); Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, false); + role, false, null, SchemaVersion.Latest, 0, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer); assertTrue(topic.delete().isCompletedExceptionally()); @@ -1030,7 +1031,7 @@ public Object answer(InvocationOnMock invocationOnMock) throws Throwable { String role = "appid1"; Thread.sleep(10); /* delay to ensure that the delete gets executed first */ Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, false); + role, false, null, SchemaVersion.Latest, 0, false, PulsarApi.CommandProducer.GroupMode.Exclusive); topic.addProducer(producer); fail("Should have failed"); } catch (BrokerServiceException e) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ResendRequestTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ResendRequestTest.java index e6bdd71ad4de5..b2f353ef4c07a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ResendRequestTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ResendRequestTest.java @@ -77,7 +77,7 @@ public void testExclusiveSingleAckedNormalTopic() throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. Create consumer Consumer consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName) @@ -165,7 +165,7 @@ public void testSharedSingleAckedNormalTopic() throws Exception { .create(); PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. Create consumer Consumer consumer1 = pulsarClient.newConsumer().topic(topicName) @@ -252,7 +252,7 @@ public void testFailoverSingleAckedNormalTopic() throws Exception { .create(); PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. Create consumer ConsumerBuilder consumerBuilder = pulsarClient.newConsumer().topic(topicName) @@ -372,7 +372,7 @@ public void testExclusiveCumulativeAckedNormalTopic() throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. Create consumer Consumer consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName) @@ -672,7 +672,7 @@ public void testFailoverInactiveConsumer() throws Exception { .create(); PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // 2. Create consumer ConsumerBuilder consumerBuilder = pulsarClient.newConsumer().topic(topicName) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java index 1a4231e83e98b..3d10fc93fef4e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java @@ -412,7 +412,7 @@ public void testProducerCommand() throws Exception { PersistentTopic topicRef = (PersistentTopic) brokerService.getTopicReference(successTopicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); // test PRODUCER error case clientCommand = Commands.newProducer(failTopicName, 2, 2, @@ -423,7 +423,7 @@ public void testProducerCommand() throws Exception { assertFalse(pulsar.getBrokerService().getTopicReference(failTopicName).isPresent()); channel.finish(); - assertEquals(topicRef.getProducers().size(), 0); + assertEquals(topicRef.getProducers().count(), 0); } @Test(timeOut = 5000) @@ -492,10 +492,10 @@ public void testProducerCommandWithAuthorizationPositive() throws Exception { PersistentTopic topicRef = (PersistentTopic) brokerService.getTopicReference(successTopicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); channel.finish(); - assertEquals(topicRef.getProducers().size(), 0); + assertEquals(topicRef.getProducers().count(), 0); } @Test(timeOut = 30000) @@ -584,7 +584,7 @@ public void testNonExistentTopicSuperUserAccess() throws Exception { PersistentTopic topicRef = (PersistentTopic) brokerService.getTopicReference(nonExistentTopicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); channel.finish(); // Test consumer creation @@ -1299,7 +1299,7 @@ public void testProducerSuccessOnEncryptionRequiredTopic() throws Exception { assertEquals(response.getClass(), CommandProducerSuccess.class); PersistentTopic topicRef = (PersistentTopic) brokerService.getTopicReference(encryptionRequiredTopicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); channel.finish(); } @@ -1329,7 +1329,7 @@ public void testProducerFailureOnEncryptionRequiredTopic() throws Exception { assertEquals(errorResponse.getError(), ServerError.MetadataError); PersistentTopic topicRef = (PersistentTopic) brokerService.getTopicReference(encryptionRequiredTopicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 0); + assertEquals(topicRef.getProducers().count(), 0); channel.finish(); } 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 bbcb6b194b19b..cc6285f0e9df4 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 @@ -65,7 +65,7 @@ public void testSeek() throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); assertEquals(topicRef.getSubscriptions().size(), 1); List messageIds = new ArrayList<>(); @@ -122,7 +122,7 @@ public void testSeekTime() throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); assertEquals(topicRef.getSubscriptions().size(), 1); PersistentSubscription sub = topicRef.getSubscription("my-subscription"); @@ -162,7 +162,7 @@ public void testSeekTimeOnPartitionedTopic() throws Exception { PersistentTopic topicRef = (PersistentTopic) pulsar.getBrokerService() .getTopicReference(topicName + TopicName.PARTITIONED_TOPIC_SUFFIX + i).get(); assertNotNull(topicRef); - assertEquals(topicRef.getProducers().size(), 1); + assertEquals(topicRef.getProducers().count(), 1); assertEquals(topicRef.getSubscriptions().size(), 1); PersistentSubscription sub = topicRef.getSubscription("my-subscription"); assertNotNull(sub); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java index b13b61d4d74b2..f50c9da0ee3b5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java @@ -92,4 +92,93 @@ public void testGeneratedNameProducerReconnect(TopicDomain domain) throws Pulsar Assert.assertEquals(producer.getConnectionHandler().getEpoch(), 1); Assert.assertTrue(producer.isConnected()); } + + @Test(dataProvider = "topicDomainProvider") + public void testParallelProducersGeneratedNameFails(TopicDomain domain) { + try { + pulsarClient.newProducer() + .topic("testParallelProducersGeneratedNameFails") + .groupMode(ProducerGroupMode.Parallel) + .create(); + Assert.fail("previous statement should have failed"); + } catch (PulsarClientException e) { + // ok here + String msg = e.getMessage(); + Assert.assertTrue(msg.endsWith("producerName must be specified in non-exclusive group modes"), msg); + } + } + + @Test(dataProvider = "topicDomainProvider") + public void testParallelProducers(TopicDomain domain) throws PulsarClientException { + Producer producer1 = pulsarClient.newProducer() + .topic(TopicName.get(domain.value(), "public", "default", "testParallelProducers").toString()) + .producerName("p-name-1") + .groupMode(ProducerGroupMode.Parallel) + .create(); + + Assert.assertNotNull(producer1); + + Producer producer2 = pulsarClient.newProducer() + .topic("testParallelProducers") + .producerName("p-name-1") + .groupMode(ProducerGroupMode.Parallel) + .create(); + + Assert.assertNotNull(producer2); + } + + @Test(dataProvider = "topicDomainProvider") + public void testParallelProducerCannotJoinExclusiveGroup(TopicDomain domain) throws PulsarClientException { + Producer producer1 = pulsarClient.newProducer() + .topic(TopicName.get(domain.value(), "public", "default", "testParallelProducerCannotJoinExclusiveGroup").toString()) + .producerName("p-name-1") + .create(); + + Assert.assertNotNull(producer1); + + try { + pulsarClient.newProducer() + .topic(TopicName.get(domain.value(), "public", "default", "testParallelProducerCannotJoinExclusiveGroup").toString()) + .producerName("p-name-1") + .groupMode(ProducerGroupMode.Parallel) + .create(); + Assert.fail("previous statement should have failed"); + } catch (PulsarClientException e) { + // ok here + String msg = e.getMessage(); + Assert.assertTrue(msg.endsWith("Exclusive Producer with name 'p-name-1' is already connected to topic"), msg); + } + } + + @Test(dataProvider = "topicDomainProvider") + public void testExclusiveProducerEvictsParallelGroup(TopicDomain domain) throws PulsarClientException { + Producer producer1 = pulsarClient.newProducer() + .topic(TopicName.get(domain.value(), "public", "default", "testParallelProducerCannotJoinExclusiveGroup").toString()) + .producerName("p-name-1") + .groupMode(ProducerGroupMode.Parallel) + .create(); + + Assert.assertNotNull(producer1); + + Producer producer2 = pulsarClient.newProducer() + .topic(TopicName.get(domain.value(), "public", "default", "testParallelProducerCannotJoinExclusiveGroup").toString()) + .producerName("p-name-1") + .groupMode(ProducerGroupMode.Parallel) + .create(); + + Assert.assertNotNull(producer2); + + Assert.assertTrue(producer1.isConnected()); + Assert.assertTrue(producer2.isConnected()); + + Producer producer3 = pulsarClient.newProducer() + .topic(TopicName.get(domain.value(), "public", "default", "testParallelProducerCannotJoinExclusiveGroup").toString()) + .producerName("p-name-1") + .create(); + + Assert.assertNotNull(producer3); + + producer3.send("my-message".getBytes()); + } + } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessagePublishThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessagePublishThrottlingTest.java index 8653cf13446ee..348cf6cfe90bc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessagePublishThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessagePublishThrottlingTest.java @@ -93,7 +93,7 @@ public void testSimplePublishMessageThrottling() throws Exception { 200); Assert.assertNotEquals(topic.getTopicPublishRateLimiter(), PublishRateLimiter.DISABLED_RATE_LIMITER); - Producer prod = topic.getProducers().values().iterator().next(); + Producer prod = topic.getProducers().findFirst().get(); // reset counter prod.updateRates(); int total = 200; @@ -157,7 +157,7 @@ public void testSimplePublishByteThrottling() throws Exception { 200); Assert.assertNotEquals(topic.getTopicPublishRateLimiter(), PublishRateLimiter.DISABLED_RATE_LIMITER); - Producer prod = topic.getProducers().values().iterator().next(); + Producer prod = topic.getProducers().findFirst().get(); // reset counter prod.updateRates(); int total = 100; @@ -232,7 +232,7 @@ public void testBrokerPublishMessageThrottling() throws Exception { Assert.assertNotEquals(topic.getBrokerPublishRateLimiter(), PublishRateLimiter.DISABLED_RATE_LIMITER); - Producer prod = topic.getProducers().values().iterator().next(); + Producer prod = topic.getProducers().findFirst().get(); // reset counter prod.updateRates(); int total = 100; @@ -309,7 +309,7 @@ public void testBrokerPublishByteThrottling() throws Exception { Assert.assertNotEquals(topic.getBrokerPublishRateLimiter(), PublishRateLimiter.DISABLED_RATE_LIMITER); - Producer prod = topic.getProducers().values().iterator().next(); + Producer prod = topic.getProducers().findFirst().get(); // reset counter prod.updateRates(); int numMessage = 20; @@ -406,7 +406,7 @@ public void testBrokerTopicPublishByteThrottling() throws Exception { Assert.assertNotEquals(topic.getBrokerPublishRateLimiter(), PublishRateLimiter.DISABLED_RATE_LIMITER); Assert.assertNotEquals(topic.getTopicPublishRateLimiter(), PublishRateLimiter.DISABLED_RATE_LIMITER); - Producer prod = topic.getProducers().values().iterator().next(); + Producer prod = topic.getProducers().findFirst().get(); // reset counter prod.updateRates(); int numMessage = 40; @@ -462,7 +462,7 @@ public void testBrokerTopicPublishByteThrottling() throws Exception { int id = index.incrementAndGet(); ProducerImpl iProducer = producers.get(id); PersistentTopic iTopic = topics.get(id); - Producer iProd = iTopic.getProducers().values().iterator().next(); + Producer iProd = iTopic.getProducers().findFirst().get(); // reset counter iProd.updateRates(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TopicPublishThrottlingInitTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TopicPublishThrottlingInitTest.java index efcbc51a4da24..48cd6678ee7bf 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TopicPublishThrottlingInitTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TopicPublishThrottlingInitTest.java @@ -84,7 +84,7 @@ public void testBrokerPublishMessageThrottlingInit() throws Exception { Assert.assertNotEquals(topic.getBrokerPublishRateLimiter(), PublishRateLimiter.DISABLED_RATE_LIMITER); - Producer prod = topic.getProducers().values().iterator().next(); + Producer prod = topic.getProducers().findFirst().get(); // reset counter prod.updateRates(); int total = 100; diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java index c4a250e901490..b97c482c003ea 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java @@ -490,4 +490,14 @@ public interface ProducerBuilder extends Cloneable { * @since 2.5.0 */ ProducerBuilder enableMultiSchema(boolean multiSchema); + + /** + * Controls the behavior when multiple producers with the same producerName connect to the same topic. + * default: Exclusive + * + * @see ProducerGroupMode + * @param groupMode the producer group mode for this producer + * @return the producer builder instance + */ + ProducerBuilder groupMode(ProducerGroupMode groupMode); } diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerGroupMode.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerGroupMode.java new file mode 100644 index 0000000000000..a7b063188d1c2 --- /dev/null +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerGroupMode.java @@ -0,0 +1,44 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.api; + +/** + * Controls the behavior when multiple producers with the same producerName connect to a topic. + */ +public enum ProducerGroupMode { + /** + * Only one producer can be active at any one point in time. + * + *

Producers trying to connect with the same producerName and topic will be rejected. + */ + Exclusive, + + /** + * Multiple producers can be active and producing in parallel. + * + *

This can be used in active/active producer scenarios, when deduplication needs to work across multiple + * producers. + */ + Parallel, +// +// /** +// * Concurrently connecting producer will be blocked until the active producer fails. +// */ +// Failover, +} \ No newline at end of file diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBase.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBase.java index 50f3fa15139a1..2187c8fb1f56f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBase.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBase.java @@ -22,16 +22,11 @@ import java.util.concurrent.CompletableFuture; -import org.apache.pulsar.client.api.Message; -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.Schema; -import org.apache.pulsar.client.api.SchemaSerializationException; -import org.apache.pulsar.client.api.TypedMessageBuilder; +import org.apache.pulsar.client.api.*; import org.apache.pulsar.client.api.transaction.Transaction; import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; import org.apache.pulsar.client.impl.transaction.TransactionImpl; +import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.protocol.schema.SchemaHash; import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; @@ -129,6 +124,20 @@ public void flush() throws PulsarClientException { abstract void triggerFlush(); + protected PulsarApi.CommandProducer.GroupMode getGroupMode() { + ProducerGroupMode type = conf.getGroupMode(); + switch (type) { + case Exclusive: + return PulsarApi.CommandProducer.GroupMode.Exclusive; + + case Parallel: + return PulsarApi.CommandProducer.GroupMode.Parallel; + } + + // Should not happen since we cover all cases above + return null; + } + @Override public void close() throws PulsarClientException { try { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java index 12a52e3ad5b7c..3e61cfcb6da57 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java @@ -39,6 +39,7 @@ import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerBuilder; +import org.apache.pulsar.client.api.ProducerGroupMode; import org.apache.pulsar.client.api.ProducerCryptoFailureAction; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; @@ -291,6 +292,12 @@ public ProducerBuilder enableMultiSchema(boolean multiSchema) { return this; } + @Override + public ProducerBuilder groupMode(ProducerGroupMode groupMode) { + conf.setGroupMode(groupMode); + return this; + } + private void setMessageRoutingMode() throws PulsarClientException { if(conf.getMessageRoutingMode() == null && conf.getCustomMessageRouter() == null) { messageRoutingMode(MessageRoutingMode.RoundRobinPartition); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 7805fa8ddfb17..b00b5ebed9a4a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -66,6 +66,7 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; import org.apache.pulsar.client.impl.schema.JSONSchema; +import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; import org.apache.pulsar.common.api.proto.PulsarApi.ProtocolVersion; import org.apache.pulsar.common.compression.CompressionCodec; @@ -103,6 +104,7 @@ public class ProducerImpl extends ProducerBase implements TimerTask, Conne // Globally unique producer name private String producerName; private boolean userProvidedProducerName = false; + private PulsarApi.CommandProducer.GroupMode groupMode; private String connectionId; private String connectedSince; @@ -138,6 +140,7 @@ public ProducerImpl(PulsarClientImpl client, String topic, ProducerConfiguration if (StringUtils.isNotBlank(producerName)) { this.userProvidedProducerName = true; } + this.groupMode = getGroupMode(); this.partitionIndex = partitionIndex; this.pendingMessages = Queues.newArrayBlockingQueue(conf.getMaxPendingMessages()); this.pendingCallbacks = Queues.newArrayBlockingQueue(conf.getMaxPendingMessages()); @@ -1087,7 +1090,7 @@ public void connectionOpened(final ClientCnx cnx) { cnx.sendRequestWithId( Commands.newProducer(topic, producerId, requestId, producerName, conf.isEncryptionEnabled(), metadata, - schemaInfo, connectionHandler.epoch, userProvidedProducerName), + schemaInfo, connectionHandler.epoch, userProvidedProducerName, groupMode), requestId).thenAccept(response -> { String producerName = response.getProducerName(); long lastSequenceId = response.getLastSequenceId(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java index 654c15fdb0741..fe5400086e67d 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java @@ -35,6 +35,7 @@ import org.apache.pulsar.client.api.MessageRouter; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.ProducerCryptoFailureAction; +import org.apache.pulsar.client.api.ProducerGroupMode; import com.fasterxml.jackson.annotation.JsonIgnore; import com.google.common.collect.Maps; @@ -92,6 +93,8 @@ public class ProducerConfigurationData implements Serializable, Cloneable { private boolean multiSchema = true; + private ProducerGroupMode groupMode = ProducerGroupMode.Exclusive; + private SortedMap properties = new TreeMap<>(); /** diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java index 1e9ffd608fa6d..38a091d9ea3ad 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java @@ -14507,6 +14507,10 @@ public interface CommandProducerOrBuilder // optional bool user_provided_producer_name = 9 [default = true]; boolean hasUserProvidedProducerName(); boolean getUserProvidedProducerName(); + + // optional .pulsar.proto.CommandProducer.GroupMode group_mode = 10 [default = Exclusive]; + boolean hasGroupMode(); + org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode getGroupMode(); } public static final class CommandProducer extends org.apache.pulsar.shaded.com.google.protobuf.v241.GeneratedMessageLite @@ -14542,6 +14546,47 @@ public CommandProducer getDefaultInstanceForType() { return defaultInstance; } + public enum GroupMode + implements org.apache.pulsar.shaded.com.google.protobuf.v241.Internal.EnumLite { + Exclusive(0, 0), + Parallel(1, 1), + ; + + public static final int Exclusive_VALUE = 0; + public static final int Parallel_VALUE = 1; + + + public final int getNumber() { return value; } + + public static GroupMode valueOf(int value) { + switch (value) { + case 0: return Exclusive; + case 1: return Parallel; + default: return null; + } + } + + public static org.apache.pulsar.shaded.com.google.protobuf.v241.Internal.EnumLiteMap + internalGetValueMap() { + return internalValueMap; + } + private static org.apache.pulsar.shaded.com.google.protobuf.v241.Internal.EnumLiteMap + internalValueMap = + new org.apache.pulsar.shaded.com.google.protobuf.v241.Internal.EnumLiteMap() { + public GroupMode findValueByNumber(int number) { + return GroupMode.valueOf(number); + } + }; + + private final int value; + + private GroupMode(int index, int value) { + this.value = value; + } + + // @@protoc_insertion_point(enum_scope:pulsar.proto.CommandProducer.GroupMode) + } + private int bitField0_; // required string topic = 1; public static final int TOPIC_FIELD_NUMBER = 1; @@ -14688,6 +14733,16 @@ public boolean getUserProvidedProducerName() { return userProvidedProducerName_; } + // optional .pulsar.proto.CommandProducer.GroupMode group_mode = 10 [default = Exclusive]; + public static final int GROUP_MODE_FIELD_NUMBER = 10; + private org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode groupMode_; + public boolean hasGroupMode() { + return ((bitField0_ & 0x00000100) == 0x00000100); + } + public org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode getGroupMode() { + return groupMode_; + } + private void initFields() { topic_ = ""; producerId_ = 0L; @@ -14698,6 +14753,7 @@ private void initFields() { schema_ = org.apache.pulsar.common.api.proto.PulsarApi.Schema.getDefaultInstance(); epoch_ = 0L; userProvidedProducerName_ = true; + groupMode_ = org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode.Exclusive; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -14767,6 +14823,9 @@ public void writeTo(org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStr if (((bitField0_ & 0x00000080) == 0x00000080)) { output.writeBool(9, userProvidedProducerName_); } + if (((bitField0_ & 0x00000100) == 0x00000100)) { + output.writeEnum(10, groupMode_.getNumber()); + } } private int memoizedSerializedSize = -1; @@ -14811,6 +14870,10 @@ public int getSerializedSize() { size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream .computeBoolSize(9, userProvidedProducerName_); } + if (((bitField0_ & 0x00000100) == 0x00000100)) { + size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream + .computeEnumSize(10, groupMode_.getNumber()); + } memoizedSerializedSize = size; return size; } @@ -14942,6 +15005,8 @@ public Builder clear() { bitField0_ = (bitField0_ & ~0x00000080); userProvidedProducerName_ = true; bitField0_ = (bitField0_ & ~0x00000100); + groupMode_ = org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode.Exclusive; + bitField0_ = (bitField0_ & ~0x00000200); return this; } @@ -15012,6 +15077,10 @@ public org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer buildPartial to_bitField0_ |= 0x00000080; } result.userProvidedProducerName_ = userProvidedProducerName_; + if (((from_bitField0_ & 0x00000200) == 0x00000200)) { + to_bitField0_ |= 0x00000100; + } + result.groupMode_ = groupMode_; result.bitField0_ = to_bitField0_; return result; } @@ -15052,6 +15121,9 @@ public Builder mergeFrom(org.apache.pulsar.common.api.proto.PulsarApi.CommandPro if (other.hasUserProvidedProducerName()) { setUserProvidedProducerName(other.getUserProvidedProducerName()); } + if (other.hasGroupMode()) { + setGroupMode(other.getGroupMode()); + } return this; } @@ -15156,6 +15228,15 @@ public Builder mergeFrom( userProvidedProducerName_ = input.readBool(); break; } + case 80: { + int rawValue = input.readEnum(); + org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode value = org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode.valueOf(rawValue); + if (value != null) { + bitField0_ |= 0x00000200; + groupMode_ = value; + } + break; + } } } } @@ -15471,6 +15552,30 @@ public Builder clearUserProvidedProducerName() { return this; } + // optional .pulsar.proto.CommandProducer.GroupMode group_mode = 10 [default = Exclusive]; + private org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode groupMode_ = org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode.Exclusive; + public boolean hasGroupMode() { + return ((bitField0_ & 0x00000200) == 0x00000200); + } + public org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode getGroupMode() { + return groupMode_; + } + public Builder setGroupMode(org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode value) { + if (value == null) { + throw new NullPointerException(); + } + bitField0_ |= 0x00000200; + groupMode_ = value; + + return this; + } + public Builder clearGroupMode() { + bitField0_ = (bitField0_ & ~0x00000200); + groupMode_ = org.apache.pulsar.common.api.proto.PulsarApi.CommandProducer.GroupMode.Exclusive; + + return this; + } + // @@protoc_insertion_point(builder_scope:pulsar.proto.CommandProducer) } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index 5b8c91321279c..f8f1527b20f77 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -684,7 +684,8 @@ public static ByteBuf newProducer(String topic, long producerId, long requestId, public static ByteBuf newProducer(String topic, long producerId, long requestId, String producerName, boolean encrypted, Map metadata) { - return newProducer(topic, producerId, requestId, producerName, encrypted, metadata, null, 0, false); + return newProducer(topic, producerId, requestId, producerName, encrypted, metadata, null, 0, false, + CommandProducer.GroupMode.Exclusive); } private static Schema.Type getSchemaType(SchemaType type) { @@ -724,7 +725,7 @@ private static Schema getSchema(SchemaInfo schemaInfo) { public static ByteBuf newProducer(String topic, long producerId, long requestId, String producerName, boolean encrypted, Map metadata, SchemaInfo schemaInfo, - long epoch, boolean userProvidedProducerName) { + long epoch, boolean userProvidedProducerName, CommandProducer.GroupMode groupMode) { CommandProducer.Builder producerBuilder = CommandProducer.newBuilder(); producerBuilder.setTopic(topic); producerBuilder.setProducerId(producerId); @@ -735,6 +736,7 @@ public static ByteBuf newProducer(String topic, long producerId, long requestId, } producerBuilder.setUserProvidedProducerName(userProvidedProducerName); producerBuilder.setEncrypted(encrypted); + producerBuilder.setGroupMode(groupMode); producerBuilder.addAllMetadata(CommandUtils.toKeyValueList(metadata)); diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index 8a44fca60875a..cf5512cb895c1 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -395,6 +395,12 @@ message CommandLookupTopicResponse { /// Create a new Producer on a topic, assigning the given producer_id, /// all messages sent with this producer_id will be persisted on the topic message CommandProducer { + // Controls the behavior when multiple producers with the same producerName connect to a topic + enum GroupMode { + Exclusive = 0; // only one producer can be active at any one point in time + Parallel = 1; // multiple producers can be active and producing in parallel + // TODO: Failover = 2; // concurrently connecting producer will be blocked until the active producer fails + } required string topic = 1; required uint64 producer_id = 2; required uint64 request_id = 3; @@ -416,6 +422,8 @@ message CommandProducer { // Indicate the name of the producer is generated or user provided // Use default true here is in order to be forward compatible with the client optional bool user_provided_producer_name = 9 [default = true]; + + optional GroupMode group_mode = 10 [default = Exclusive]; } message CommandSend { From 6b4b0e1e6794edd04a84d157ecf2a2b3fede6f46 Mon Sep 17 00:00:00 2001 From: Eugen Dueck Date: Mon, 17 Feb 2020 11:58:03 +0900 Subject: [PATCH 4/4] fix build after merge --- .../pulsar/broker/service/AbstractTopic.java | 20 +++++++++++++++++++ .../pulsar/broker/service/BrokerService.java | 2 +- .../MessagePublishBufferThrottleTest.java | 16 +++++++-------- 3 files changed, 29 insertions(+), 9 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index 58e3df0a9768e..7919525bff9ce 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -111,6 +111,26 @@ public AbstractTopic(String topic, BrokerService brokerService) { updatePublishDispatcher(policies); } + /** + * For testing purposes. Performs linear scan of all producers in the producer group. + */ + Producer getProducer(String producerName) { + return getProducer("", producerName); + } + + /** + * For testing purposes. Performs linear scan of all producers in the producer group. + */ + Producer getProducer(String groupName, String producerName) { + Set producers = producerGroups.get(groupName); + for (Producer producer : producers) { + if (producer.getProducerName().equals(producerName)) { + return producer; + } + } + return null; + } + protected boolean isProducersExceeded() { Policies policies; try { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 7d48f76b71e90..de7151bd5f3e4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -2050,7 +2050,7 @@ private void checkMessagePublishBuffer() { private void foreachProducer(Consumer consumer) { topics.forEach((n, t) -> { Optional topic = extractTopic(t); - topic.ifPresent(value -> value.getProducers().values().forEach(consumer)); + topic.ifPresent(value -> value.getProducers().forEach(consumer)); }); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java index 5397725a99c67..6695052f46501 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessagePublishBufferThrottleTest.java @@ -56,7 +56,7 @@ public void testMessagePublishBufferThrottleDisabled() throws Exception { .create(); Topic topicRef = pulsar.getBrokerService().getTopicReference(topic).get(); Assert.assertNotNull(topicRef); - ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); + ((AbstractTopic)topicRef).getProducer("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); Thread.sleep(20); Assert.assertFalse(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); List> futures = new ArrayList<>(); @@ -87,7 +87,7 @@ public void testMessagePublishBufferThrottleEnable() throws Exception { .create(); Topic topicRef = pulsar.getBrokerService().getTopicReference(topic).get(); Assert.assertNotNull(topicRef); - ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); + ((AbstractTopic)topicRef).getProducer("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); Thread.sleep(4); Assert.assertTrue(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); // The first message can publish success, but the second message should be blocked @@ -101,7 +101,7 @@ public void testMessagePublishBufferThrottleEnable() throws Exception { } Assert.assertNull(messageId); - ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(0L); + ((AbstractTopic)topicRef).getProducer("producer-name").getCnx().setMessagePublishBufferSize(0L); Thread.sleep(4); List> futures = new ArrayList<>(); @@ -132,12 +132,12 @@ public void testBlockByPublishRateLimiting() throws Exception { .create(); Topic topicRef = pulsar.getBrokerService().getTopicReference(topic).get(); Assert.assertNotNull(topicRef); - ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); + ((AbstractTopic)topicRef).getProducer("producer-name").getCnx().setMessagePublishBufferSize(Long.MAX_VALUE / 2); producer.sendAsync(new byte[1024]).get(1, TimeUnit.SECONDS); Thread.sleep(4); - ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setAutoReadDisabledRateLimiting(true); - ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setMessagePublishBufferSize(0); + ((AbstractTopic)topicRef).getProducer("producer-name").getCnx().setAutoReadDisabledRateLimiting(true); + ((AbstractTopic)topicRef).getProducer("producer-name").getCnx().setMessagePublishBufferSize(0); Thread.sleep(4); Assert.assertFalse(pulsar.getBrokerService().isReachMessagePublishBufferThreshold()); MessageId messageId = null; @@ -149,8 +149,8 @@ public void testBlockByPublishRateLimiting() throws Exception { } Assert.assertNull(messageId); - ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().setAutoReadDisabledRateLimiting(false); - ((AbstractTopic)topicRef).producers.get("producer-name").getCnx().enableCnxAutoRead(); + ((AbstractTopic)topicRef).getProducer("producer-name").getCnx().setAutoReadDisabledRateLimiting(false); + ((AbstractTopic)topicRef).getProducer("producer-name").getCnx().enableCnxAutoRead(); List> futures = new ArrayList<>(); // Make sure the producer can publish succeed.