From b388a2f1f4d72863af321c813309eae736bec830 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Tue, 13 Dec 2022 19:01:09 +0800 Subject: [PATCH 1/4] [cherry-pick][branch-2.10] Fix issue where unexpected ack timeout occurred (#17503) --- .../lib/MultiTopicsConsumerImpl.cc | 21 ++++-- .../lib/MultiTopicsConsumerImpl.h | 1 + pulsar-client-cpp/tests/ConsumerTest.cc | 71 +++++++++++++++++++ 3 files changed, 87 insertions(+), 6 deletions(-) diff --git a/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc b/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc index c35d84b6db668..5a5cf55627bda 100644 --- a/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc +++ b/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc @@ -451,14 +451,14 @@ void MultiTopicsConsumerImpl::messageReceived(Consumer consumer, const Message& ReceiveCallback callback = pendingReceives_.front(); pendingReceives_.pop(); lock.unlock(); - unAckedMessageTrackerPtr_->add(msg.getMessageId()); - listenerExecutor_->postWork(std::bind(callback, ResultOk, msg)); + listenerExecutor_->postWork(std::bind(&MultiTopicsConsumerImpl::notifyPendingReceivedCallback, + shared_from_this(), ResultOk, msg, callback)); } else { if (messages_.full()) { lock.unlock(); } - messages_.push(msg); - if (messageListener_) { + + if (messages_.push(msg) && messageListener_) { unAckedMessageTrackerPtr_->add(msg.getMessageId()); listenerExecutor_->postWork( std::bind(&MultiTopicsConsumerImpl::internalListener, shared_from_this(), consumer)); @@ -469,7 +469,7 @@ void MultiTopicsConsumerImpl::messageReceived(Consumer consumer, const Message& void MultiTopicsConsumerImpl::internalListener(Consumer consumer) { Message m; messages_.pop(m); - + unAckedMessageTrackerPtr_->add(m.getMessageId()); try { messageListener_(Consumer(shared_from_this()), m); } catch (const std::exception& e) { @@ -535,11 +535,20 @@ void MultiTopicsConsumerImpl::failPendingReceiveCallback() { while (!pendingReceives_.empty()) { ReceiveCallback callback = pendingReceives_.front(); pendingReceives_.pop(); - listenerExecutor_->postWork(std::bind(callback, ResultAlreadyClosed, msg)); + listenerExecutor_->postWork(std::bind(&MultiTopicsConsumerImpl::notifyPendingReceivedCallback, + shared_from_this(), ResultAlreadyClosed, msg, callback)); } lock.unlock(); } +void MultiTopicsConsumerImpl::notifyPendingReceivedCallback(Result result, Message& msg, + const ReceiveCallback& callback) { + if (result == ResultOk) { + unAckedMessageTrackerPtr_->add(msg.getMessageId()); + } + callback(result, msg); +} + void MultiTopicsConsumerImpl::acknowledgeAsync(const MessageId& msgId, ResultCallback callback) { if (state_ != Ready) { callback(ResultAlreadyClosed); diff --git a/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.h b/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.h index 95c24f68c5b78..8769d59b9908e 100644 --- a/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.h +++ b/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.h @@ -128,6 +128,7 @@ class MultiTopicsConsumerImpl : public ConsumerImplBase, void internalListener(Consumer consumer); void receiveMessages(); void failPendingReceiveCallback(); + void notifyPendingReceivedCallback(Result result, Message& message, const ReceiveCallback& callback); void handleOneTopicSubscribed(Result result, Consumer consumer, const std::string& topic, std::shared_ptr> topicsNeedCreate); diff --git a/pulsar-client-cpp/tests/ConsumerTest.cc b/pulsar-client-cpp/tests/ConsumerTest.cc index b1fc11cec8fe3..e672da33c0990 100644 --- a/pulsar-client-cpp/tests/ConsumerTest.cc +++ b/pulsar-client-cpp/tests/ConsumerTest.cc @@ -475,6 +475,77 @@ TEST(ConsumerTest, testPartitionedConsumerUnAckedMessageRedelivery) { client.close(); } +TEST(ConsumerTest, testPartitionedConsumerUnexpectedAckTimeout) { + ClientConfiguration clientConfig; + clientConfig.setMessageListenerThreads(1); + Client client(lookupUrl, clientConfig); + + const std::string partitionedTopic = + "testPartitionedConsumerUnexpectedAckTimeout" + std::to_string(time(nullptr)); + std::string subName = "sub"; + constexpr int numPartitions = 2; + constexpr int numOfMessages = 3; + constexpr int unAckedMessagesTimeoutMs = 10000; + constexpr int tickDurationInMs = 1000; + pulsar::Latch latch(numOfMessages); + std::vector messages; + std::mutex mtx; + + int res = + makePutRequest(adminUrl + "admin/v2/persistent/public/default/" + partitionedTopic + "/partitions", + std::to_string(numPartitions)); + ASSERT_TRUE(res == 204 || res == 409) << "res: " << res; + + Consumer consumer; + ConsumerConfiguration consumerConfig; + consumerConfig.setConsumerType(ConsumerShared); + consumerConfig.setUnAckedMessagesTimeoutMs(unAckedMessagesTimeoutMs); + consumerConfig.setTickDurationInMs(tickDurationInMs); + consumerConfig.setMessageListener([&](Consumer cons, const Message& msg) { + // acknowledge received messages immediately, so no ack timeout is expected + ASSERT_EQ(ResultOk, cons.acknowledge(msg.getMessageId())); + ASSERT_EQ(0, msg.getRedeliveryCount()); + + { + std::lock_guard lock(mtx); + messages.emplace_back(msg); + } + + if (latch.getCount() > 0) { + std::this_thread::sleep_for( + std::chrono::milliseconds(unAckedMessagesTimeoutMs + tickDurationInMs * 2)); + latch.countdown(); + } + }); + ASSERT_EQ(ResultOk, client.subscribe(partitionedTopic, subName, consumerConfig, consumer)); + + // send messages + ProducerConfiguration producerConfig; + producerConfig.setBatchingEnabled(false); + producerConfig.setBlockIfQueueFull(true); + producerConfig.setPartitionsRoutingMode(ProducerConfiguration::UseSinglePartition); + Producer producer; + ASSERT_EQ(ResultOk, client.createProducer(partitionedTopic, producerConfig, producer)); + std::string prefix = "message-"; + for (int i = 0; i < numOfMessages; i++) { + std::string messageContent = prefix + std::to_string(i); + Message msg = MessageBuilder().setContent(messageContent).build(); + ASSERT_EQ(ResultOk, producer.send(msg)); + } + producer.close(); + + bool wasUnblocked = latch.wait( + std::chrono::milliseconds((unAckedMessagesTimeoutMs + tickDurationInMs * 2) * numOfMessages + 5000)); + ASSERT_TRUE(wasUnblocked); + + std::this_thread::sleep_for(std::chrono::milliseconds(5000)); + // messages are expected not to be redelivered + ASSERT_EQ(numOfMessages, messages.size()); + + consumer.close(); + client.close(); +} + TEST(ConsumerTest, testMultiTopicsConsumerUnAckedMessageRedelivery) { Client client(lookupUrl); const std::string nonPartitionedTopic = From 7eb1b3d43c497f0201e9e30ce644f55be53a8790 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Tue, 13 Dec 2022 19:40:32 +0800 Subject: [PATCH 2/4] unused imported --- .../org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java index f0d9e2bf241ac..b9533c6f29192 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java @@ -27,7 +27,6 @@ import java.nio.ByteBuffer; import java.time.Clock; import java.util.List; -import java.util.concurrent.CompletableFuture; import java.util.stream.Collectors; import javax.ws.rs.container.AsyncResponse; import javax.ws.rs.core.MediaType; From 68ab2621731d1e2bb45848ce0c46d40ed8108c28 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Tue, 13 Dec 2022 21:14:25 +0800 Subject: [PATCH 3/4] fix --- pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc b/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc index 5a5cf55627bda..c78c12eaec0fb 100644 --- a/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc +++ b/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc @@ -457,8 +457,8 @@ void MultiTopicsConsumerImpl::messageReceived(Consumer consumer, const Message& if (messages_.full()) { lock.unlock(); } - - if (messages_.push(msg) && messageListener_) { + messages_.push(msg); + if (messageListener_) { unAckedMessageTrackerPtr_->add(msg.getMessageId()); listenerExecutor_->postWork( std::bind(&MultiTopicsConsumerImpl::internalListener, shared_from_this(), consumer)); From be7b6ede5e0e03ef3694f3f22e790d7e0aec7ce0 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Tue, 13 Dec 2022 21:14:54 +0800 Subject: [PATCH 4/4] fix --- pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc b/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc index c78c12eaec0fb..18723684afe42 100644 --- a/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc +++ b/pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc @@ -459,7 +459,6 @@ void MultiTopicsConsumerImpl::messageReceived(Consumer consumer, const Message& } messages_.push(msg); if (messageListener_) { - unAckedMessageTrackerPtr_->add(msg.getMessageId()); listenerExecutor_->postWork( std::bind(&MultiTopicsConsumerImpl::internalListener, shared_from_this(), consumer)); }