diff --git a/pulsar-client-cpp/lib/ClientConnection.cc b/pulsar-client-cpp/lib/ClientConnection.cc index f79cd423fc1cc..077f36fdf583d 100644 --- a/pulsar-client-cpp/lib/ClientConnection.cc +++ b/pulsar-client-cpp/lib/ClientConnection.cc @@ -1338,11 +1338,14 @@ void ClientConnection::close() { if (keepAliveTimer_) { keepAliveTimer_->cancel(); + keepAliveTimer_.reset(); } if (consumerStatsRequestTimer_) { consumerStatsRequestTimer_->cancel(); + consumerStatsRequestTimer_.reset(); } + for (ProducersMap::iterator it = producers.begin(); it != producers.end(); ++it) { HandlerBase::handleDisconnection(ResultConnectError, shared_from_this(), it->second); } diff --git a/pulsar-client-cpp/lib/ConnectionPool.cc b/pulsar-client-cpp/lib/ConnectionPool.cc index 86c89489979d4..124f81ff855ca 100644 --- a/pulsar-client-cpp/lib/ConnectionPool.cc +++ b/pulsar-client-cpp/lib/ConnectionPool.cc @@ -42,6 +42,19 @@ ConnectionPool::ConnectionPool(const ClientConfiguration& conf, ExecutorServiceP poolConnections_(poolConnections), mutex_() {} +ConnectionPool::~ConnectionPool() { + std::unique_lock lock(mutex_); + if (poolConnections_) { + for (auto cnxIt = pool_.begin(); cnxIt != pool_.end(); cnxIt++) { + ClientConnectionPtr cnx = cnxIt->second.lock(); + if (cnx && !cnx->isClosed()) { + cnx->close(); + } + } + pool_.clear(); + } +} + Future ConnectionPool::getConnectionAsync( const std::string& logicalAddress, const std::string& physicalAddress) { std::unique_lock lock(mutex_); diff --git a/pulsar-client-cpp/lib/ConnectionPool.h b/pulsar-client-cpp/lib/ConnectionPool.h index f4c1f68f93acf..21011c5bdbf85 100644 --- a/pulsar-client-cpp/lib/ConnectionPool.h +++ b/pulsar-client-cpp/lib/ConnectionPool.h @@ -36,6 +36,8 @@ class PULSAR_PUBLIC ConnectionPool { ConnectionPool(const ClientConfiguration& conf, ExecutorServiceProviderPtr executorProvider, const AuthenticationPtr& authentication, bool poolConnections = true); + ~ConnectionPool(); + /** * Get a connection from the pool. *

diff --git a/pulsar-client-cpp/lib/ProducerImpl.cc b/pulsar-client-cpp/lib/ProducerImpl.cc index 666df3b27f184..8488c3d01b161 100644 --- a/pulsar-client-cpp/lib/ProducerImpl.cc +++ b/pulsar-client-cpp/lib/ProducerImpl.cc @@ -88,9 +88,6 @@ ProducerImpl::ProducerImpl(ClientImplPtr client, const std::string& topic, const ProducerImpl::~ProducerImpl() { LOG_DEBUG(getName() << "~ProducerImpl"); - if (dataKeyGenTImer_) { - dataKeyGenTImer_->cancel(); - } closeAsync(ResultCallback()); printStats(); } @@ -473,6 +470,8 @@ void ProducerImpl::printStats() { void ProducerImpl::closeAsync(CloseCallback callback) { Lock lock(mutex_); + cancelTimers(); + if (state_ != Ready) { lock.unlock(); if (callback) { @@ -682,10 +681,20 @@ void ProducerImpl::start() { HandlerBase::start(); } void ProducerImpl::shutdown() { Lock lock(mutex_); state_ = Closed; + cancelTimers(); + producerCreatedPromise_.setFailed(ResultAlreadyClosed); +} + +void ProducerImpl::cancelTimers() { + if (dataKeyGenTImer_) { + dataKeyGenTImer_->cancel(); + dataKeyGenTImer_.reset(); + } + if (sendTimer_) { sendTimer_->cancel(); + sendTimer_.reset(); } - producerCreatedPromise_.setFailed(ResultAlreadyClosed); } bool ProducerImplCmp::operator()(const ProducerImplPtr& a, const ProducerImplPtr& b) const { diff --git a/pulsar-client-cpp/lib/ProducerImpl.h b/pulsar-client-cpp/lib/ProducerImpl.h index cb2a8a6e8a990..3a3780417d30f 100644 --- a/pulsar-client-cpp/lib/ProducerImpl.h +++ b/pulsar-client-cpp/lib/ProducerImpl.h @@ -139,6 +139,8 @@ class ProducerImpl : public HandlerBase, bool encryptMessage(proto::MessageMetadata& metadata, SharedBuffer& payload, SharedBuffer& encryptedPayload); + void cancelTimers(); + typedef std::unique_lock Lock; ProducerConfiguration conf_;