From 60c702d339eba393895a46ad3c81906667ff4a59 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Sat, 7 May 2022 16:27:15 +0200 Subject: [PATCH 1/6] [enh] Exclusive Producer: ability to fence out an existing Producer (ExclusiveWithFencing mode) --- .../pulsar/broker/service/AbstractTopic.java | 38 ++++++++++++++++++- .../broker/service/ExclusiveProducerTest.java | 19 +++++++++- .../pulsar/client/api/ProducerAccessMode.java | 5 +++ .../pulsar/client/api/ProducerBuilder.java | 2 + .../pulsar/common/protocol/Commands.java | 4 ++ pulsar-common/src/main/proto/PulsarApi.proto | 1 + 6 files changed, 67 insertions(+), 2 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 2c0d884100841..61a3a78a790a0 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 @@ -738,7 +738,43 @@ protected CompletableFuture> incrementTopicEpochIfNeeded(Producer return topicEpoch; }); } - + case ExclusiveWithFencing: + if (hasExclusiveProducer) { + producers.forEach((k, currentProducer) -> { + log.info("Fencing out producer {}", currentProducer); + currentProducer.close(true); + }); + } + if (producer.getTopicEpoch().isPresent() + && producer.getTopicEpoch().get() < topicEpoch.orElse(-1L)) { + // If a producer reconnects, but all the topic epoch has already moved forward, + // this producer needs to be fenced, because a new producer had been present in between. + hasExclusiveProducer = false; + return FutureUtil.failedFuture(new ProducerFencedException( + String.format("Topic epoch has already moved. Current epoch: %d, Producer epoch: %d", + topicEpoch.get(), producer.getTopicEpoch().get()))); + } else { + // There are currently no existing producers + hasExclusiveProducer = true; + exclusiveProducerName = producer.getProducerName(); + + CompletableFuture future; + if (producer.getTopicEpoch().isPresent()) { + future = setTopicEpoch(producer.getTopicEpoch().get()); + } else { + future = incrementTopicEpoch(topicEpoch); + } + future.exceptionally(ex -> { + hasExclusiveProducer = false; + exclusiveProducerName = null; + return null; + }); + + return future.thenApply(epoch -> { + topicEpoch = Optional.of(epoch); + return topicEpoch; + }); + } case WaitForExclusive: { if (hasExclusiveProducer || !producers.isEmpty()) { CompletableFuture> future = new CompletableFuture<>(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java index 4d2b16a6a0e75..2ed63a8991f63 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java @@ -73,6 +73,8 @@ public static Object[][] accessMode() { // ProducerAccessMode, partitioned { ProducerAccessMode.Exclusive, Boolean.TRUE}, { ProducerAccessMode.Exclusive, Boolean.FALSE }, + { ProducerAccessMode.ExclusiveWithFencing, Boolean.TRUE}, + { ProducerAccessMode.ExclusiveWithFencing, Boolean.FALSE }, { ProducerAccessMode.WaitForExclusive, Boolean.TRUE }, { ProducerAccessMode.WaitForExclusive, Boolean.FALSE }, }; @@ -118,7 +120,22 @@ private void simpleTest(String topic) throws Exception { .topic(topic) .accessMode(ProducerAccessMode.Exclusive) .create(); - p2.close(); + + Producer p3 = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .accessMode(ProducerAccessMode.ExclusiveWithFencing) + .create(); + + try { + p2.send("test"); + fail("Should have failed"); + } catch (ProducerFencedException expected) { + } + + // this should work + p3.send("test"); + p3.close(); + } @Test(dataProvider = "topics") diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java index 85199e7cf4c27..c7852214fdee6 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java @@ -33,6 +33,11 @@ public enum ProducerAccessMode { */ Exclusive, + /** + * Require exclusive access for producer. Fence out the old producer. + */ + ExclusiveWithFencing, + /** * Producer creation is pending until it can acquire exclusive access. */ 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 cab4d497307f2..89c330338d8d2 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 @@ -133,6 +133,8 @@ public interface ProducerBuilder extends Cloneable { *
  • {@link ProducerAccessMode#Shared}: By default multiple producers can publish on a topic *
  • {@link ProducerAccessMode#Exclusive}: Require exclusive access for producer. Fail immediately if there's * already a producer connected. + *
  • {@link ProducerAccessMode#ExclusiveWithFencing}: Require exclusive access for producer. Fence out any + * producer that is connected. *
  • {@link ProducerAccessMode#WaitForExclusive}: Producer creation is pending until it can acquire exclusive * access * 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 d1cbe3ad96141..2d8e043058dd7 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 @@ -1853,6 +1853,8 @@ private static org.apache.pulsar.common.api.proto.ProducerAccessMode convertProd return org.apache.pulsar.common.api.proto.ProducerAccessMode.Shared; case WaitForExclusive: return org.apache.pulsar.common.api.proto.ProducerAccessMode.WaitForExclusive; + case ExclusiveWithFencing: + return org.apache.pulsar.common.api.proto.ProducerAccessMode.ExclusiveWithFencing; default: throw new IllegalArgumentException("Unknonw access mode: " + accessMode); } @@ -1867,6 +1869,8 @@ public static ProducerAccessMode convertProducerAccessMode( return ProducerAccessMode.Shared; case WaitForExclusive: return ProducerAccessMode.WaitForExclusive; + case ExclusiveWithFencing: + return ProducerAccessMode.ExclusiveWithFencing; default: throw new IllegalArgumentException("Unknonw access mode: " + accessMode); } diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index a5d97e51acfd7..a2bfd76c92c4a 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -99,6 +99,7 @@ enum ProducerAccessMode { Shared = 0; // By default multiple producers can publish on a topic Exclusive = 1; // Require exclusive access for producer. Fail immediately if there's already a producer connected. WaitForExclusive = 2; // Producer creation is pending until it can acquire exclusive access + ExclusiveWithFencing = 3; // Require exclusive access for producer. Fence out old producer. } message MessageMetadata { From 227b7b960dd5bd098d29ea9ecc9b8bf9093b96ec Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 9 May 2022 10:28:28 +0200 Subject: [PATCH 2/6] Fix behaviour in presence of Waiting Producers --- .../pulsar/broker/service/AbstractTopic.java | 17 +++- .../broker/service/ExclusiveProducerTest.java | 78 +++++++++++++++++++ .../pulsar/client/api/ProducerAccessMode.java | 2 +- .../pulsar/client/api/ProducerBuilder.java | 4 +- 4 files changed, 96 insertions(+), 5 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 61a3a78a790a0..e101dc4db3ba1 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 @@ -22,6 +22,8 @@ import static org.apache.bookkeeper.mledger.impl.ManagedLedgerMBeanImpl.ENTRY_LATENCY_BUCKETS_USEC; import com.google.common.base.MoreObjects; import com.google.common.collect.Lists; + +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.EnumSet; @@ -739,9 +741,20 @@ protected CompletableFuture> incrementTopicEpochIfNeeded(Producer }); } case ExclusiveWithFencing: - if (hasExclusiveProducer) { + if (hasExclusiveProducer || !producers.isEmpty()) { + // clear all waiting producers + // otherwise closing any producer will trigger the promotion + // of the next pending producer + List>>> waitingExclusiveProducersCopy = + new ArrayList<>(waitingExclusiveProducers); + waitingExclusiveProducers.clear(); + waitingExclusiveProducersCopy.forEach((Pair>> handle) -> { + log.info("[{}] Failing waiting producer {}", topic, handle.getKey()); + handle.getValue().completeExceptionally(new ProducerFencedException("Fenced out")); + handle.getKey().close(true); + }); producers.forEach((k, currentProducer) -> { - log.info("Fencing out producer {}", currentProducer); + log.info("[{}] Fencing out producer {}", topic, currentProducer); currentProducer.close(true); }); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java index 2ed63a8991f63..71bf10493d3be 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java @@ -19,10 +19,12 @@ package org.apache.pulsar.broker.service; import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; import java.util.Optional; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import io.netty.util.HashedWheelTimer; @@ -90,12 +92,14 @@ private void simpleTest(String topic) throws Exception { Producer p1 = pulsarClient.newProducer(Schema.STRING) .topic(topic) + .producerName("p1") .accessMode(ProducerAccessMode.Exclusive) .create(); try { pulsarClient.newProducer(Schema.STRING) .topic(topic) + .producerName("p-fail-1") .accessMode(ProducerAccessMode.Exclusive) .create(); fail("Should have failed"); @@ -106,6 +110,7 @@ private void simpleTest(String topic) throws Exception { try { pulsarClient.newProducer(Schema.STRING) .topic(topic) + .producerName("p-fail-2") .accessMode(ProducerAccessMode.Shared) .create(); fail("Should have failed"); @@ -118,11 +123,13 @@ private void simpleTest(String topic) throws Exception { // Now producer should be allowed to get in Producer p2 = pulsarClient.newProducer(Schema.STRING) .topic(topic) + .producerName("p2") .accessMode(ProducerAccessMode.Exclusive) .create(); Producer p3 = pulsarClient.newProducer(Schema.STRING) .topic(topic) + .producerName("p3") .accessMode(ProducerAccessMode.ExclusiveWithFencing) .create(); @@ -136,6 +143,77 @@ private void simpleTest(String topic) throws Exception { p3.send("test"); p3.close(); + // test now WaitForExclusive vs ExclusiveWithFencing + + // use two different Clients, because sometimes fencing a Producer triggers connection close + // making the test unreliable. + + @Cleanup + PulsarClient pulsarClient2 = PulsarClient.builder() + .serviceUrl(pulsar.getBrokerServiceUrl()) + .operationTimeout(2, TimeUnit.SECONDS) + .build(); + + Producer p4 = pulsarClient2.newProducer(Schema.STRING) + .topic(topic) + .producerName("p4") + .accessMode(ProducerAccessMode.Exclusive) + .create(); + + p4.send("test"); + + // p5 will be waiting for the lock to be released + CompletableFuture> p5 = pulsarClient2.newProducer(Schema.STRING) + .topic(topic) + .producerName("p5") + .accessMode(ProducerAccessMode.WaitForExclusive) + .createAsync(); + + // p6 fences out all the current producers, even p5 + Producer p6 = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .producerName("p6") + .accessMode(ProducerAccessMode.ExclusiveWithFencing) + .create(); + + p6.send("test"); + + // p7 is enqueued after p6 + CompletableFuture> p7 = pulsarClient2.newProducer(Schema.STRING) + .topic(topic) + .producerName("p7") + .accessMode(ProducerAccessMode.WaitForExclusive) + .createAsync(); + + // this should work, p6 is the owner + p6.send("test"); + + try { + p4.send("test"); + fail("Should have failed"); + } catch (ProducerFencedException expected) { + } + + // this should work, p6 is the owner + p6.send("test"); + + // p5 fails + try { + p5.get(); + fail("Should have failed"); + } catch (ExecutionException expected) { + assertTrue(expected.getCause() instanceof ProducerFencedException, + "unexpected exception " + expected.getCause()); + } + + // this should work, p6 is the owner + p6.send("test"); + + p6.close(); + + // p7 finally acquires the lock + p7.get().send("test"); + p7.get().close(); } @Test(dataProvider = "topics") diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java index c7852214fdee6..cfe3594c52ca8 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java @@ -34,7 +34,7 @@ public enum ProducerAccessMode { Exclusive, /** - * Require exclusive access for producer. Fence out the old producer. + * Acquire exclusive access for the producer. Any existing producer will be immediately. */ ExclusiveWithFencing, 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 89c330338d8d2..1892b6650f229 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 @@ -133,8 +133,8 @@ public interface ProducerBuilder extends Cloneable { *
  • {@link ProducerAccessMode#Shared}: By default multiple producers can publish on a topic *
  • {@link ProducerAccessMode#Exclusive}: Require exclusive access for producer. Fail immediately if there's * already a producer connected. - *
  • {@link ProducerAccessMode#ExclusiveWithFencing}: Require exclusive access for producer. Fence out any - * producer that is connected. + *
  • {@link ProducerAccessMode#ExclusiveWithFencing}: Require exclusive access for the producer. + * Any existing producer will be immediately.. *
  • {@link ProducerAccessMode#WaitForExclusive}: Producer creation is pending until it can acquire exclusive * access * From a53b1df11ca433a0c3ce60d9602a236c4a2f0666 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Mon, 9 May 2022 11:02:25 +0200 Subject: [PATCH 3/6] checkstyle --- .../java/org/apache/pulsar/broker/service/AbstractTopic.java | 4 ++-- 1 file changed, 2 insertions(+), 2 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 e101dc4db3ba1..89bffcca3c041 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 @@ -22,7 +22,6 @@ import static org.apache.bookkeeper.mledger.impl.ManagedLedgerMBeanImpl.ENTRY_LATENCY_BUCKETS_USEC; import com.google.common.base.MoreObjects; import com.google.common.collect.Lists; - import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -748,7 +747,8 @@ protected CompletableFuture> incrementTopicEpochIfNeeded(Producer List>>> waitingExclusiveProducersCopy = new ArrayList<>(waitingExclusiveProducers); waitingExclusiveProducers.clear(); - waitingExclusiveProducersCopy.forEach((Pair>> handle) -> { + waitingExclusiveProducersCopy.forEach((Pair>> handle) -> { log.info("[{}] Failing waiting producer {}", topic, handle.getKey()); handle.getValue().completeExceptionally(new ProducerFencedException("Fenced out")); handle.getKey().close(true); From 42808ea1a076daaca7ed63df8501c2fb4e6fc9bd Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 10 May 2022 14:43:03 +0200 Subject: [PATCH 4/6] added some docs --- .../apache/pulsar/broker/service/ExclusiveProducerTest.java | 2 +- .../java/org/apache/pulsar/client/api/ProducerAccessMode.java | 3 ++- .../java/org/apache/pulsar/client/api/ProducerBuilder.java | 2 +- site2/docs/concepts-messaging.md | 1 + 4 files changed, 5 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java index 71bf10493d3be..e9e623be2e761 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java @@ -78,7 +78,7 @@ public static Object[][] accessMode() { { ProducerAccessMode.ExclusiveWithFencing, Boolean.TRUE}, { ProducerAccessMode.ExclusiveWithFencing, Boolean.FALSE }, { ProducerAccessMode.WaitForExclusive, Boolean.TRUE }, - { ProducerAccessMode.WaitForExclusive, Boolean.FALSE }, + { ProducerAccessMode.WaitForExclusive, Boolean.FALSE } }; } diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java index cfe3594c52ca8..33492802d161b 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerAccessMode.java @@ -34,7 +34,8 @@ public enum ProducerAccessMode { Exclusive, /** - * Acquire exclusive access for the producer. Any existing producer will be immediately. + * Acquire exclusive access for the producer. Any existing producer will be removed and + * invalidated immediately. */ ExclusiveWithFencing, 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 1892b6650f229..d1778685f84c1 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 @@ -134,7 +134,7 @@ public interface ProducerBuilder extends Cloneable { *
  • {@link ProducerAccessMode#Exclusive}: Require exclusive access for producer. Fail immediately if there's * already a producer connected. *
  • {@link ProducerAccessMode#ExclusiveWithFencing}: Require exclusive access for the producer. - * Any existing producer will be immediately.. + * Any existing producer will be removed and invalidated immediately. *
  • {@link ProducerAccessMode#WaitForExclusive}: Producer creation is pending until it can acquire exclusive * access * diff --git a/site2/docs/concepts-messaging.md b/site2/docs/concepts-messaging.md index a98e5eb3bd2f9..b4fdd7708e86c 100644 --- a/site2/docs/concepts-messaging.md +++ b/site2/docs/concepts-messaging.md @@ -65,6 +65,7 @@ You can have different types of access modes on topics for producers. |---|--- `Shared`|Multiple producers can publish on a topic.

    This is the **default** setting. `Exclusive`|Only one producer can publish on a topic.

    If there is already a producer connected, other producers trying to publish on this topic get errors immediately.

    The “old” producer is evicted and a “new” producer is selected to be the next exclusive producer if the “old” producer experiences a network partition with the broker. +`ExclusiveWithFencing`|Only one producer can publish on a topic.

    If there is already a producer connected, is will be removed and invalidated immediately. `WaitForExclusive`|If there is already a producer connected, the producer creation is pending (rather than timing out) until the producer gets the `Exclusive` access.

    The producer that succeeds in becoming the exclusive one is treated as the leader. Consequently, if you want to implement the leader election scheme for your application, you can use this access mode. > **Note** From fdd0ecf4826d8d7ee82962da5bd60c0342939f1d Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 10 May 2022 15:30:54 +0200 Subject: [PATCH 5/6] Make the test more reliable --- .../apache/pulsar/broker/service/ExclusiveProducerTest.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java index e9e623be2e761..604abd8d7095f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ExclusiveProducerTest.java @@ -169,6 +169,9 @@ private void simpleTest(String topic) throws Exception { .accessMode(ProducerAccessMode.WaitForExclusive) .createAsync(); + // wait for all the Producers to be enqueued in order to prevent races + Thread.sleep(2000); + // p6 fences out all the current producers, even p5 Producer p6 = pulsarClient.newProducer(Schema.STRING) .topic(topic) From cd434ffb25871b935f28df18df68d9520d7419c8 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Wed, 11 May 2022 09:30:45 +0200 Subject: [PATCH 6/6] Update site2/docs/concepts-messaging.md Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- site2/docs/concepts-messaging.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/site2/docs/concepts-messaging.md b/site2/docs/concepts-messaging.md index b4fdd7708e86c..02303ea93b7fb 100644 --- a/site2/docs/concepts-messaging.md +++ b/site2/docs/concepts-messaging.md @@ -65,7 +65,7 @@ You can have different types of access modes on topics for producers. |---|--- `Shared`|Multiple producers can publish on a topic.

    This is the **default** setting. `Exclusive`|Only one producer can publish on a topic.

    If there is already a producer connected, other producers trying to publish on this topic get errors immediately.

    The “old” producer is evicted and a “new” producer is selected to be the next exclusive producer if the “old” producer experiences a network partition with the broker. -`ExclusiveWithFencing`|Only one producer can publish on a topic.

    If there is already a producer connected, is will be removed and invalidated immediately. +`ExclusiveWithFencing`|Only one producer can publish on a topic.

    If there is already a producer connected, it will be removed and invalidated immediately. `WaitForExclusive`|If there is already a producer connected, the producer creation is pending (rather than timing out) until the producer gets the `Exclusive` access.

    The producer that succeeds in becoming the exclusive one is treated as the leader. Consequently, if you want to implement the leader election scheme for your application, you can use this access mode. > **Note**