rdkafka詳解

一、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-hant/n/246856.html

(0)
打賞 微信掃一掃 微信掃一掃 支付寶掃一掃 支付寶掃一掃
小藍的頭像小藍
上一篇 2024-12-12 13:18
下一篇 2024-12-12 13:18

相關推薦

  • 神經網絡代碼詳解

    神經網絡作為一種人工智能技術,被廣泛應用於語音識別、圖像識別、自然語言處理等領域。而神經網絡的模型編寫,離不開代碼。本文將從多個方面詳細闡述神經網絡模型編寫的代碼技術。 一、神經網…

    編程 2025-04-25
  • Linux sync詳解

    一、sync概述 sync是Linux中一個非常重要的命令,它可以將文件系統緩存中的內容,強制寫入磁盤中。在執行sync之前,所有的文件系統更新將不會立即寫入磁盤,而是先緩存在內存…

    編程 2025-04-25
  • Python安裝OS庫詳解

    一、OS簡介 OS庫是Python標準庫的一部分,它提供了跨平台的操作系統功能,使得Python可以進行文件操作、進程管理、環境變量讀取等系統級操作。 OS庫中包含了大量的文件和目…

    編程 2025-04-25
  • 詳解eclipse設置

    一、安裝與基礎設置 1、下載eclipse並進行安裝。 2、打開eclipse,選擇對應的工作空間路徑。 File -> Switch Workspace -> [選擇…

    編程 2025-04-25
  • MPU6050工作原理詳解

    一、什麼是MPU6050 MPU6050是一種六軸慣性傳感器,能夠同時測量加速度和角速度。它由三個傳感器組成:一個三軸加速度計和一個三軸陀螺儀。這個組合提供了非常精細的姿態解算,其…

    編程 2025-04-25
  • Linux修改文件名命令詳解

    在Linux系統中,修改文件名是一個很常見的操作。Linux提供了多種方式來修改文件名,這篇文章將介紹Linux修改文件名的詳細操作。 一、mv命令 mv命令是Linux下的常用命…

    編程 2025-04-25
  • Java BigDecimal 精度詳解

    一、基礎概念 Java BigDecimal 是一個用於高精度計算的類。普通的 double 或 float 類型只能精確表示有限的數字,而對於需要高精度計算的場景,BigDeci…

    編程 2025-04-25
  • git config user.name的詳解

    一、為什麼要使用git config user.name? git是一個非常流行的分布式版本控制系統,很多程序員都會用到它。在使用git commit提交代碼時,需要記錄commi…

    編程 2025-04-25
  • Python輸入輸出詳解

    一、文件讀寫 Python中文件的讀寫操作是必不可少的基本技能之一。讀寫文件分別使用open()函數中的’r’和’w’參數,讀取文件…

    編程 2025-04-25
  • nginx與apache應用開發詳解

    一、概述 nginx和apache都是常見的web服務器。nginx是一個高性能的反向代理web服務器,將負載均衡和緩存集成在了一起,可以動靜分離。apache是一個可擴展的web…

    編程 2025-04-25

發表回復

登錄後才能評論