Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
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 @@ -84,6 +84,7 @@
import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.pulsar.PulsarVersion;
import org.apache.pulsar.broker.BrokerTestUtil;
import org.apache.pulsar.broker.PulsarService;
import org.apache.pulsar.broker.service.persistent.PersistentTopic;
import org.apache.pulsar.client.admin.PulsarAdminException;
Expand All @@ -107,7 +108,9 @@
import org.apache.pulsar.common.compression.CompressionCodec;
import org.apache.pulsar.common.compression.CompressionCodecProvider;
import org.apache.pulsar.common.naming.TopicName;
import org.apache.pulsar.common.policies.data.ConsumerStats;
import org.apache.pulsar.common.policies.data.PublisherStats;
import org.apache.pulsar.common.policies.data.SubscriptionStats;
import org.apache.pulsar.common.policies.data.TopicStats;
import org.apache.pulsar.common.protocol.Commands;
import org.apache.pulsar.common.schema.SchemaType;
Expand Down Expand Up @@ -190,6 +193,16 @@ public Object[][] ackReceiptEnabledAndSubscriptionTypes() {
};
}

@DataProvider(name = "subscriptionTypes")
public Object[][] subType() {
return new Object[][] {
{SubscriptionType.Shared},
{SubscriptionType.Key_Shared},
{SubscriptionType.Exclusive},
{SubscriptionType.Failover}
};
}

@AfterClass(alwaysRun = true)
@Override
protected void cleanup() throws Exception {
Expand Down Expand Up @@ -340,6 +353,58 @@ public Object[][] codecProvider() {
return new Object[][] { { 0 }, { 1000 } };
}

@Test(dataProvider = "subscriptionTypes")
public void testConsumerReconnectTwice(SubscriptionType subscriptionType) throws Exception {
final String topicName = BrokerTestUtil.newUniqueName("persistent://my-property/my-ns/tp_");
final String subscriptionName = "subscription1";
admin.topics().createNonPartitionedTopic(topicName);
admin.topics().createSubscription(topicName, subscriptionName, MessageId.earliest);
// Create producer and consumer.
ConsumerImpl<String> consumer = (ConsumerImpl<String>) pulsarClient.newConsumer(Schema.STRING)
.subscriptionType(subscriptionType)
.receiverQueueSize(1000).topic(topicName).subscriptionName(subscriptionName).subscribe();
Producer<String> producer = pulsarClient.newProducer(Schema.STRING).enableBatching(false)
.topic(topicName).create();
int sendMessageCount = 10;
for (int i = 0; i < sendMessageCount; i++){
producer.send("msg- " + i);
}
Awaitility.await().untilAsserted(() -> {
assertEquals(consumer.numMessagesInQueue(), sendMessageCount);
});
printConsumerStats(topicName, subscriptionName);

// Do the second subscribe.
consumer.connectionOpened(consumer.getClientCnx());

// Verify messages are not lost.
List<Message<String>> messages = new ArrayList<>();
while (true) {
Message<String> message = consumer.receive(2, TimeUnit.SECONDS);
if (message == null) {
break;
}
messages.add(message);
}
printConsumerStats(topicName, subscriptionName);
assertEquals(messages.size(), sendMessageCount);

// cleanup.
consumer.close();
producer.close();
admin.topics().delete(topicName, false);
}

private void printConsumerStats(String topicName, String subscriptionName) throws Exception {
SubscriptionStats subscriptionStats =
admin.topics().getStats(topicName).getSubscriptions().get(subscriptionName);
ConsumerStats consumerStats =
admin.topics().getStats(topicName).getSubscriptions().get(subscriptionName).getConsumers().get(0);
log.info("msgBacklog: {}, msgOutCounter: {}, unackedMessages: {}, availablePermits: {}",
subscriptionStats.getMsgBacklog(), consumerStats.getMsgOutCounter(),
consumerStats.getUnackedMessages(), consumerStats.getAvailablePermits());
}

@Test(timeOut = 100000, dataProvider = "batch")
public void testSyncProducerAndConsumer(int batchMessageDelayMs) throws Exception {
log.info("-- Starting {} test --", methodName);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,7 @@ public class ConsumerImpl<T> extends ConsumerBase<T> implements ConnectionHandle

protected volatile MessageId lastDequeuedMessageId = MessageId.earliest;
private volatile MessageId lastMessageIdInBroker = MessageId.earliest;
protected volatile CompletableFuture<Void> inProgressSubscribeFuture;
Comment thread
BewareMyPower marked this conversation as resolved.

private final long lookupDeadline;

Expand Down Expand Up @@ -761,6 +762,18 @@ public void negativeAcknowledge(Message<?> message) {

@Override
public void connectionOpened(final ClientCnx cnx) {
synchronized (this) {
// Wait the previous subscribe done.
if (inProgressSubscribeFuture != null && !inProgressSubscribeFuture.isDone()){
return;
}
// If success.
if (getState() == State.Ready) {
return;
}
// Do subscribe if previous was failed.
inProgressSubscribeFuture = new CompletableFuture<>();
}
previousExceptions.clear();

if (getState() == State.Closing || getState() == State.Closed) {
Expand Down Expand Up @@ -906,6 +919,12 @@ public void connectionOpened(final ClientCnx cnx) {
reconnectLater(e.getCause());
}
return null;
}).whenComplete((ignore, ex) -> {
if (ex == null) {
inProgressSubscribeFuture.complete(null);
} else {
inProgressSubscribeFuture.completeExceptionally(ex);
}
});
}
}
Expand Down