This is an automated email from the ASF dual-hosted git repository.
mmerli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new 8ab6c34 Cpp client: add readCompacted in consumer config (#1858)
8ab6c34 is described below
commit 8ab6c34541eac5ee14e38764c2d4e3db13000da5
Author: Jia Zhai <[email protected]>
AuthorDate: Tue Jun 12 01:18:23 2018 +0800
Cpp client: add readCompacted in consumer config (#1858)
* add readCompacted in consumer config
* change following comments
* change following comments
* add config in c wrapper
* add e2e test in pulsar_test
* add readerconfig in c wrapper
* change following @ivan's comments
* fix rebase error
---
.../include/pulsar/ConsumerConfiguration.h | 3 +
.../include/pulsar/ReaderConfiguration.h | 3 +
.../include/pulsar/c/consumer_configuration.h | 5 +
.../include/pulsar/c/reader_configuration.h | 5 +
pulsar-client-cpp/lib/ClientImpl.cc | 6 ++
pulsar-client-cpp/lib/Commands.cc | 3 +-
pulsar-client-cpp/lib/Commands.h | 3 +-
pulsar-client-cpp/lib/ConsumerConfiguration.cc | 4 +
pulsar-client-cpp/lib/ConsumerConfigurationImpl.h | 4 +-
pulsar-client-cpp/lib/ConsumerImpl.cc | 10 +-
pulsar-client-cpp/lib/ConsumerImpl.h | 2 +
pulsar-client-cpp/lib/ReaderConfiguration.cc | 4 +
pulsar-client-cpp/lib/ReaderConfigurationImpl.h | 7 +-
pulsar-client-cpp/lib/ReaderImpl.cc | 1 +
pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc | 9 ++
pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc | 9 ++
pulsar-client-cpp/python/pulsar/__init__.py | 5 +-
pulsar-client-cpp/python/pulsar_test.py | 55 +++++++++++
pulsar-client-cpp/python/src/config.cc | 4 +
.../tests/ConsumerConfigurationTest.cc | 105 +++++++++++++++++++++
20 files changed, 239 insertions(+), 8 deletions(-)
diff --git a/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
b/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
index 274abd5..c9584e3 100644
--- a/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
+++ b/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
@@ -149,6 +149,9 @@ class ConsumerConfiguration {
ConsumerCryptoFailureAction getCryptoFailureAction() const;
ConsumerConfiguration& setCryptoFailureAction(ConsumerCryptoFailureAction
action);
+ bool isReadCompacted() const;
+ void setReadCompacted(bool compacted);
+
friend class PulsarWrapper;
private:
diff --git a/pulsar-client-cpp/include/pulsar/ReaderConfiguration.h
b/pulsar-client-cpp/include/pulsar/ReaderConfiguration.h
index 3102cf5..8d365ab 100644
--- a/pulsar-client-cpp/include/pulsar/ReaderConfiguration.h
+++ b/pulsar-client-cpp/include/pulsar/ReaderConfiguration.h
@@ -86,6 +86,9 @@ class ReaderConfiguration {
void setSubscriptionRolePrefix(const std::string& subscriptionRolePrefix);
const std::string& getSubscriptionRolePrefix() const;
+ void setReadCompacted(bool compacted);
+ bool isReadCompacted() const;
+
private:
boost::shared_ptr<ReaderConfigurationImpl> impl_;
};
diff --git a/pulsar-client-cpp/include/pulsar/c/consumer_configuration.h
b/pulsar-client-cpp/include/pulsar/c/consumer_configuration.h
index fcefc97..7299867 100644
--- a/pulsar-client-cpp/include/pulsar/c/consumer_configuration.h
+++ b/pulsar-client-cpp/include/pulsar/c/consumer_configuration.h
@@ -147,6 +147,11 @@ long
pulsar_consumer_get_unacked_messages_timeout_ms(pulsar_consumer_configurati
int pulsar_consumer_is_encryption_enabled(pulsar_consumer_configuration_t
*consumer_configuration);
+int pulsar_consumer_is_read_compacted(pulsar_consumer_configuration_t
*consumer_configuration);
+
+void pulsar_consumer_set_read_compacted(pulsar_consumer_configuration_t
*consumer_configuration,
+ int compacted);
+
// const CryptoKeyReaderPtr getCryptoKeyReader()
//
// const;
diff --git a/pulsar-client-cpp/include/pulsar/c/reader_configuration.h
b/pulsar-client-cpp/include/pulsar/c/reader_configuration.h
index faac37c..327fef1 100644
--- a/pulsar-client-cpp/include/pulsar/c/reader_configuration.h
+++ b/pulsar-client-cpp/include/pulsar/c/reader_configuration.h
@@ -79,6 +79,11 @@ void
pulsar_reader_configuration_set_subscription_role_prefix(pulsar_reader_conf
const char *pulsar_reader_configuration_get_subscription_role_prefix(
pulsar_reader_configuration_t *configuration);
+void
pulsar_reader_configuration_set_read_compacted(pulsar_reader_configuration_t
*configuration,
+ int readCompacted);
+
+int
pulsar_reader_configuration_is_read_compacted(pulsar_reader_configuration_t
*configuration);
+
#pragma GCC visibility pop
#ifdef __cplusplus
diff --git a/pulsar-client-cpp/lib/ClientImpl.cc
b/pulsar-client-cpp/lib/ClientImpl.cc
index 80f947b..d8d07e7 100644
--- a/pulsar-client-cpp/lib/ClientImpl.cc
+++ b/pulsar-client-cpp/lib/ClientImpl.cc
@@ -222,6 +222,12 @@ void ClientImpl::subscribeAsync(const std::string& topic,
const std::string& con
lock.unlock();
callback(ResultInvalidTopicName, Consumer());
return;
+ } else if (conf.isReadCompacted() &&
(topicName->getDomain().compare("persistent") != 0 ||
+ (conf.getConsumerType() !=
ConsumerExclusive &&
+ conf.getConsumerType() !=
ConsumerFailover))) {
+ lock.unlock();
+ callback(ResultInvalidConfiguration, Consumer());
+ return;
}
}
diff --git a/pulsar-client-cpp/lib/Commands.cc
b/pulsar-client-cpp/lib/Commands.cc
index 677fce0..fcf8117 100644
--- a/pulsar-client-cpp/lib/Commands.cc
+++ b/pulsar-client-cpp/lib/Commands.cc
@@ -185,7 +185,7 @@ SharedBuffer Commands::newConnect(const AuthenticationPtr&
authentication, const
SharedBuffer Commands::newSubscribe(const std::string& topic, const
std::string& subscription,
uint64_t consumerId, uint64_t requestId,
CommandSubscribe_SubType subType,
const std::string& consumerName,
SubscriptionMode subscriptionMode,
- Optional<MessageId> startMessageId) {
+ Optional<MessageId> startMessageId, bool
readCompacted) {
BaseCommand cmd;
cmd.set_type(BaseCommand::SUBSCRIBE);
CommandSubscribe* subscribe = cmd.mutable_subscribe();
@@ -196,6 +196,7 @@ SharedBuffer Commands::newSubscribe(const std::string&
topic, const std::string&
subscribe->set_request_id(requestId);
subscribe->set_consumer_name(consumerName);
subscribe->set_durable(subscriptionMode == SubscriptionModeDurable);
+ subscribe->set_read_compacted(readCompacted);
if (startMessageId.is_present()) {
MessageIdData& messageIdData = *subscribe->mutable_start_message_id();
messageIdData.set_ledgerid(startMessageId.value().ledgerId());
diff --git a/pulsar-client-cpp/lib/Commands.h b/pulsar-client-cpp/lib/Commands.h
index 38c498a..5018116 100644
--- a/pulsar-client-cpp/lib/Commands.h
+++ b/pulsar-client-cpp/lib/Commands.h
@@ -77,7 +77,8 @@ class Commands {
static SharedBuffer newSubscribe(const std::string& topic, const
std::string& subscription,
uint64_t consumerId, uint64_t requestId,
proto::CommandSubscribe_SubType subType,
const std::string& consumerName,
- SubscriptionMode subscriptionMode,
Optional<MessageId> startMessageId);
+ SubscriptionMode subscriptionMode,
Optional<MessageId> startMessageId,
+ bool readCompacted);
static SharedBuffer newUnsubscribe(uint64_t consumerId, uint64_t
requestId);
diff --git a/pulsar-client-cpp/lib/ConsumerConfiguration.cc
b/pulsar-client-cpp/lib/ConsumerConfiguration.cc
index 1b99f58..0c145c1 100644
--- a/pulsar-client-cpp/lib/ConsumerConfiguration.cc
+++ b/pulsar-client-cpp/lib/ConsumerConfiguration.cc
@@ -101,4 +101,8 @@ ConsumerConfiguration&
ConsumerConfiguration::setCryptoFailureAction(ConsumerCry
return *this;
}
+bool ConsumerConfiguration::isReadCompacted() const { return
impl_->readCompacted; }
+
+void ConsumerConfiguration::setReadCompacted(bool compacted) {
impl_->readCompacted = compacted; }
+
} // namespace pulsar
diff --git a/pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
b/pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
index 91434a2..eb0c374 100644
--- a/pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
+++ b/pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
@@ -34,6 +34,7 @@ struct ConsumerConfigurationImpl {
long brokerConsumerStatsCacheTimeInMs;
CryptoKeyReaderPtr cryptoKeyReader;
ConsumerCryptoFailureAction cryptoFailureAction;
+ bool readCompacted;
ConsumerConfigurationImpl()
: unAckedMessagesTimeoutMs(0),
consumerType(ConsumerExclusive),
@@ -43,7 +44,8 @@ struct ConsumerConfigurationImpl {
receiverQueueSize(1000),
maxTotalReceiverQueueSizeAcrossPartitions(50000),
cryptoKeyReader(),
- cryptoFailureAction(ConsumerCryptoFailureAction::FAIL) {}
+ cryptoFailureAction(ConsumerCryptoFailureAction::FAIL),
+ readCompacted(false) {}
};
} // namespace pulsar
#endif /* LIB_CONSUMERCONFIGURATIONIMPL_H_ */
diff --git a/pulsar-client-cpp/lib/ConsumerImpl.cc
b/pulsar-client-cpp/lib/ConsumerImpl.cc
index 0425238..a50394f 100644
--- a/pulsar-client-cpp/lib/ConsumerImpl.cc
+++ b/pulsar-client-cpp/lib/ConsumerImpl.cc
@@ -57,7 +57,8 @@ ConsumerImpl::ConsumerImpl(const ClientImplPtr client, const
std::string& topic,
batchAcknowledgementTracker_(topic_, subscription, (long)consumerId_),
brokerConsumerStats_(),
consumerStatsBasePtr_(),
- msgCrypto_() {
+ msgCrypto_(),
+ readCompacted_(conf.isReadCompacted()) {
std::stringstream consumerStrStream;
consumerStrStream << "[" << topic_ << ", " << subscription_ << ", " <<
consumerId_ << "] ";
consumerStr_ = consumerStrStream.str();
@@ -135,8 +136,9 @@ void ConsumerImpl::connectionOpened(const
ClientConnectionPtr& cnx) {
ClientImplPtr client = client_.lock();
uint64_t requestId = client->newRequestId();
- SharedBuffer cmd = Commands::newSubscribe(topic_, subscription_,
consumerId_, requestId, getSubType(),
- consumerName_,
subscriptionMode_, startMessageId_);
+ SharedBuffer cmd =
+ Commands::newSubscribe(topic_, subscription_, consumerId_, requestId,
getSubType(), consumerName_,
+ subscriptionMode_, startMessageId_,
readCompacted_);
cnx->sendRequestWithId(cmd, requestId)
.addListener(boost::bind(&ConsumerImpl::handleCreateConsumer,
shared_from_this(), cnx, _1));
}
@@ -919,4 +921,6 @@ void ConsumerImpl::seekAsync(const MessageId& msgId,
ResultCallback callback) {
callback(ResultNotConnected);
}
+bool ConsumerImpl::isReadCompacted() { return readCompacted_; }
+
} /* namespace pulsar */
diff --git a/pulsar-client-cpp/lib/ConsumerImpl.h
b/pulsar-client-cpp/lib/ConsumerImpl.h
index e35f444..d16c8b35 100644
--- a/pulsar-client-cpp/lib/ConsumerImpl.h
+++ b/pulsar-client-cpp/lib/ConsumerImpl.h
@@ -104,6 +104,7 @@ class ConsumerImpl : public ConsumerImplBase,
virtual void getBrokerConsumerStatsAsync(BrokerConsumerStatsCallback
callback);
void handleSeek(Result result, ResultCallback callback);
virtual void seekAsync(const MessageId& msgId, ResultCallback callback);
+ virtual bool isReadCompacted();
protected:
void connectionOpened(const ClientConnectionPtr& cnx);
@@ -165,6 +166,7 @@ class ConsumerImpl : public ConsumerImplBase,
BrokerConsumerStatsImpl brokerConsumerStats_;
MessageCryptoPtr msgCrypto_;
+ const bool readCompacted_;
friend class PulsarFriend;
};
diff --git a/pulsar-client-cpp/lib/ReaderConfiguration.cc
b/pulsar-client-cpp/lib/ReaderConfiguration.cc
index e44714f..a649b33 100644
--- a/pulsar-client-cpp/lib/ReaderConfiguration.cc
+++ b/pulsar-client-cpp/lib/ReaderConfiguration.cc
@@ -56,4 +56,8 @@ const std::string&
ReaderConfiguration::getSubscriptionRolePrefix() const {
void ReaderConfiguration::setSubscriptionRolePrefix(const std::string&
subscriptionRolePrefix) {
impl_->subscriptionRolePrefix = subscriptionRolePrefix;
}
+
+bool ReaderConfiguration::isReadCompacted() const { return
impl_->readCompacted; }
+
+void ReaderConfiguration::setReadCompacted(bool compacted) {
impl_->readCompacted = compacted; }
} // namespace pulsar
diff --git a/pulsar-client-cpp/lib/ReaderConfigurationImpl.h
b/pulsar-client-cpp/lib/ReaderConfigurationImpl.h
index 67ddb3f..5dca8e3 100644
--- a/pulsar-client-cpp/lib/ReaderConfigurationImpl.h
+++ b/pulsar-client-cpp/lib/ReaderConfigurationImpl.h
@@ -29,8 +29,13 @@ struct ReaderConfigurationImpl {
int receiverQueueSize;
std::string readerName;
std::string subscriptionRolePrefix;
+ bool readCompacted;
ReaderConfigurationImpl()
- : hasReaderListener(false), receiverQueueSize(1000), readerName(),
subscriptionRolePrefix() {}
+ : hasReaderListener(false),
+ receiverQueueSize(1000),
+ readerName(),
+ subscriptionRolePrefix(),
+ readCompacted(false) {}
};
} // namespace pulsar
#endif /* LIB_READERCONFIGURATIONIMPL_H_ */
diff --git a/pulsar-client-cpp/lib/ReaderImpl.cc
b/pulsar-client-cpp/lib/ReaderImpl.cc
index 5e69a72..45c509b 100644
--- a/pulsar-client-cpp/lib/ReaderImpl.cc
+++ b/pulsar-client-cpp/lib/ReaderImpl.cc
@@ -32,6 +32,7 @@ void ReaderImpl::start(const MessageId& startMessageId) {
ConsumerConfiguration consumerConf;
consumerConf.setConsumerType(ConsumerExclusive);
consumerConf.setReceiverQueueSize(readerConf_.getReceiverQueueSize());
+ consumerConf.setReadCompacted(readerConf_.isReadCompacted());
if (readerConf_.getReaderName().length() > 0) {
consumerConf.setConsumerName(readerConf_.getReaderName());
diff --git a/pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc
b/pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc
index 42c7c73..c8d5453 100644
--- a/pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc
+++ b/pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc
@@ -104,3 +104,12 @@ long pulsar_consumer_get_unacked_messages_timeout_ms(
int pulsar_consumer_is_encryption_enabled(pulsar_consumer_configuration_t
*consumer_configuration) {
return consumer_configuration->consumerConfiguration.isEncryptionEnabled();
}
+
+int pulsar_consumer_is_read_compacted(pulsar_consumer_configuration_t
*consumer_configuration) {
+ return consumer_configuration->consumerConfiguration.isReadCompacted();
+}
+
+void pulsar_consumer_set_read_compacted(pulsar_consumer_configuration_t
*consumer_configuration,
+ int compacted) {
+ consumer_configuration->consumerConfiguration.setReadCompacted(compacted);
+}
diff --git a/pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc
b/pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc
index 55a6dc5..9b764f5 100644
--- a/pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc
+++ b/pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc
@@ -75,4 +75,13 @@ void
pulsar_reader_configuration_set_subscription_role_prefix(pulsar_reader_conf
const char *pulsar_reader_configuration_get_subscription_role_prefix(
pulsar_reader_configuration_t *configuration) {
return configuration->conf.getSubscriptionRolePrefix().c_str();
+}
+
+void
pulsar_reader_configuration_set_read_compacted(pulsar_reader_configuration_t
*configuration,
+ int readCompacted) {
+ configuration->conf.setReadCompacted(readCompacted);
+}
+
+int
pulsar_reader_configuration_is_read_compacted(pulsar_reader_configuration_t
*configuration) {
+ return configuration->conf.isReadCompacted();
}
\ No newline at end of file
diff --git a/pulsar-client-cpp/python/pulsar/__init__.py
b/pulsar-client-cpp/python/pulsar/__init__.py
index d8254d3..8c36c08 100644
--- a/pulsar-client-cpp/python/pulsar/__init__.py
+++ b/pulsar-client-cpp/python/pulsar/__init__.py
@@ -357,7 +357,8 @@ class Client:
receiver_queue_size=1000,
consumer_name=None,
unacked_messages_timeout_ms=None,
- broker_consumer_stats_cache_time_ms=30000
+ broker_consumer_stats_cache_time_ms=30000,
+ is_read_compacted=False
):
"""
Subscribe to the given topic and subscription combination.
@@ -415,9 +416,11 @@ class Client:
_check_type_or_none(str, consumer_name, 'consumer_name')
_check_type_or_none(int, unacked_messages_timeout_ms,
'unacked_messages_timeout_ms')
_check_type(int, broker_consumer_stats_cache_time_ms,
'broker_consumer_stats_cache_time_ms')
+ _check_type(bool, is_read_compacted, 'is_read_compacted')
conf = _pulsar.ConsumerConfiguration()
conf.consumer_type(consumer_type)
+ conf.read_compacted(is_read_compacted)
if message_listener:
conf.message_listener(message_listener)
conf.receiver_queue_size(receiver_queue_size)
diff --git a/pulsar-client-cpp/python/pulsar_test.py
b/pulsar-client-cpp/python/pulsar_test.py
index e2f3e25..31163f8 100755
--- a/pulsar-client-cpp/python/pulsar_test.py
+++ b/pulsar-client-cpp/python/pulsar_test.py
@@ -39,6 +39,13 @@ def doHttpPost(url, data):
req.add_header('Content-Type', 'application/json')
urlopen(req)
+import urllib2
+def doHttpPut(url, data):
+ opener = urllib2.build_opener(urllib2.HTTPHandler)
+ request = urllib2.Request(url, data=data.encode())
+ request.add_header('Content-Type', 'application/json')
+ request.get_method = lambda: 'PUT'
+ opener.open(request)
class PulsarTest(TestCase):
@@ -405,6 +412,54 @@ class PulsarTest(TestCase):
self._check_value_error(lambda: client.create_reader(topic,
MessageId.earliest, reader_name=5))
client.close()
+ def test_publish_compact_and_consume(self):
+ client = Client(self.serviceUrl)
+ topic =
'persistent://sample/standalone/ns1/my-python-test_publish_compact_and_consume'
+ producer = client.create_producer(topic,
producer_name='my-producer-name', batching_enabled=False)
+ self.assertEqual(producer.last_sequence_id(), -1)
+ consumer = client.subscribe(topic, 'my-sub1', is_read_compacted=True)
+ consumer.close()
+ consumer2 = client.subscribe(topic, 'my-sub2', is_read_compacted=False)
+
+ # producer create 2 messages with same key.
+ producer.send('hello-0', partition_key='key0')
+ producer.send('hello-1', partition_key='key0')
+ producer.close()
+
+ # issue compact command, and wait success
+ url=self.adminUrl +
'/admin/persistent/sample/standalone/ns1/my-python-test_publish_compact_and_consume/compaction'
+ doHttpPut(url, '')
+ while True:
+ req = urllib2.Request(url)
+ response = urllib2.urlopen(req)
+ s=response.read()
+ if 'RUNNING' in s:
+ print("Compact still running")
+ print(s)
+ time.sleep(0.2)
+ else:
+ self.assertTrue('SUCCESS' in s)
+ print("Compact Complete now")
+ print(s)
+ break
+
+ # after compact, consumer with `is_read_compacted=True`, expected read
only the second message for same key.
+ consumer1 = client.subscribe(topic, 'my-sub1', is_read_compacted=True)
+ msg0 = consumer1.receive()
+ self.assertEqual(msg0.data(), b'hello-1')
+ consumer1.acknowledge(msg0)
+ consumer1.close()
+
+ # after compact, consumer with `is_read_compacted=False`, expected
read 2 messages for same key.
+ msg0 = consumer2.receive()
+ self.assertEqual(msg0.data(), b'hello-0')
+ consumer2.acknowledge(msg0)
+ msg1 = consumer2.receive()
+ self.assertEqual(msg1.data(), b'hello-1')
+ consumer2.acknowledge(msg1)
+ consumer2.close()
+ client.close()
+
def _check_value_error(self, fun):
try:
fun()
diff --git a/pulsar-client-cpp/python/src/config.cc
b/pulsar-client-cpp/python/src/config.cc
index f03c883..9149add 100644
--- a/pulsar-client-cpp/python/src/config.cc
+++ b/pulsar-client-cpp/python/src/config.cc
@@ -134,6 +134,8 @@ void export_config() {
.def("unacked_messages_timeout_ms",
&ConsumerConfiguration::setUnAckedMessagesTimeoutMs)
.def("broker_consumer_stats_cache_time_ms",
&ConsumerConfiguration::getBrokerConsumerStatsCacheTimeInMs)
.def("broker_consumer_stats_cache_time_ms",
&ConsumerConfiguration::setBrokerConsumerStatsCacheTimeInMs)
+ .def("read_compacted", &ConsumerConfiguration::isReadCompacted)
+ .def("read_compacted", &ConsumerConfiguration::setReadCompacted)
;
class_<ReaderConfiguration>("ReaderConfiguration")
@@ -144,5 +146,7 @@ void export_config() {
.def("reader_name", &ReaderConfiguration::setReaderName)
.def("subscription_role_prefix",
&ReaderConfiguration::getSubscriptionRolePrefix,
return_value_policy<copy_const_reference>())
.def("subscription_role_prefix",
&ReaderConfiguration::setSubscriptionRolePrefix)
+ .def("read_compacted", &ReaderConfiguration::isReadCompacted)
+ .def("read_compacted", &ReaderConfiguration::setReadCompacted)
;
}
diff --git a/pulsar-client-cpp/tests/ConsumerConfigurationTest.cc
b/pulsar-client-cpp/tests/ConsumerConfigurationTest.cc
new file mode 100644
index 0000000..cede638
--- /dev/null
+++ b/pulsar-client-cpp/tests/ConsumerConfigurationTest.cc
@@ -0,0 +1,105 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+#include <pulsar/Client.h>
+#include <gtest/gtest.h>
+
+#include "../lib/Future.h"
+#include "../lib/Utils.h"
+
+using namespace pulsar;
+
+TEST(ConsumerConfigurationTest, testReadCompactPersistentExclusive) {
+ std::string lookupUrl = "pulsar://localhost:8885";
+ std::string topicName = "persistent://prop/unit/ns1/persist-topic";
+ std::string subName = "test-persist-exclusive";
+
+ Result result;
+
+ ConsumerConfiguration config;
+ config.setReadCompacted(true);
+ config.setConsumerType(ConsumerExclusive);
+
+ ClientConfiguration clientConfig;
+ Client client(lookupUrl, clientConfig);
+
+ Consumer consumer;
+ result = client.subscribe(topicName, subName, config, consumer);
+ ASSERT_EQ(ResultOk, result);
+ consumer.close();
+}
+
+TEST(ConsumerConfigurationTest, testReadCompactPersistentFailover) {
+ std::string lookupUrl = "pulsar://localhost:8885";
+ std::string topicName = "persistent://prop/unit/ns1/persist-topic";
+ std::string subName = "test-persist-fail-over";
+
+ Result result;
+
+ ConsumerConfiguration config;
+ config.setReadCompacted(true);
+ config.setConsumerType(ConsumerFailover);
+
+ ClientConfiguration clientConfig;
+ Client client(lookupUrl, clientConfig);
+
+ Consumer consumer;
+ result = client.subscribe(topicName, subName, config, consumer);
+ ASSERT_EQ(ResultOk, result);
+ consumer.close();
+}
+
+TEST(ConsumerConfigurationTest, testReadCompactPersistentShared) {
+ std::string lookupUrl = "pulsar://localhost:8885";
+ std::string topicName = "persistent://prop/unit/ns1/persist-topic";
+ std::string subName = "test-persist-shared";
+
+ Result result;
+
+ ConsumerConfiguration config;
+ config.setReadCompacted(true);
+ config.setConsumerType(ConsumerShared);
+
+ ClientConfiguration clientConfig;
+ Client client(lookupUrl, clientConfig);
+
+ Consumer consumer;
+ result = client.subscribe(topicName, subName, config, consumer);
+ ASSERT_EQ(ResultInvalidConfiguration, result);
+ consumer.close();
+}
+
+TEST(ConsumerConfigurationTest, testReadCompactNonPersistentExclusive) {
+ std::string lookupUrl = "pulsar://localhost:8885";
+ std::string topicName =
"non-persistent://prop/unit/ns1/testNonPersistentTopic";
+ std::string subName = "test-non-persist-exclusive";
+
+ Result result;
+
+ ConsumerConfiguration config;
+ config.setReadCompacted(true);
+ config.setConsumerType(ConsumerExclusive);
+
+ ClientConfiguration clientConfig;
+ Client client(lookupUrl, clientConfig);
+
+ Consumer consumer;
+ result = client.subscribe(topicName, subName, config, consumer);
+ ASSERT_EQ(ResultInvalidConfiguration, result);
+ consumer.close();
+}
\ No newline at end of file
--
To stop receiving notification emails like this one, please contact
[email protected].