一、rdkafka簡介
rdkafka是一個高性能的消息隊列中間件,它基於C++編寫而成,具有很好的可擴展性和可靠性,適用於各類分散式場景下的數據傳輸等場景,目前已應用於多個大型在線系統中。
rdkafka主要特點如下:
1、高吞吐量:消息傳輸快速
2、支持多語言:除了C++,還支持Java、Python等語言
3、支持多種傳輸協議:比如TCP、UDP等協議
4、高可靠性:消息傳輸過程中保證冪等性和事務性
5、易擴展:支持集群和分散式部署
二、rdkafka使用方法
1. 安裝rdkafka
在Linux系統下,可以通過以下命令安裝rdkafka:
git clone https://github.com/edenhill/librdkafka.git
cd librdkafka
./configure
make
make install
2. rdafka C++ API的使用
2.1 生產者示例
以下是rdkafka生產者示例的代碼:
#include <librdkafka/rdkafkacpp.h>
#include <string>
#include <iostream>
int main() {
std::string errstr;
RdKafka::Producer *producer = RdKafka::Producer::create({
{"bootstrap.servers", "localhost"},
{"message.timeout.ms", "3000"}
}, errstr);
if (!producer) {
std::cerr << "Failed to create producer: " << errstr << std::endl;
return 1;
}
RdKafka::Topic *topic = RdKafka::Topic::create(producer, "test", nullptr, errstr);
if (!topic) {
std::cerr << "Failed to create topic: " << errstr <produce(topic, RdKafka::Topic::PARTITION_UA /*指定分區,這裡使用默認值*/, RdKafka::Producer::RK_MSG_COPY /*複製消息體內容*/, message, nullptr);
// 同步等待生產者發送消息的結果
producer->flush(RdKafka::Topic::PARTITION_UA, 3000);
}
2.2 消費者示例
以下是rdkafka消費者示例的代碼:
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
class MyEventCb : public RdKafka::EventCb {
public:
void event_cb (RdKafka::Event &event) {
switch (event.type())
{
case RdKafka::Event::EVENT_ERROR:
if (event.err()) {
std::cerr << "ERROR (event): " << RdKafka::err2str(event.err()) << std::endl;
}
break;
case RdKafka::Event::EVENT_LOG:
std::cout << "LOG-" << event.severity() << "-" << event.fac() << ": " << event.str() << std::endl;
break;
default:
std::cerr << "Unhandled event type: " << event.type() <set("group.id", "test", errstr) != RdKafka::Conf::CONF_OK) {
std::cerr << errstr <set("event_cb", &eventCb, errstr);
// 創建消費者實例
RdKafka::KafkaConsumer *consumer = RdKafka::KafkaConsumer::create(conf, errstr);
if (!consumer) {
std::cerr << "Failed to create consumer: " << errstr <subscribe(topics) != RdKafka::ERR_NO_ERROR) {
std::cerr << "Failed to subscribe to topics: " << errstr <consume(1000);
if (!message) {
continue;
}
if (message->err() == RdKafka::ERR_NO_ERROR) {
std::cout << "Received message:" << std::endl;
std::cout << "payload:" <payload()), message->len()) << std::endl;
}
delete message;
}
return 0;
}
三、rdkafka集成到項目中的示例
以下是rdkafka集成到項目中的示例代碼:
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <thread>
#include <chrono>
class KafkaConsumer {
public:
KafkaConsumer() {
conf_ = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
std::string errstr;
// 指定Kafka集群地址
if (conf_->set("bootstrap.servers", "localhost:9092", errstr) != RdKafka::Conf::CONF_OK) {
throw std::runtime_error("Failed to set broker configuration: " + errstr);
}
// 配置消費者ID
if (conf_->set("group.id", "test-group", errstr) != RdKafka::Conf::CONF_OK) {
throw std::runtime_error("Failed to set group.id configuration: " + errstr);
}
// 創建消費者實例
consumer_ = RdKafka::KafkaConsumer::create(conf_.get(), errstr);
if (!consumer_) {
throw std::runtime_error("Failed to create KafkaConsumer: " + errstr);
}
// 訂閱主題
std::vector<std::string> topics = {"test-topic"};
if (consumer_->subscribe(topics) != RdKafka::ERR_NO_ERROR) {
throw std::runtime_error("Failed to subscribe to topics: " + errstr);
}
}
~KafkaConsumer() {
consumer_->close();
delete consumer_;
}
void consumeMessage() {
while (true) {
// 從Kafka伺服器接收消息
RdKafka::Message *message = consumer_->consume(1000);
if (!message) {
continue;
}
if (message->err() == RdKafka::ERR_NO_ERROR) {
std::cout << "Received message: "
<payload()), message->len()) <set("bootstrap.servers", "localhost:9092", errstr) != RdKafka::Conf::CONF_OK) {
throw std::runtime_error("Failed to set broker configuration: " + errstr);
}
// 創建生產者實例
producer_ = RdKafka::Producer::create(conf_.get(), errstr);
if (!producer_) {
throw std::runtime_error("Failed to create KafkaProducer: " + errstr);
}
}
~KafkaProducer() {
producer_->flush(0);
delete producer_;
}
void produceMessage(const std::string& message) {
RdKafka::Topic *topic = RdKafka::Topic::create(producer_, "test-topic", nullptr, errstr_);
if (!topic) {
throw std::runtime_error("Failed to create topic: " + errstr_);
}
// 創建消息
RdKafka::Producer::Message* msg = RdKafka::Producer::Message::create(
topic,
RdKafka::Producer::RK_MSG_COPY /*複製消息體內容*/,
const_cast<char*>(message.c_str()), message.size(), nullptr, nullptr);
// 發送消息,不需要等待ack
producer_->produce(topic, RdKafka::Topic::PARTITION_UA /*指定分區,這裡使用默認值*/, RdKafka::Producer::RK_MSG_COPY /*複製消息體內容*/, msg, nullptr);
delete msg;
}
private:
std::unique_ptr<RdKafka::Conf> conf_;
RdKafka::Producer *producer_;
std::string errstr_;
};
int main() {
KafkaProducer producer;
producer.produceMessage("Hello World!");
// 延遲1s,確保消息已經被消費
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
KafkaConsumer consumer;
consumer.consumeMessage();
return 0;
}
原創文章,作者:小藍,如若轉載,請註明出處:https://www.506064.com/zh-tw/n/246856.html