diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index d76553877fbef..d90624ebd1825 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -899,7 +899,8 @@ public void operationComplete() { log.info("[{}] reset position to {} skipping from current read position {} on cursor {}", ledger.getName(), newPosition, oldReadPosition, name); } - readPosition = newPosition; + + rewind(); } finally { lock.writeLock().unlock(); } @@ -910,8 +911,9 @@ public void operationComplete() { ledger.getName(), newPosition, name); } } - callback.resetComplete(newPosition); + callback.resetComplete(newPosition); + notifyEntriesAvailable(); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index 200d06e63cec4..4a31be589c469 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -683,70 +683,43 @@ private void resetCursor(Position finalPosition, CompletableFuture future) return; } - final CompletableFuture disconnectFuture; - - // Lock the Subscription object before locking the Dispatcher object to avoid deadlocks - synchronized (this) { - if (dispatcher != null && dispatcher.isConsumerConnected()) { - disconnectFuture = dispatcher.disconnectAllConsumers(true); - } else { - disconnectFuture = CompletableFuture.completedFuture(null); - } - } - - disconnectFuture.whenComplete((aVoid, throwable) -> { - if (dispatcher != null) { - dispatcher.resetCloseFuture(); - } - - if (throwable != null) { - log.error("[{}][{}] Failed to disconnect consumer from subscription", topicName, subName, throwable); - IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); - future.completeExceptionally( - new SubscriptionBusyException("Failed to disconnect consumers from subscription")); - return; - } - - log.info("[{}][{}] Successfully disconnected consumers from subscription, proceeding with cursor reset", - topicName, subName); + try { + cursor.asyncResetCursor(finalPosition, new AsyncCallbacks.ResetCursorCallback() { + @Override + public void resetComplete(Object ctx) { + if (log.isDebugEnabled()) { + log.debug("[{}][{}] Successfully reset subscription to position {}", topicName, subName, + finalPosition); + } - try { - cursor.asyncResetCursor(finalPosition, new AsyncCallbacks.ResetCursorCallback() { - @Override - public void resetComplete(Object ctx) { - if (log.isDebugEnabled()) { - log.debug("[{}][{}] Successfully reset subscription to position {}", topicName, subName, - finalPosition); - } - if (dispatcher != null) { - dispatcher.cursorIsReset(); - } - IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); - future.complete(null); + if (dispatcher != null) { + dispatcher.cursorIsReset(); } + IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); + future.complete(null); + } - @Override - public void resetFailed(ManagedLedgerException exception, Object ctx) { - log.error("[{}][{}] Failed to reset subscription to position {}", topicName, subName, - finalPosition, exception); - IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); - // todo - retry on InvalidCursorPositionException - // or should we just ask user to retry one more time? - if (exception instanceof InvalidCursorPositionException) { - future.completeExceptionally(new SubscriptionInvalidCursorPosition(exception.getMessage())); - } else if (exception instanceof ConcurrentFindCursorPositionException) { - future.completeExceptionally(new SubscriptionBusyException(exception.getMessage())); - } else { - future.completeExceptionally(new BrokerServiceException(exception)); - } + @Override + public void resetFailed(ManagedLedgerException exception, Object ctx) { + log.error("[{}][{}] Failed to reset subscription to position {}", topicName, subName, + finalPosition, exception); + IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); + // todo - retry on InvalidCursorPositionException + // or should we just ask user to retry one more time? + if (exception instanceof InvalidCursorPositionException) { + future.completeExceptionally(new SubscriptionInvalidCursorPosition(exception.getMessage())); + } else if (exception instanceof ConcurrentFindCursorPositionException) { + future.completeExceptionally(new SubscriptionBusyException(exception.getMessage())); + } else { + future.completeExceptionally(new BrokerServiceException(exception)); } - }); - } catch (Exception e) { - log.error("[{}][{}] Error while resetting cursor", topicName, subName, e); - IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); - future.completeExceptionally(new BrokerServiceException(e)); - } - }); + } + }); + } catch (Exception e) { + log.error("[{}][{}] Error while resetting cursor", topicName, subName, e); + IS_FENCED_UPDATER.set(PersistentSubscription.this, FALSE); + future.completeExceptionally(new BrokerServiceException(e)); + } } @Override diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java index b83c5c29b6af4..729ab65eb9533 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java @@ -1748,7 +1748,6 @@ public void persistentTopicsInvalidCursorReset() throws Exception { } admin.topics().resetCursor(topicName, "my-sub", System.currentTimeMillis() + 90000); - consumer = client.newConsumer().topic(topicName).subscriptionName("my-sub").subscribe(); consumer.close(); client.close(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/v1/V1_AdminApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/v1/V1_AdminApiTest.java index cf5db2b660220..4e4a2881ec338 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/v1/V1_AdminApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/v1/V1_AdminApiTest.java @@ -1665,7 +1665,6 @@ public void persistentTopicsInvalidCursorReset() throws Exception { } admin.topics().resetCursor(topicName, "my-sub", System.currentTimeMillis() + 90000); - consumer = client.newConsumer().topic(topicName).subscriptionName("my-sub").subscribe(); consumer.close(); client.close(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/TopicReaderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/TopicReaderTest.java index 69139b5d04d72..5a2bb33691760 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/TopicReaderTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/TopicReaderTest.java @@ -25,17 +25,22 @@ import com.google.common.collect.Lists; import com.google.common.collect.Sets; + import java.io.IOException; import java.nio.file.Files; import java.nio.file.Paths; -import java.util.*; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.Random; +import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; + import org.apache.pulsar.client.impl.BatchMessageIdImpl; import org.apache.pulsar.client.impl.MessageImpl; import org.apache.pulsar.client.impl.ReaderImpl; import org.apache.pulsar.common.policies.data.TopicStats; -import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.RelativeTimeUtil; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -638,7 +643,7 @@ public void testReaderIsAbleToSeekWithMessageIdOnMiddleOfTopic() throws Exceptio // Read all halved messages after seek() Set messageSetB = Sets.newHashSet(); - for (int i = halfMessages + 1; i < numOfMessage; i++) { + for (int i = halfMessages; i < numOfMessage; i++) { Message message = reader.readNext(); String receivedMessage = new String(message.getData()); String expectedMessage = String.format("msg num %d", i);