diff --git a/include/cppkafka/consumer.h b/include/cppkafka/consumer.h index 5c90985..332b942 100644 --- a/include/cppkafka/consumer.h +++ b/include/cppkafka/consumer.h @@ -403,7 +403,7 @@ public: * \return A Queue object * * \remark Note that this call will disable forwarding to the consumer_queue. - * To restore forwarding (if desired) call Queue::forward_to_queue(consumer_queue) + * To restore forwarding if desired, call Queue::forward_to_queue(consumer_queue) */ Queue get_main_queue() const; @@ -424,7 +424,7 @@ public: * \return A Queue object * * \remark Note that this call will disable forwarding to the consumer_queue. - * To restore forwarding (if desired) call Queue::forward_to_queue(consumer_queue) + * To restore forwarding if desired, call Queue::forward_to_queue(consumer_queue) */ Queue get_partition_queue(const TopicPartition& partition) const; private: diff --git a/src/consumer.cpp b/src/consumer.cpp index d35f453..d2eb797 100644 --- a/src/consumer.cpp +++ b/src/consumer.cpp @@ -47,6 +47,17 @@ using std::equal; namespace cppkafka { +// See: https://github.com/edenhill/librdkafka/issues/1792 +const int rd_kafka_queue_refcount_bug_version = 0x000b0500; +Queue get_queue(rd_kafka_queue_t* handle) { + if (rd_kafka_version() <= rd_kafka_queue_refcount_bug_version) { + return Queue::make_non_owning(handle); + } + else { + return Queue(handle); + } +} + void Consumer::rebalance_proxy(rd_kafka_t*, rd_kafka_resp_err_t error, rd_kafka_topic_partition_list_t *partitions, void *opaque) { TopicPartitionList list = convert(partitions); @@ -262,19 +273,19 @@ MessageList Consumer::poll_batch(size_t max_batch_size, milliseconds timeout) { } Queue Consumer::get_main_queue() const { - Queue queue(Queue::make_non_owning(rd_kafka_queue_get_main(get_handle()))); + Queue queue(get_queue(rd_kafka_queue_get_main(get_handle()))); queue.disable_queue_forwarding(); return queue; } Queue Consumer::get_consumer_queue() const { - return Queue::make_non_owning(rd_kafka_queue_get_consumer(get_handle())); + return get_queue(rd_kafka_queue_get_consumer(get_handle())); } Queue Consumer::get_partition_queue(const TopicPartition& partition) const { - Queue queue(Queue::make_non_owning(rd_kafka_queue_get_partition(get_handle(), - partition.get_topic().c_str(), - partition.get_partition()))); + Queue queue(get_queue(rd_kafka_queue_get_partition(get_handle(), + partition.get_topic().c_str(), + partition.get_partition()))); queue.disable_queue_forwarding(); return queue; }