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
1 change: 1 addition & 0 deletions pulsar-client-cpp/.gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ Makefile
cmake_install.cmake
CMakeFiles
CMakeCache.txt
.cmake

pulsar-dist
install_manifest.txt
Expand Down
21 changes: 21 additions & 0 deletions pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
Original file line number Diff line number Diff line change
Expand Up @@ -499,6 +499,27 @@ class PULSAR_PUBLIC ConsumerConfiguration {
*/
bool isAutoAckOldestChunkedMessageOnQueueFull() const;

/**
* If this is enabled, consumer receiver queue size is init as a very small value, 1 by default,
* and it will double itself until it reaches the value set by {@link #receiverQueueSize(int)}, if and only if
* 1) User calls receive() and there is no messages in receiver queue.
* 2) The last message we put in the receiver queue took the last space available in receiver queue.
*
* This is disabled by default and currentReceiverQueueSize is init as maxReceiverQueueSize.
*
* The feature should be able to reduce client memory usage.
*
* Default: false
*
* @param enabled whether to enable AutoScaledReceiverQueueSize.
*/
ConsumerConfiguration& setAutoScaledReceiverQueueSizeEnabled(bool enabled);

/**
* The associated getter of autoScaledReceiverQueueSizeEnabled
*/
bool isAutoScaledReceiverQueueSizeEnabled() const;

friend class PulsarWrapper;

private:
Expand Down
7 changes: 6 additions & 1 deletion pulsar-client-cpp/lib/ClientImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,12 @@ ClientImpl::ClientImpl(const std::string& serviceUrl, const ClientConfiguration&
state_(Open),
serviceUrl_(serviceUrl),
clientConfiguration_(detectTls(serviceUrl, clientConfiguration)),
memoryLimitController_(clientConfiguration.getMemoryLimit()),
memoryLimitController_(
clientConfiguration_.getMemoryLimit(),
clientConfiguration_.getMemoryLimit() * MEMORY_THRESHOLD_FOR_RECEIVER_QUEUE_SIZE_EXPANSION5,
[this]() {
std::for_each(consumers_.begin(), consumers_.end(), [](ConsumerImplBaseWeakPtr consumer) { consumer.lock()->reduceCurrentReceiverQueueSize(); });
}),
ioExecutorProvider_(std::make_shared<ExecutorServiceProvider>(clientConfiguration_.getIOThreads())),
listenerExecutorProvider_(
std::make_shared<ExecutorServiceProvider>(clientConfiguration_.getMessageListenerThreads())),
Expand Down
8 changes: 8 additions & 0 deletions pulsar-client-cpp/lib/ConsumerConfiguration.cc
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
#include <lib/ConsumerConfigurationImpl.h>

#include <stdexcept>
#include "pulsar/ConsumerConfiguration.h"

namespace pulsar {

Expand Down Expand Up @@ -259,5 +260,12 @@ ConsumerConfiguration& ConsumerConfiguration::setAutoAckOldestChunkedMessageOnQu
bool ConsumerConfiguration::isAutoAckOldestChunkedMessageOnQueueFull() const {
return impl_->autoAckOldestChunkedMessageOnQueueFull;
}
ConsumerConfiguration& ConsumerConfiguration::setAutoScaledReceiverQueueSizeEnabled(bool enabled) {
impl_->autoScaledReceiverQueueSizeEnabled = enabled;
return *this;
}
bool ConsumerConfiguration::isAutoScaledReceiverQueueSizeEnabled() const {
return impl_->autoScaledReceiverQueueSizeEnabled;
}

} // namespace pulsar
1 change: 1 addition & 0 deletions pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ struct ConsumerConfigurationImpl {
KeySharedPolicy keySharedPolicy;
size_t maxPendingChunkedMessage{10};
bool autoAckOldestChunkedMessageOnQueueFull{false};
bool autoScaledReceiverQueueSizeEnabled{false};
};
} // namespace pulsar
#endif /* LIB_CONSUMERCONFIGURATIONIMPL_H_ */
58 changes: 56 additions & 2 deletions pulsar-client-cpp/lib/ConsumerImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -31,11 +31,14 @@
#include "AckGroupingTrackerDisabled.h"
#include <exception>
#include <algorithm>
#include <math.h>

namespace pulsar {

DECLARE_LOG_OBJECT()

const static int INITIAL_RECEIVER_QUEUE_SIZE = 1;

ConsumerImpl::ConsumerImpl(const ClientImplPtr client, const std::string& topic,
const std::string& subscriptionName, const ConsumerConfiguration& conf,
const ExecutorServicePtr listenerExecutor /* = NULL by default */,
Expand All @@ -55,6 +58,7 @@ ConsumerImpl::ConsumerImpl(const ClientImplPtr client, const std::string& topic,
// This is the initial capacity of the queue
incomingMessages_(std::max(config_.getReceiverQueueSize(), 1)),
availablePermits_(0),
scaleReceiverQueueHint(false),
receiverQueueRefillThreshold_(config_.getReceiverQueueSize() / 2),
consumerId_(client->newConsumerId()),
consumerName_(config_.getConsumerName()),
Expand Down Expand Up @@ -103,6 +107,7 @@ ConsumerImpl::ConsumerImpl(const ClientImplPtr client, const std::string& topic,
if (conf.isEncryptionEnabled()) {
msgCrypto_ = std::make_shared<MessageCrypto>(consumerStr_, false);
}
initReceiverQueueSize();
}

ConsumerImpl::~ConsumerImpl() {
Expand Down Expand Up @@ -880,10 +885,12 @@ Optional<MessageId> ConsumerImpl::clearReceiveQueue() {
void ConsumerImpl::increaseAvailablePermits(const ClientConnectionPtr& currentCnx, int delta) {
int newAvailablePermits = availablePermits_.fetch_add(delta) + delta;

while (newAvailablePermits >= receiverQueueRefillThreshold_ && messageListenerRunning_) {
while (newAvailablePermits >= getCurrentReceiverQueueSize() / 2 && messageListenerRunning_) {
if (availablePermits_.compare_exchange_weak(newAvailablePermits, 0)) {
sendFlowPermitsToBroker(currentCnx, newAvailablePermits);
break;
} else {
newAvailablePermits = availablePermits_;
}
}
}
Expand Down Expand Up @@ -1385,4 +1392,51 @@ bool ConsumerImpl::isConnected() const {

uint64_t ConsumerImpl::getNumberOfConnectedConsumer() { return isConnected() ? 1 : 0; }

} /* namespace pulsar */
MemoryLimitController& ConsumerImpl::getMemoryLimitController() {
return client_.lock()->getMemoryLimitController();
}

void ConsumerImpl::reduceCurrentReceiverQueueSize() {
if (!config_.isAutoScaledReceiverQueueSizeEnabled()) {
return ;
}
int oldSize = getCurrentReceiverQueueSize();
int newSize = std::max(minReceiverQueueSize(), oldSize / 2);
if (oldSize > newSize) {
setCurrentReceiverQueueSize(newSize);
}
}

void ConsumerImpl::expectMoreIncomingMessages() {
if (!config_.isAutoScaledReceiverQueueSizeEnabled()) {
return ;
}
double usage = getMemoryLimitController().currentUsagePercent();
if (bool expectedState = true && usage < MEMORY_THRESHOLD_FOR_RECEIVER_QUEUE_SIZE_EXPANSION5
&& scaleReceiverQueueHint.compare_exchange_strong(expectedState, false)) {
int oldSize = getCurrentReceiverQueueSize();
int newSize = std::min(config_.getReceiverQueueSize(), oldSize * 2);
setCurrentReceiverQueueSize(newSize);
}
}
void ConsumerImpl::initReceiverQueueSize() {
if (config_.isAutoScaledReceiverQueueSizeEnabled()) {
int size = minReceiverQueueSize();
currentReceiverQueueSize_.exchange(size);
} else {
currentReceiverQueueSize_.exchange(config_.getReceiverQueueSize());
}
}

void ConsumerImpl::setCurrentReceiverQueueSize(int newSize) {
currentReceiverQueueSize_.fetch_xor(newSize);
}

int ConsumerImpl::getCurrentReceiverQueueSize() {
return currentReceiverQueueSize_;
}
int ConsumerImpl::minReceiverQueueSize() {
int size = std::min(INITIAL_RECEIVER_QUEUE_SIZE, config_.getReceiverQueueSize());
return size;
}
} /* namespace pulsar */
11 changes: 11 additions & 0 deletions pulsar-client-cpp/lib/ConsumerImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,7 @@ class ConsumerImpl : public ConsumerImplBase,
void negativeAcknowledge(const MessageId& msgId) override;
bool isConnected() const override;
uint64_t getNumberOfConnectedConsumer() override;
void reduceCurrentReceiverQueueSize() override;

virtual void disconnectConsumer();
Result fetchSingleMessageFromBroker(Message& msg);
Expand All @@ -141,6 +142,7 @@ class ConsumerImpl : public ConsumerImplBase,
virtual bool isReadCompacted();
virtual void hasMessageAvailableAsync(HasMessageAvailableCallback callback);
virtual void getLastMessageIdAsync(BrokerGetLastMessageIdCallback callback);
int getCurrentReceiverQueueSize();

protected:
// overrided methods from HandlerBase
Expand All @@ -157,6 +159,11 @@ class ConsumerImpl : public ConsumerImplBase,
void handleClose(Result result, ResultCallback callback, ConsumerImplPtr consumer);
ConsumerStatsBasePtr consumerStatsBasePtr_;

void setCurrentReceiverQueueSize(int newSize);
void expectMoreIncomingMessages();
void initReceiverQueueSize();
int minReceiverQueueSize();

private:
bool waitingForZeroQueueSizeMessage;
bool uncompressMessageIfNeeded(const ClientConnectionPtr& cnx, const proto::MessageIdData& messageIdData,
Expand All @@ -173,6 +180,8 @@ class ConsumerImpl : public ConsumerImplBase,
bool decryptMessageIfNeeded(const ClientConnectionPtr& cnx, const proto::CommandMessage& msg,
const proto::MessageMetadata& metadata, SharedBuffer& payload);

MemoryLimitController& getMemoryLimitController();

// TODO - Convert these functions to lambda when we move to C++11
Result receiveHelper(Message& msg);
Result receiveHelper(Message& msg, int timeout);
Expand All @@ -199,6 +208,8 @@ class ConsumerImpl : public ConsumerImplBase,
UnboundedBlockingQueue<Message> incomingMessages_;
std::queue<ReceiveCallback> pendingReceives_;
std::atomic_int availablePermits_;
std::atomic_int currentReceiverQueueSize_;
std::atomic_bool scaleReceiverQueueHint;
const int receiverQueueRefillThreshold_;
uint64_t consumerId_;
std::string consumerName_;
Expand Down
1 change: 1 addition & 0 deletions pulsar-client-cpp/lib/ConsumerImplBase.h
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ class ConsumerImplBase {
virtual void negativeAcknowledge(const MessageId& msgId) = 0;
virtual bool isConnected() const = 0;
virtual uint64_t getNumberOfConnectedConsumer() = 0;
virtual void reduceCurrentReceiverQueueSize() = 0;

private:
virtual void setNegativeAcknowledgeEnabledForTesting(bool enabled) = 0;
Expand Down
30 changes: 29 additions & 1 deletion pulsar-client-cpp/lib/MemoryLimitController.cc
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,12 @@
namespace pulsar {

MemoryLimitController::MemoryLimitController(uint64_t memoryLimit)
: memoryLimit_(memoryLimit), currentUsage_(0), mutex_(), condition_() {}
: memoryLimit_(memoryLimit), currentUsage_(0), mutex_(), condition_(),
triggerThreshold_(0), trigger_(), triggerRunning(false) {}

MemoryLimitController::MemoryLimitController(uint64_t memoryLimit, uint64_t triggerThreshold, Trigger trigger)
: memoryLimit_(memoryLimit), currentUsage_(0), mutex_(), condition_(), triggerThreshold_(triggerThreshold),
trigger_(trigger), triggerRunning(false) {}

bool MemoryLimitController::tryReserveMemory(uint64_t size) {
// Avoid CAS operation when size is 0
Expand All @@ -40,11 +45,30 @@ bool MemoryLimitController::tryReserveMemory(uint64_t size) {
}

if (currentUsage_.compare_exchange_strong(current, newUsage)) {
checkTrigger(current, newUsage);
return true;
}
}
}

void MemoryLimitController::forceReserveMemory(uint64_t size) {
uint64_t newUsage = currentUsage_.fetch_add(size);
checkTrigger(newUsage - size, newUsage);
}

void MemoryLimitController::checkTrigger(uint64_t preUsage, uint64_t newUsage) {
if (newUsage >= triggerThreshold_ && preUsage < triggerThreshold_ && trigger_) {
bool expectedState = false;
if (triggerRunning.compare_exchange_strong(expectedState, true)) {
try {
trigger_();
} catch (const std::exception exception) {
}
triggerRunning.exchange(false);
}
}
}

bool MemoryLimitController::reserveMemory(uint64_t size) {
if (!tryReserveMemory(size)) {
std::unique_lock<std::mutex> lock(mutex_);
Expand Down Expand Up @@ -83,4 +107,8 @@ void MemoryLimitController::close() {
condition_.notify_all();
}

double MemoryLimitController::currentUsagePercent() const {
return 1.0 * currentUsage_ / memoryLimit_;
}

} // namespace pulsar
10 changes: 10 additions & 0 deletions pulsar-client-cpp/lib/MemoryLimitController.h
Original file line number Diff line number Diff line change
Expand Up @@ -25,14 +25,20 @@
#include <mutex>

namespace pulsar {
typedef std::function<void(void)> Trigger;

const static double MEMORY_THRESHOLD_FOR_RECEIVER_QUEUE_SIZE_EXPANSION5 = 0.75;

class MemoryLimitController {
public:
explicit MemoryLimitController(uint64_t memoryLimit);
MemoryLimitController(uint64_t memoryLimit, uint64_t triggerThreshold, Trigger trigger);
void forceReserveMemory(uint64_t size);
bool tryReserveMemory(uint64_t size);
bool reserveMemory(uint64_t size);
void releaseMemory(uint64_t size);
uint64_t currentUsage() const;
double currentUsagePercent() const;

void close();

Expand All @@ -42,6 +48,10 @@ class MemoryLimitController {
std::mutex mutex_;
std::condition_variable condition_;
bool isClosed_ = false;
const uint64_t triggerThreshold_;
const Trigger trigger_;
std::atomic_bool triggerRunning;
void checkTrigger(uint64_t preUsage, uint64_t newUsage);
};

} // namespace pulsar
1 change: 1 addition & 0 deletions pulsar-client-cpp/lib/MultiTopicsConsumerImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -738,3 +738,4 @@ uint64_t MultiTopicsConsumerImpl::getNumberOfConnectedConsumer() {
});
return numberOfConnectedConsumer;
}
void MultiTopicsConsumerImpl::reduceCurrentReceiverQueueSize() {}
1 change: 1 addition & 0 deletions pulsar-client-cpp/lib/MultiTopicsConsumerImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,7 @@ class MultiTopicsConsumerImpl : public ConsumerImplBase,
void negativeAcknowledge(const MessageId& msgId) override;
bool isConnected() const override;
uint64_t getNumberOfConnectedConsumer() override;
void reduceCurrentReceiverQueueSize() override;

void handleGetConsumerStats(Result, BrokerConsumerStats, LatchPtr, MultiTopicsBrokerConsumerStatsPtr,
size_t, BrokerConsumerStatsCallback);
Expand Down
1 change: 1 addition & 0 deletions pulsar-client-cpp/lib/PartitionedConsumerImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -642,5 +642,6 @@ uint64_t PartitionedConsumerImpl::getNumberOfConnectedConsumer() {
}
return numberOfConnectedConsumer;
}
void PartitionedConsumerImpl::reduceCurrentReceiverQueueSize() {}

} // namespace pulsar
1 change: 1 addition & 0 deletions pulsar-client-cpp/lib/PartitionedConsumerImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ class PartitionedConsumerImpl : public ConsumerImplBase,
void negativeAcknowledge(const MessageId& msgId) override;
bool isConnected() const override;
uint64_t getNumberOfConnectedConsumer() override;
void reduceCurrentReceiverQueueSize() override;

void handleGetConsumerStats(Result, BrokerConsumerStats, LatchPtr, PartitionedBrokerConsumerStatsPtr,
size_t, BrokerConsumerStatsCallback);
Expand Down
25 changes: 25 additions & 0 deletions pulsar-client-cpp/tests/MemoryLimitControllerTest.cc
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,31 @@ TEST(MemoryLimitControllerTest, testLimit) {
ASSERT_EQ(mlc.currentUsage(), 101);
}

TEST(MemoryLimitControllerTest, testTrigger) {
int num = 0;
MemoryLimitController mlc(100, 95, [&num]() {
num++;
});

mlc.tryReserveMemory(95);
ASSERT_EQ(num , 1);

mlc.releaseMemory(95);
ASSERT_EQ(mlc.currentUsage(), 0);

std::thread t1([&]() {
mlc.forceReserveMemory(95);
});

std::thread t2([&]() {
mlc.forceReserveMemory(95);
});

t1.join();
t2.join();
ASSERT_EQ(num, 2);
}

TEST(MemoryLimitControllerTest, testBlocking) {
MemoryLimitController mlc(100);

Expand Down