Skip to content
Merged
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
10 changes: 5 additions & 5 deletions pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
Original file line number Diff line number Diff line change
Expand Up @@ -450,11 +450,11 @@ class PULSAR_PUBLIC ConsumerConfiguration {
* Buffering large number of outstanding uncompleted chunked messages can create memory pressure and it
* can be guarded by providing this maxPendingChunkedMessage threshold. Once, consumer reaches this
* threshold, it drops the outstanding unchunked-messages by silently acking or asking broker to redeliver
* later by marking it unacked. See setAutoOldestChunkedMessageOnQueueFull.
* later by marking it unacked. See setAutoAckOldestChunkedMessageOnQueueFull.
*
* If it's zero, the pending chunked messages will not be limited.
*
* Default: 100
* Default: 10
*
* @param maxPendingChunkedMessage the number of max pending chunked messages
*/
Expand All @@ -475,13 +475,13 @@ class PULSAR_PUBLIC ConsumerConfiguration {
*
* @param autoAckOldestChunkedMessageOnQueueFull whether to ack the discarded chunked message
*/
ConsumerConfiguration& setAutoOldestChunkedMessageOnQueueFull(
ConsumerConfiguration& setAutoAckOldestChunkedMessageOnQueueFull(
bool autoAckOldestChunkedMessageOnQueueFull);

/**
* The associated getter of setAutoOldestChunkedMessageOnQueueFull
* The associated getter of setAutoAckOldestChunkedMessageOnQueueFull
*/
bool isAutoOldestChunkedMessageOnQueueFull() const;
bool isAutoAckOldestChunkedMessageOnQueueFull() const;

friend class PulsarWrapper;

Expand Down
4 changes: 2 additions & 2 deletions pulsar-client-cpp/lib/ConsumerConfiguration.cc
Original file line number Diff line number Diff line change
Expand Up @@ -238,13 +238,13 @@ ConsumerConfiguration& ConsumerConfiguration::setMaxPendingChunkedMessage(size_t

size_t ConsumerConfiguration::getMaxPendingChunkedMessage() const { return impl_->maxPendingChunkedMessage; }

ConsumerConfiguration& ConsumerConfiguration::setAutoOldestChunkedMessageOnQueueFull(
ConsumerConfiguration& ConsumerConfiguration::setAutoAckOldestChunkedMessageOnQueueFull(
bool autoAckOldestChunkedMessageOnQueueFull) {
impl_->autoAckOldestChunkedMessageOnQueueFull = autoAckOldestChunkedMessageOnQueueFull;
return *this;
}

bool ConsumerConfiguration::isAutoOldestChunkedMessageOnQueueFull() const {
bool ConsumerConfiguration::isAutoAckOldestChunkedMessageOnQueueFull() const {
return impl_->autoAckOldestChunkedMessageOnQueueFull;
}

Expand Down
2 changes: 1 addition & 1 deletion pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ struct ConsumerConfigurationImpl {
std::map<std::string, std::string> properties;
int priorityLevel{0};
KeySharedPolicy keySharedPolicy;
size_t maxPendingChunkedMessage{100};
size_t maxPendingChunkedMessage{10};
bool autoAckOldestChunkedMessageOnQueueFull{false};
};
} // namespace pulsar
Expand Down
2 changes: 1 addition & 1 deletion pulsar-client-cpp/lib/ConsumerImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ ConsumerImpl::ConsumerImpl(const ClientImplPtr client, const std::string& topic,
readCompacted_(conf.isReadCompacted()),
startMessageId_(startMessageId),
maxPendingChunkedMessage_(conf.getMaxPendingChunkedMessage()),
autoAckOldestChunkedMessageOnQueueFull_(conf.isAutoOldestChunkedMessageOnQueueFull()) {
autoAckOldestChunkedMessageOnQueueFull_(conf.isAutoAckOldestChunkedMessageOnQueueFull()) {
std::stringstream consumerStrStream;
consumerStrStream << "[" << topic_ << ", " << subscription_ << ", " << consumerId_ << "] ";
consumerStr_ = consumerStrStream.str();
Expand Down
8 changes: 4 additions & 4 deletions pulsar-client-cpp/tests/ConsumerConfigurationTest.cc
Original file line number Diff line number Diff line change
Expand Up @@ -59,8 +59,8 @@ TEST(ConsumerConfigurationTest, testDefaultConfig) {
ASSERT_EQ(conf.isReplicateSubscriptionStateEnabled(), false);
ASSERT_EQ(conf.getProperties().empty(), true);
ASSERT_EQ(conf.getPriorityLevel(), 0);
ASSERT_EQ(conf.getMaxPendingChunkedMessage(), 100);
ASSERT_EQ(conf.isAutoOldestChunkedMessageOnQueueFull(), false);
ASSERT_EQ(conf.getMaxPendingChunkedMessage(), 10);
ASSERT_EQ(conf.isAutoAckOldestChunkedMessageOnQueueFull(), false);
}

TEST(ConsumerConfigurationTest, testCustomConfig) {
Expand Down Expand Up @@ -145,8 +145,8 @@ TEST(ConsumerConfigurationTest, testCustomConfig) {
conf.setMaxPendingChunkedMessage(500);
ASSERT_EQ(conf.getMaxPendingChunkedMessage(), 500);

conf.setAutoOldestChunkedMessageOnQueueFull(true);
ASSERT_TRUE(conf.isAutoOldestChunkedMessageOnQueueFull());
conf.setAutoAckOldestChunkedMessageOnQueueFull(true);
ASSERT_TRUE(conf.isAutoAckOldestChunkedMessageOnQueueFull());
}

TEST(ConsumerConfigurationTest, testReadCompactPersistentExclusive) {
Expand Down