From 0440e9a297dd05623d7ee503f0215f89f4ff68aa Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Wed, 19 Feb 2020 22:13:03 +0800 Subject: [PATCH 1/7] Creating a topic does not wait for creating cursor when creating replicator --- .../broker/service/AbstractReplicator.java | 80 +++++++++++-------- .../NonPersistentReplicator.java | 5 ++ .../persistent/PersistentReplicator.java | 37 ++++++++- .../service/persistent/PersistentTopic.java | 37 +++------ 4 files changed, 96 insertions(+), 63 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java index d9b2f8eee24fa..8602d9f42360a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java @@ -64,7 +64,7 @@ protected enum State { } public AbstractReplicator(String topicName, String replicatorPrefix, String localCluster, String remoteCluster, - BrokerService brokerService) throws NamingException { + BrokerService brokerService) throws NamingException { validatePartitionedTopic(topicName, brokerService); this.brokerService = brokerService; this.topicName = topicName; @@ -100,52 +100,62 @@ public String getRemoteCluster() { // This method needs to be synchronized with disconnects else if there is a disconnect followed by startProducer // the end result can be disconnect. public synchronized void startProducer() { - if (STATE_UPDATER.get(this) == State.Stopping) { - long waitTimeMs = backOff.next(); - if (log.isDebugEnabled()) { - log.debug( + prepareStartProducer().thenAccept(v -> { + if (STATE_UPDATER.get(this) == State.Stopping) { + long waitTimeMs = backOff.next(); + if (log.isDebugEnabled()) { + log.debug( "[{}][{} -> {}] waiting for producer to close before attempting to reconnect, retrying in {} s", topicName, localCluster, remoteCluster, waitTimeMs / 1000.0); - } - // BackOff before retrying - brokerService.executor().schedule(this::startProducer, waitTimeMs, TimeUnit.MILLISECONDS); - return; - } - State state = STATE_UPDATER.get(this); - if (!STATE_UPDATER.compareAndSet(this, State.Stopped, State.Starting)) { - if (state == State.Started) { - // Already running - if (log.isDebugEnabled()) { - log.debug("[{}][{} -> {}] Replicator was already running", topicName, localCluster, remoteCluster); } - } else { - log.info("[{}][{} -> {}] Replicator already being started. Replicator state: {}", topicName, - localCluster, remoteCluster, state); + // BackOff before retrying + brokerService.executor().schedule(this::startProducer, waitTimeMs, TimeUnit.MILLISECONDS); + return; } + State state = STATE_UPDATER.get(this); + if (!STATE_UPDATER.compareAndSet(this, State.Stopped, State.Starting)) { + if (state == State.Started) { + // Already running + if (log.isDebugEnabled()) { + log.debug("[{}][{} -> {}] Replicator was already running", topicName, localCluster, remoteCluster); + } + } else { + log.info("[{}][{} -> {}] Replicator already being started. Replicator state: {}", topicName, + localCluster, remoteCluster, state); + } - return; - } + return; + } - log.info("[{}][{} -> {}] Starting replicator", topicName, localCluster, remoteCluster); - producerBuilder.createAsync().thenAccept(producer -> { - readEntries(producer); + log.info("[{}][{} -> {}] Starting replicator", topicName, localCluster, remoteCluster); + producerBuilder.createAsync().thenAccept(producer -> { + readEntries(producer); + }).exceptionally(ex -> { + retryCreateProducer(ex); + return null; + }); }).exceptionally(ex -> { - if (STATE_UPDATER.compareAndSet(this, State.Starting, State.Stopped)) { - long waitTimeMs = backOff.next(); - log.warn("[{}][{} -> {}] Failed to create remote producer ({}), retrying in {} s", topicName, - localCluster, remoteCluster, ex.getMessage(), waitTimeMs / 1000.0); - - // BackOff before retrying - brokerService.executor().schedule(this::startProducer, waitTimeMs, TimeUnit.MILLISECONDS); - } else { - log.warn("[{}][{} -> {}] Failed to create remote producer. Replicator state: {}", topicName, - localCluster, remoteCluster, STATE_UPDATER.get(this), ex); - } + retryCreateProducer(ex); return null; }); + } + + private void retryCreateProducer(Throwable ex) { + if (STATE_UPDATER.compareAndSet(this, State.Starting, State.Stopped)) { + long waitTimeMs = backOff.next(); + log.warn("[{}][{} -> {}] Failed to create remote producer ({}), retrying in {} s", topicName, + localCluster, remoteCluster, ex.getMessage(), waitTimeMs / 1000.0); + // BackOff before retrying + brokerService.executor().schedule(this::startProducer, waitTimeMs, TimeUnit.MILLISECONDS); + } else { + log.warn("[{}][{} -> {}] Failed to create remote producer. Replicator state: {}", topicName, + localCluster, remoteCluster, STATE_UPDATER.get(this), ex); + } } + protected abstract CompletableFuture prepareStartProducer(); + protected synchronized CompletableFuture closeProducerAsync() { if (producer == null) { STATE_UPDATER.set(this, State.Stopped); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentReplicator.java index b6ea53ae9c682..663607407f722 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentReplicator.java @@ -251,6 +251,11 @@ protected void disableReplicatorRead() { // No-op } + @Override + protected CompletableFuture prepareStartProducer() { + return CompletableFuture.completedFuture(null); + } + @Override public boolean isConnected() { ProducerImpl producer = this.producer; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index c4a19c460b529..8001c93b8d559 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -33,11 +33,13 @@ import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.AsyncCallbacks.ClearBacklogCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteCallback; +import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenCursorCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntryCallback; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedCursor.IndividualDeletedEntries; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerException.CursorAlreadyClosedException; import org.apache.bookkeeper.mledger.ManagedLedgerException.TooManyRequestsException; @@ -46,6 +48,7 @@ import org.apache.pulsar.broker.service.AbstractReplicator; import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.BrokerServiceException.NamingException; +import org.apache.pulsar.broker.service.BrokerServiceException.PersistenceException; import org.apache.pulsar.broker.service.BrokerServiceException.TopicBusyException; import org.apache.pulsar.broker.service.Replicator; import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter.Type; @@ -55,6 +58,7 @@ import org.apache.pulsar.client.impl.MessageImpl; import org.apache.pulsar.client.impl.ProducerImpl; import org.apache.pulsar.client.impl.SendCallback; +import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosition; import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.api.proto.PulsarMarkers.MarkerType; import org.apache.pulsar.common.policies.data.ReplicatorStats; @@ -65,7 +69,9 @@ public class PersistentReplicator extends AbstractReplicator implements Replicator, ReadEntriesCallback, DeleteCallback { private final PersistentTopic topic; - private final ManagedCursor cursor; + private final String replicatorName; + private final ManagedLedger ledger; + protected ManagedCursor cursor; private Optional dispatchRateLimiter = Optional.empty(); @@ -97,11 +103,12 @@ public class PersistentReplicator extends AbstractReplicator implements Replicat private final ReplicatorStats stats = new ReplicatorStats(); - public PersistentReplicator(PersistentTopic topic, ManagedCursor cursor, String localCluster, String remoteCluster, - BrokerService brokerService) throws NamingException { + public PersistentReplicator(PersistentTopic topic, String replicatorName, String localCluster, String remoteCluster, + BrokerService brokerService, ManagedLedger ledger) throws NamingException { super(topic.getName(), topic.getReplicatorPrefix(), localCluster, remoteCluster, brokerService); + this.replicatorName = replicatorName; + this.ledger = ledger; this.topic = topic; - this.cursor = cursor; this.expiryMonitor = new PersistentMessageExpiryMonitor(topicName, Codec.decode(cursor.getName()), cursor); HAVE_PENDING_READ_UPDATER.set(this, FALSE); PENDING_MESSAGES_UPDATER.set(this, 0); @@ -158,6 +165,28 @@ protected void disableReplicatorRead() { this.cursor.setInactive(); } + @Override + protected synchronized CompletableFuture prepareStartProducer() { + if (cursor != null) { + return CompletableFuture.completedFuture(null); + } + CompletableFuture res = new CompletableFuture<>(); + ledger.asyncOpenCursor(replicatorName, InitialPosition.Earliest, new OpenCursorCallback() { + @Override + public void openCursorComplete(ManagedCursor cursor, Object ctx) { + PersistentReplicator.this.cursor = cursor; + res.complete(null); + } + + @Override + public void openCursorFailed(ManagedLedgerException exception, Object ctx) { + res.completeExceptionally(new PersistenceException(exception)); + } + + }, null); + return res; + } + /** * Calculate available permits for read entries. 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 335be8cf3a847..c950ef267b3a7 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 @@ -224,7 +224,7 @@ public PersistentTopic(String topic, ManagedLedger ledger, BrokerService brokerS if (cursor.getName().startsWith(replicatorPrefix)) { String localCluster = brokerService.pulsar().getConfiguration().getClusterName(); String remoteCluster = PersistentReplicator.getRemoteCluster(cursor.getName()); - boolean isReplicatorStarted = addReplicationCluster(remoteCluster, this, cursor, localCluster); + boolean isReplicatorStarted = addReplicationCluster(remoteCluster, this, cursor.getName(), localCluster); if (!isReplicatorStarted) { throw new NamingException( PersistentTopic.this.getName() + " Failed to start replicator " + remoteCluster); @@ -1189,37 +1189,26 @@ CompletableFuture startReplicator(String remoteCluster) { log.info("[{}] Starting replicator to remote: {}", topic, remoteCluster); final CompletableFuture future = new CompletableFuture<>(); - String name = PersistentReplicator.getReplicatorName(replicatorPrefix, remoteCluster); - ledger.asyncOpenCursor(name, new OpenCursorCallback() { - @Override - public void openCursorComplete(ManagedCursor cursor, Object ctx) { - String localCluster = brokerService.pulsar().getConfiguration().getClusterName(); - boolean isReplicatorStarted = addReplicationCluster(remoteCluster, PersistentTopic.this, cursor, localCluster); - if (isReplicatorStarted) { - future.complete(null); - } else { - future.completeExceptionally(new NamingException( - PersistentTopic.this.getName() + " Failed to start replicator " + remoteCluster)); - } - } - - @Override - public void openCursorFailed(ManagedLedgerException exception, Object ctx) { - future.completeExceptionally(new PersistenceException(exception)); - } - - }, null); + String replicatorName = PersistentReplicator.getReplicatorName(replicatorPrefix, remoteCluster); + String localCluster = brokerService.pulsar().getConfiguration().getClusterName(); + boolean isReplicatorStarted = addReplicationCluster(remoteCluster, PersistentTopic.this, replicatorName, localCluster); + if (isReplicatorStarted) { + future.complete(null); + } else { + future.completeExceptionally(new NamingException( + PersistentTopic.this.getName() + " Failed to start replicator " + remoteCluster)); + } return future; } - protected boolean addReplicationCluster(String remoteCluster, PersistentTopic persistentTopic, ManagedCursor cursor, + protected boolean addReplicationCluster(String remoteCluster, PersistentTopic persistentTopic, String replicatorName, String localCluster) { AtomicBoolean isReplicatorStarted = new AtomicBoolean(true); replicators.computeIfAbsent(remoteCluster, r -> { try { - return new PersistentReplicator(PersistentTopic.this, cursor, localCluster, remoteCluster, - brokerService); + return new PersistentReplicator(PersistentTopic.this, replicatorName, localCluster, remoteCluster, + brokerService, ledger); } catch (NamingException e) { isReplicatorStarted.set(false); log.error("[{}] Replicator startup failed due to partitioned-topic {}", topic, remoteCluster); From a562bbe28cc8c9250b36a76bdd13ae59a3f51db5 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Wed, 19 Feb 2020 22:26:16 +0800 Subject: [PATCH 2/7] Fix unit test. --- .../persistent/PersistentReplicator.java | 26 ++++++++++++++++++- 1 file changed, 25 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index 8001c93b8d559..3095ae16858ba 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -19,6 +19,8 @@ package org.apache.pulsar.broker.service.persistent; import static org.apache.pulsar.broker.service.persistent.PersistentTopic.MESSAGE_RATE_BACKOFF_MS; + +import com.google.common.annotations.VisibleForTesting; import io.netty.buffer.ByteBuf; import io.netty.util.Recycler; import io.netty.util.Recycler.Handle; @@ -103,8 +105,30 @@ public class PersistentReplicator extends AbstractReplicator implements Replicat private final ReplicatorStats stats = new ReplicatorStats(); + // Only for test + public PersistentReplicator(PersistentTopic topic, ManagedCursor cursor, String localCluster, String remoteCluster, + BrokerService brokerService) throws NamingException { + super(topic.getName(), topic.getReplicatorPrefix(), localCluster, remoteCluster, brokerService); + this.replicatorName = cursor.getName(); + this.ledger = cursor.getManagedLedger(); + this.cursor = cursor; + this.topic = topic; + this.expiryMonitor = new PersistentMessageExpiryMonitor(topicName, Codec.decode(cursor.getName()), cursor); + HAVE_PENDING_READ_UPDATER.set(this, FALSE); + PENDING_MESSAGES_UPDATER.set(this, 0); + + readBatchSize = Math.min( + producerQueueSize, + topic.getBrokerService().pulsar().getConfiguration().getDispatcherMaxReadBatchSize()); + producerQueueThreshold = (int) (producerQueueSize * 0.9); + + this.initializeDispatchRateLimiterIfNeeded(Optional.empty()); + + startProducer(); + } + public PersistentReplicator(PersistentTopic topic, String replicatorName, String localCluster, String remoteCluster, - BrokerService brokerService, ManagedLedger ledger) throws NamingException { + BrokerService brokerService, ManagedLedger ledger) throws NamingException { super(topic.getName(), topic.getReplicatorPrefix(), localCluster, remoteCluster, brokerService); this.replicatorName = replicatorName; this.ledger = ledger; From 86e728e9d6aca44a9540cb6892e5862c607b96af Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Wed, 19 Feb 2020 23:54:03 +0800 Subject: [PATCH 3/7] Fix tests. --- .../service/persistent/PersistentReplicator.java | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index 3095ae16858ba..7966e48d432ea 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -133,7 +133,6 @@ public PersistentReplicator(PersistentTopic topic, String replicatorName, String this.replicatorName = replicatorName; this.ledger = ledger; this.topic = topic; - this.expiryMonitor = new PersistentMessageExpiryMonitor(topicName, Codec.decode(cursor.getName()), cursor); HAVE_PENDING_READ_UPDATER.set(this, FALSE); PENDING_MESSAGES_UPDATER.set(this, 0); @@ -192,6 +191,9 @@ protected void disableReplicatorRead() { @Override protected synchronized CompletableFuture prepareStartProducer() { if (cursor != null) { + if (expiryMonitor == null) { + this.expiryMonitor = new PersistentMessageExpiryMonitor(topicName, Codec.decode(cursor.getName()), cursor); + } return CompletableFuture.completedFuture(null); } CompletableFuture res = new CompletableFuture<>(); @@ -199,6 +201,7 @@ protected synchronized CompletableFuture prepareStartProducer() { @Override public void openCursorComplete(ManagedCursor cursor, Object ctx) { PersistentReplicator.this.cursor = cursor; + PersistentReplicator.this.expiryMonitor = new PersistentMessageExpiryMonitor(topicName, Codec.decode(cursor.getName()), cursor); res.complete(null); } @@ -654,7 +657,9 @@ public void updateRates() { msgExpired.calculateRate(); stats.msgRateOut = msgOut.getRate(); stats.msgThroughputOut = msgOut.getValueRate(); - stats.msgRateExpired = msgExpired.getRate() + expiryMonitor.getMessageExpiryRate(); + if (expiryMonitor != null) { + stats.msgRateExpired = msgExpired.getRate() + expiryMonitor.getMessageExpiryRate(); + } } public ReplicatorStats getStats() { @@ -692,7 +697,9 @@ public void expireMessages(int messageTTLInSeconds) { // don't do anything for almost caught-up connected subscriptions return; } - expiryMonitor.expireMessages(messageTTLInSeconds); + if (expiryMonitor != null) { + expiryMonitor.expireMessages(messageTTLInSeconds); + } } @Override From 370217cd0586ce8e3c7238208fd62b6221977769 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Thu, 20 Feb 2020 08:54:45 +0800 Subject: [PATCH 4/7] Add log for open cursor in replicator --- .../broker/service/AbstractReplicator.java | 63 +++++++++---------- .../NonPersistentReplicator.java | 2 +- .../persistent/PersistentReplicator.java | 6 +- 3 files changed, 36 insertions(+), 35 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java index 8602d9f42360a..b4270a30d68c3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java @@ -100,41 +100,38 @@ public String getRemoteCluster() { // This method needs to be synchronized with disconnects else if there is a disconnect followed by startProducer // the end result can be disconnect. public synchronized void startProducer() { - prepareStartProducer().thenAccept(v -> { - if (STATE_UPDATER.get(this) == State.Stopping) { - long waitTimeMs = backOff.next(); - if (log.isDebugEnabled()) { - log.debug( - "[{}][{} -> {}] waiting for producer to close before attempting to reconnect, retrying in {} s", - topicName, localCluster, remoteCluster, waitTimeMs / 1000.0); - } - // BackOff before retrying - brokerService.executor().schedule(this::startProducer, waitTimeMs, TimeUnit.MILLISECONDS); - return; + log.info("[{}][{} -> {}] Starting replicator", topicName, localCluster, remoteCluster); + if (STATE_UPDATER.get(this) == State.Stopping) { + long waitTimeMs = backOff.next(); + if (log.isDebugEnabled()) { + log.debug( + "[{}][{} -> {}] waiting for producer to close before attempting to reconnect, retrying in {} s", + topicName, localCluster, remoteCluster, waitTimeMs / 1000.0); } - State state = STATE_UPDATER.get(this); - if (!STATE_UPDATER.compareAndSet(this, State.Stopped, State.Starting)) { - if (state == State.Started) { - // Already running - if (log.isDebugEnabled()) { - log.debug("[{}][{} -> {}] Replicator was already running", topicName, localCluster, remoteCluster); - } - } else { - log.info("[{}][{} -> {}] Replicator already being started. Replicator state: {}", topicName, - localCluster, remoteCluster, state); + // BackOff before retrying + brokerService.executor().schedule(this::startProducer, waitTimeMs, TimeUnit.MILLISECONDS); + return; + } + State state = STATE_UPDATER.get(this); + if (!STATE_UPDATER.compareAndSet(this, State.Stopped, State.Starting)) { + if (state == State.Started) { + // Already running + if (log.isDebugEnabled()) { + log.debug("[{}][{} -> {}] Replicator was already running", topicName, localCluster, remoteCluster); } - - return; + } else { + log.info("[{}][{} -> {}] Replicator already being started. Replicator state: {}", topicName, + localCluster, remoteCluster, state); } - - log.info("[{}][{} -> {}] Starting replicator", topicName, localCluster, remoteCluster); - producerBuilder.createAsync().thenAccept(producer -> { - readEntries(producer); - }).exceptionally(ex -> { - retryCreateProducer(ex); - return null; - }); - }).exceptionally(ex -> { + return; + } + openCursorAsync().thenAccept(v -> + producerBuilder.createAsync() + .thenAccept(this::readEntries) + .exceptionally(ex -> { + retryCreateProducer(ex); + return null; + })).exceptionally(ex -> { retryCreateProducer(ex); return null; }); @@ -154,7 +151,7 @@ private void retryCreateProducer(Throwable ex) { } } - protected abstract CompletableFuture prepareStartProducer(); + protected abstract CompletableFuture openCursorAsync(); protected synchronized CompletableFuture closeProducerAsync() { if (producer == null) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentReplicator.java index 663607407f722..c109560f9f8ec 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentReplicator.java @@ -252,7 +252,7 @@ protected void disableReplicatorRead() { } @Override - protected CompletableFuture prepareStartProducer() { + protected CompletableFuture openCursorAsync() { return CompletableFuture.completedFuture(null); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index 7966e48d432ea..b12e92dbbb42b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -189,8 +189,10 @@ protected void disableReplicatorRead() { } @Override - protected synchronized CompletableFuture prepareStartProducer() { + protected synchronized CompletableFuture openCursorAsync() { + log.info("[{}][{} -> {}] Starting open cursor for replicator", topicName, localCluster, remoteCluster); if (cursor != null) { + log.info("[{}][{} -> {}] Using the exists cursor for replicator", topicName, localCluster, remoteCluster); if (expiryMonitor == null) { this.expiryMonitor = new PersistentMessageExpiryMonitor(topicName, Codec.decode(cursor.getName()), cursor); } @@ -200,6 +202,7 @@ protected synchronized CompletableFuture prepareStartProducer() { ledger.asyncOpenCursor(replicatorName, InitialPosition.Earliest, new OpenCursorCallback() { @Override public void openCursorComplete(ManagedCursor cursor, Object ctx) { + log.info("[{}][{} -> {}] Open cursor succeed for replicator", topicName, localCluster, remoteCluster); PersistentReplicator.this.cursor = cursor; PersistentReplicator.this.expiryMonitor = new PersistentMessageExpiryMonitor(topicName, Codec.decode(cursor.getName()), cursor); res.complete(null); @@ -207,6 +210,7 @@ public void openCursorComplete(ManagedCursor cursor, Object ctx) { @Override public void openCursorFailed(ManagedLedgerException exception, Object ctx) { + log.warn("[{}][{} -> {}] Open cursor failed for replicator", topicName, localCluster, remoteCluster, exception); res.completeExceptionally(new PersistenceException(exception)); } From dc450ea36439048b4c0a5e7e03d688deea592bab Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Thu, 20 Feb 2020 08:59:17 +0800 Subject: [PATCH 5/7] re-format code --- .../apache/pulsar/broker/service/AbstractReplicator.java | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java index b4270a30d68c3..6ace803cdc98f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java @@ -100,13 +100,12 @@ public String getRemoteCluster() { // This method needs to be synchronized with disconnects else if there is a disconnect followed by startProducer // the end result can be disconnect. public synchronized void startProducer() { - log.info("[{}][{} -> {}] Starting replicator", topicName, localCluster, remoteCluster); if (STATE_UPDATER.get(this) == State.Stopping) { long waitTimeMs = backOff.next(); if (log.isDebugEnabled()) { log.debug( - "[{}][{} -> {}] waiting for producer to close before attempting to reconnect, retrying in {} s", - topicName, localCluster, remoteCluster, waitTimeMs / 1000.0); + "[{}][{} -> {}] waiting for producer to close before attempting to reconnect, retrying in {} s", + topicName, localCluster, remoteCluster, waitTimeMs / 1000.0); } // BackOff before retrying brokerService.executor().schedule(this::startProducer, waitTimeMs, TimeUnit.MILLISECONDS); @@ -121,10 +120,12 @@ public synchronized void startProducer() { } } else { log.info("[{}][{} -> {}] Replicator already being started. Replicator state: {}", topicName, - localCluster, remoteCluster, state); + localCluster, remoteCluster, state); } return; } + + log.info("[{}][{} -> {}] Starting replicator", topicName, localCluster, remoteCluster); openCursorAsync().thenAccept(v -> producerBuilder.createAsync() .thenAccept(this::readEntries) From e1e32d58f6ab077b94ecc36f974723f129aa0c2d Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Thu, 20 Feb 2020 09:01:34 +0800 Subject: [PATCH 6/7] re-format code --- .../org/apache/pulsar/broker/service/AbstractReplicator.java | 2 +- .../broker/service/persistent/PersistentReplicator.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java index 6ace803cdc98f..13cd0918ab58a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractReplicator.java @@ -64,7 +64,7 @@ protected enum State { } public AbstractReplicator(String topicName, String replicatorPrefix, String localCluster, String remoteCluster, - BrokerService brokerService) throws NamingException { + BrokerService brokerService) throws NamingException { validatePartitionedTopic(topicName, brokerService); this.brokerService = brokerService; this.topicName = topicName; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index b12e92dbbb42b..7fdb18a5eb7f3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -107,7 +107,7 @@ public class PersistentReplicator extends AbstractReplicator implements Replicat // Only for test public PersistentReplicator(PersistentTopic topic, ManagedCursor cursor, String localCluster, String remoteCluster, - BrokerService brokerService) throws NamingException { + BrokerService brokerService) throws NamingException { super(topic.getName(), topic.getReplicatorPrefix(), localCluster, remoteCluster, brokerService); this.replicatorName = cursor.getName(); this.ledger = cursor.getManagedLedger(); @@ -128,7 +128,7 @@ public PersistentReplicator(PersistentTopic topic, ManagedCursor cursor, String } public PersistentReplicator(PersistentTopic topic, String replicatorName, String localCluster, String remoteCluster, - BrokerService brokerService, ManagedLedger ledger) throws NamingException { + BrokerService brokerService, ManagedLedger ledger) throws NamingException { super(topic.getName(), topic.getReplicatorPrefix(), localCluster, remoteCluster, brokerService); this.replicatorName = replicatorName; this.ledger = ledger; From 28af75587cef85ef2e6fe30c52e86b230f451841 Mon Sep 17 00:00:00 2001 From: Penghui Li Date: Thu, 20 Feb 2020 09:03:03 +0800 Subject: [PATCH 7/7] remove unused import --- .../pulsar/broker/service/persistent/PersistentReplicator.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index 7fdb18a5eb7f3..2e95eb6654f66 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -20,7 +20,6 @@ import static org.apache.pulsar.broker.service.persistent.PersistentTopic.MESSAGE_RATE_BACKOFF_MS; -import com.google.common.annotations.VisibleForTesting; import io.netty.buffer.ByteBuf; import io.netty.util.Recycler; import io.netty.util.Recycler.Handle;