Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
117df78
PIP-307: Implement extensible load manager consumer changes
dragosvictor Dec 5, 2023
d2ba22f
Merge relevant changes in PR-21668
dragosvictor Dec 5, 2023
3ea12a2
Merge remote-tracking branch 'dmisca/master' into pip-307-consumer
dragosvictor Dec 5, 2023
eb069f4
Reset invocationCount on testTransferNonPersistentClientReconnectionW…
dragosvictor Dec 5, 2023
7f3ab3b
Merge remote-tracking branch 'origin/master' into pip-307-consumer
dragosvictor Dec 6, 2023
04da2b8
Speedup tests
dragosvictor Dec 6, 2023
68731c1
Add disconnectConsumers parameter to Dispatcher.close and Subscriptio…
dragosvictor Dec 7, 2023
58cb41d
Rename Subscription.disconnect to Subscription.close
dragosvictor Dec 7, 2023
4b480ec
Draft fix
dragosvictor Dec 7, 2023
94fdee4
Reimplement Subscription.disconnect
dragosvictor Dec 7, 2023
c205665
Skip reading more entries for transferring topics
dragosvictor Dec 8, 2023
bd94002
Add Javadoc comments
dragosvictor Dec 8, 2023
3c82142
Optimize dispatch condition in PersistentDispatcherMultipleConsumers#…
dragosvictor Dec 8, 2023
c2cb6ba
Add comments for PersistentSubscription.closeCursor callers
dragosvictor Dec 8, 2023
6203e3d
Ignore message acks during topic transfer
dragosvictor Dec 11, 2023
ae14bd5
Merge remote-tracking branch 'origin/master' into pip-307-consumer
dragosvictor Dec 11, 2023
52b7db3
Fix test org.apache.pulsar.broker.transaction.pendingack.PendingAckPe…
dragosvictor Dec 11, 2023
9483304
Disable stack trace collection for BrokerServiceException.TopicTransf…
dragosvictor Dec 11, 2023
2449f94
Revert "Disable stack trace collection for BrokerServiceException.Top…
dragosvictor Dec 11, 2023
2993154
Revert "Ignore message acks during topic transfer"
dragosvictor Dec 11, 2023
b376538
Ignore message acks during topic transfer
dragosvictor Dec 11, 2023
10e51ee
Factor out common logic between handleCloseProducer and handleCloseCo…
dragosvictor Dec 12, 2023
4c18e9e
Fix test org.apache.pulsar.client.impl.ClientCnxTest#testHandleCloseP…
dragosvictor Dec 12, 2023
f313140
Merge remote-tracking branch 'origin/master' into pip-307-consumer
dragosvictor Dec 12, 2023
903c494
Update pulsar-broker/src/main/java/org/apache/pulsar/broker/service/p…
dragosvictor Dec 13, 2023
7098d51
Update pulsar-broker/src/main/java/org/apache/pulsar/broker/service/p…
dragosvictor Dec 15, 2023
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import java.util.Map;
import java.util.NavigableMap;
import java.util.Objects;
import java.util.Optional;
import java.util.TreeMap;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CopyOnWriteArrayList;
Expand All @@ -32,6 +33,7 @@
import java.util.concurrent.atomic.AtomicReferenceFieldUpdater;
import org.apache.bookkeeper.mledger.ManagedCursor;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData;
import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException;
import org.apache.pulsar.broker.service.BrokerServiceException.ServerMetadataException;
import org.apache.pulsar.client.impl.Murmur3Hash32;
Expand Down Expand Up @@ -80,8 +82,6 @@ public AbstractDispatcherSingleActiveConsumer(SubType subscriptionType, int part

protected abstract void scheduleReadOnActiveConsumer();

protected abstract void readMoreEntries(Consumer consumer);

protected abstract void cancelPendingRead();

protected void notifyActiveConsumerChanged(Consumer activeConsumer) {
Expand Down Expand Up @@ -257,9 +257,12 @@ public synchronized boolean canUnsubscribe(Consumer consumer) {
return (consumers.size() == 1) && Objects.equals(consumer, ACTIVE_CONSUMER_UPDATER.get(this));
}

public CompletableFuture<Void> close() {
@Override
public CompletableFuture<Void> close(boolean disconnectConsumers,
Optional<BrokerLookupData> assignedBrokerLookupData) {
IS_CLOSED_UPDATER.set(this, TRUE);
return disconnectAllConsumers();
return disconnectConsumers
? disconnectAllConsumers(false, assignedBrokerLookupData) : CompletableFuture.completedFuture(null);
}

public boolean isClosed() {
Expand All @@ -268,15 +271,23 @@ public boolean isClosed() {

/**
* Disconnect all consumers on this dispatcher (server side close). This triggers channelInactive on the inbound
* handler which calls dispatcher.removeConsumer(), where the closeFuture is completed
* handler which calls dispatcher.removeConsumer(), where the closeFuture is completed.
*
Comment thread
heesung-sohn marked this conversation as resolved.
* @return
* @param isResetCursor
* Specifies if the cursor has been reset.
* @param assignedBrokerLookupData
* Optional target broker redirect information. Allows the consumer to quickly reconnect to a broker
* during bundle unloading.
*
* @return CompletableFuture indicating the completion of the operation.
*/
public synchronized CompletableFuture<Void> disconnectAllConsumers(boolean isResetCursor) {
@Override
public synchronized CompletableFuture<Void> disconnectAllConsumers(
boolean isResetCursor, Optional<BrokerLookupData> assignedBrokerLookupData) {
closeFuture = new CompletableFuture<>();

if (!consumers.isEmpty()) {
consumers.forEach(consumer -> consumer.disconnect(isResetCursor));
consumers.forEach(consumer -> consumer.disconnect(isResetCursor, assignedBrokerLookupData));
cancelPendingRead();
} else {
// no consumer connected, complete disconnect immediately
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@
import org.apache.commons.lang3.mutable.MutableInt;
import org.apache.commons.lang3.tuple.MutablePair;
import org.apache.pulsar.broker.authentication.AuthenticationDataSubscription;
import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData;
import org.apache.pulsar.broker.service.persistent.PersistentSubscription;
import org.apache.pulsar.broker.service.persistent.PersistentTopic;
import org.apache.pulsar.client.api.MessageId;
Expand Down Expand Up @@ -407,8 +408,12 @@ public void disconnect() {
}

public void disconnect(boolean isResetCursor) {
disconnect(isResetCursor, Optional.empty());
}

public void disconnect(boolean isResetCursor, Optional<BrokerLookupData> assignedBrokerLookupData) {
log.info("Disconnecting consumer: {}", this);
cnx.closeConsumer(this);
cnx.closeConsumer(this, assignedBrokerLookupData);
try {
close(isResetCursor);
} catch (BrokerServiceException e) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import org.apache.bookkeeper.mledger.impl.PositionImpl;
import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData;
import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter;
import org.apache.pulsar.common.api.proto.CommandSubscribe.SubType;
import org.apache.pulsar.common.api.proto.MessageMetadata;
Expand Down Expand Up @@ -49,7 +50,11 @@ public interface Dispatcher {
*
* @return
*/
CompletableFuture<Void> close();
default CompletableFuture<Void> close() {
return close(true, Optional.empty());
}

CompletableFuture<Void> close(boolean disconnectClients, Optional<BrokerLookupData> assignedBrokerLookupData);

boolean isClosed();

Expand All @@ -63,12 +68,17 @@ public interface Dispatcher {
*
* @return
*/
CompletableFuture<Void> disconnectAllConsumers(boolean isResetCursor);
default CompletableFuture<Void> disconnectAllConsumers(boolean isResetCursor) {
return disconnectAllConsumers(isResetCursor, Optional.empty());
}

default CompletableFuture<Void> disconnectAllConsumers() {
return disconnectAllConsumers(false);
}

CompletableFuture<Void> disconnectAllConsumers(boolean isResetCursor,
Optional<BrokerLookupData> assignedBrokerLookupData);

void resetCloseFuture();

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -497,6 +497,7 @@ public void completed(Exception exception, long ledgerId, long entryId) {
} else if (!(exception instanceof TopicClosedException)) {
// For TopicClosed exception there's no need to send explicit error, since the client was
// already notified
// For TopicClosingOrDeleting exception, a notification will be sent separately
long callBackSequenceId = Math.max(highestSequenceId, sequenceId);
producer.cnx.getCommandSender().sendSendError(producer.producerId, callBackSequenceId,
serverError, exception.getMessage());
Expand Down Expand Up @@ -718,7 +719,7 @@ public CompletableFuture<Void> disconnect() {
*/
public CompletableFuture<Void> disconnect(Optional<BrokerLookupData> assignedBrokerLookupData) {
if (!closeFuture.isDone() && isDisconnecting.compareAndSet(false, true)) {
log.info("Disconnecting producer: {}", this);
log.info("Disconnecting producer: {}, assignedBrokerLookupData: {}", this, assignedBrokerLookupData);
cnx.execute(() -> {
cnx.closeProducer(this, assignedBrokerLookupData);
closeNow(true);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import static org.apache.pulsar.broker.service.persistent.PersistentTopic.getMigratedClusterUrl;
import static org.apache.pulsar.common.api.proto.ProtocolVersion.v5;
import static org.apache.pulsar.common.protocol.Commands.DEFAULT_CONSUMER_EPOCH;
import static org.apache.pulsar.common.protocol.Commands.newCloseConsumer;
import static org.apache.pulsar.common.protocol.Commands.newLookupErrorResponse;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Strings;
Expand Down Expand Up @@ -1321,7 +1322,7 @@ protected void handleSubscribe(final CommandSubscribe subscribe) {
topicName, remoteAddress, consumerId);
}
consumers.remove(consumerId, consumerFuture);
closeConsumer(consumerId);
closeConsumer(consumerId, Optional.empty());
return null;
}
} else if (exception.getCause() instanceof BrokerServiceException) {
Expand Down Expand Up @@ -1766,7 +1767,7 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) {
// if the topic is transferring, we ignore send msg.
if (producer.getTopic().isTransferring()) {
long ignoredMsgCount = ExtensibleLoadManagerImpl.get(pulsar)
.getIgnoredSendMsgCounter().incrementAndGet();
.getIgnoredSendMsgCounter().addAndGet(send.getNumMessages());
if (log.isDebugEnabled()) {
log.debug("Ignored send msg from:{}:{} to fenced topic:{} while transferring."
+ " Ignored message count:{}.",
Expand Down Expand Up @@ -1842,6 +1843,15 @@ protected void handleAck(CommandAck ack) {

if (consumerFuture != null && consumerFuture.isDone() && !consumerFuture.isCompletedExceptionally()) {
Consumer consumer = consumerFuture.getNow(null);
Subscription subscription = consumer.getSubscription();
if (subscription.getTopic().isTransferring()) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do we have a test to cover in transferring acks?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not specifically, but the modified test covers ack behavior: https://github.com/apache/pulsar/pull/21682/files#diff-744119c61c9f6a1b786c3966acd1e0f63748985cf68a6a1a15418c8d9900a9a8R543-R547. The test does not succeed unless all messages sent during the transfer surface out on the consumers. Since there are 200 messages involved, at least one of them likely has the acknowledgement initially ignored.

// Message acks are silently ignored during topic transfer.
if (log.isDebugEnabled()) {
log.debug("[{}] [{}] Ignoring message acknowledgment during topic transfer, ack count: {}",
subscription, consumerId, ack.getMessageIdsCount());
}
return;
}
consumer.messageAcked(ack).thenRun(() -> {
if (hasRequestId) {
writeAndFlush(Commands.newAckResponse(
Expand Down Expand Up @@ -3071,15 +3081,17 @@ private void closeProducer(long producerId, long epoch, Optional<BrokerLookupDat
}

@Override
public void closeConsumer(Consumer consumer) {
public void closeConsumer(Consumer consumer, Optional<BrokerLookupData> assignedBrokerLookupData) {
// removes consumer-connection from map and send close command to consumer
safelyRemoveConsumer(consumer);
closeConsumer(consumer.consumerId());
closeConsumer(consumer.consumerId(), assignedBrokerLookupData);
}

private void closeConsumer(long consumerId) {
private void closeConsumer(long consumerId, Optional<BrokerLookupData> assignedBrokerLookupData) {
if (getRemoteEndpointProtocolVersion() >= v5.getValue()) {
writeAndFlush(Commands.newCloseConsumer(consumerId, -1L));
writeAndFlush(newCloseConsumer(consumerId, -1L,
assignedBrokerLookupData.map(BrokerLookupData::pulsarServiceUrl).orElse(null),
assignedBrokerLookupData.map(BrokerLookupData::pulsarServiceUrlTls).orElse(null)));
} else {
close();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.apache.bookkeeper.mledger.Position;
import org.apache.bookkeeper.mledger.impl.PositionImpl;
import org.apache.pulsar.broker.intercept.BrokerInterceptor;
import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData;
import org.apache.pulsar.common.api.proto.CommandAck.AckType;
import org.apache.pulsar.common.api.proto.CommandSubscribe.SubType;
import org.apache.pulsar.common.api.proto.ReplicatedSubscriptionsSnapshot;
Expand Down Expand Up @@ -64,13 +65,13 @@ default long getNumberOfEntriesDelayed() {

List<Consumer> getConsumers();

CompletableFuture<Void> close();
Comment thread
gaoran10 marked this conversation as resolved.

CompletableFuture<Void> delete();

CompletableFuture<Void> deleteForcefully();

CompletableFuture<Void> disconnect();
CompletableFuture<Void> disconnect(Optional<BrokerLookupData> assignedBrokerLookupData);

CompletableFuture<Void> close(boolean disconnectConsumers, Optional<BrokerLookupData> assignedBrokerLookupData);

CompletableFuture<Void> doUnsubscribe(Consumer consumer);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ public interface TransportCnx {

void removedConsumer(Consumer consumer);

void closeConsumer(Consumer consumer);
void closeConsumer(Consumer consumer, Optional<BrokerLookupData> assignedBrokerLookupData);

boolean isPreciseDispatcherFlowControl();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,12 @@
package org.apache.pulsar.broker.service.nonpersistent;

import java.util.List;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
import org.apache.bookkeeper.mledger.Entry;
import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData;
import org.apache.pulsar.broker.service.AbstractDispatcherMultipleConsumers;
import org.apache.pulsar.broker.service.BrokerServiceException;
import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException;
Expand Down Expand Up @@ -126,9 +128,11 @@ public synchronized boolean canUnsubscribe(Consumer consumer) {
}

@Override
public CompletableFuture<Void> close() {
public CompletableFuture<Void> close(boolean disconnectConsumers,
Optional<BrokerLookupData> assignedBrokerLookupData) {
IS_CLOSED_UPDATER.set(this, TRUE);
return disconnectAllConsumers();
return disconnectConsumers
? disconnectAllConsumers(false, assignedBrokerLookupData) : CompletableFuture.completedFuture(null);
}

@Override
Expand All @@ -147,12 +151,13 @@ public synchronized void consumerFlow(Consumer consumer, int additionalNumberOfM
}

@Override
public synchronized CompletableFuture<Void> disconnectAllConsumers(boolean isResetCursor) {
public synchronized CompletableFuture<Void> disconnectAllConsumers(
boolean isResetCursor, Optional<BrokerLookupData> assignedBrokerLookupData) {
closeFuture = new CompletableFuture<>();
if (consumerList.isEmpty()) {
closeFuture.complete(null);
} else {
consumerList.forEach(Consumer::disconnect);
Comment thread
dragosvictor marked this conversation as resolved.
consumerList.forEach(consumer -> consumer.disconnect(isResetCursor, assignedBrokerLookupData));
}
return closeFuture;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -101,11 +101,6 @@ protected void scheduleReadOnActiveConsumer() {
// No-op
}

@Override
protected void readMoreEntries(Consumer consumer) {
// No-op
}

@Override
protected void cancelPendingRead() {
// No-op
Expand Down
Loading