当前位置:首页 > 技术 > 正文内容

C++ 使用 Protobuf 实现跨语言 Kafka 数据交换

访客 技术 2026年10月5日 1

Kafka 简介与核心特性

Apache Kafka 是一个分布式流处理平台,具备以下核心功能:

  • 发布和订阅记录流,类似于消息队列系统
  • 以容错和持久化的方式存储记录流
  • 实时处理记录流

Kafka 主要应用于两类场景:

  • 构建实时数据管道,在不同系统间可靠传输数据
  • 开发实时流处理应用,对数据流进行转换或响应

librdkafka 库概述

librdkafka 是用 C 语言实现的 Kafka 协议客户端库,提供生产者、消费者和管理客户端功能。其设计注重消息传递的可靠性和高性能,可达到每秒百万级消息的处理能力。

主要特性包括:

  • 完整的一次性语义(Exactly-Once-Semantics)支持
  • 高级生产者,包含幂等和事务生产者
  • 高级平衡消费者(需 broker 0.9+ 版本)
  • 简单消费者(传统模式)
  • 管理客户端
  • 多种压缩格式支持:snappy、gzip、lz4、zstd
  • SSL 和 SASL 安全认证支持
  • 良好的跨平台兼容性

Kafka C++ 客户端实现

以下是使用 librdkafka 进行 Kafka 消息生产和消费的 C++ 示例代码:

#include "rdkafkacpp.h"
#include <signal.h>
#include <iostream>
#include <string>

static volatile sig_atomic_t keep_running = 1;

static void signal_handler(int sig) {
    keep_running = 0;
}

class DeliveryCallback : public RdKafka::DeliveryReportCb {
public:
    void dr_cb(RdKafka::Message &message) override {
        std::string status;
        switch (message.status()) {
            case RdKafka::Message::MSG_STATUS_NOT_PERSISTED:
                status = "未持久化";
                break;
            case RdKafka::Message::MSG_STATUS_POSSIBLY_PERSISTED:
                status = "可能已持久化";
                break;
            case RdKafka::Message::MSG_STATUS_PERSISTED:
                status = "已持久化";
                break;
            default:
                status = "未知状态";
                break;
        }
        
        std::cout << "消息投递状态 (" << message.len() << " 字节): " 
                  << status << ": " << message.errstr() << std::endl;
    }
};

class EventCallback : public RdKafka::EventCb {
public:
    void event_cb(RdKafka::Event &event) override {
        switch (event.type()) {
            case RdKafka::Event::EVENT_ERROR:
                if (event.fatal()) {
                    std::cerr << "致命错误 ";
                    keep_running = 0;
                }
                std::cerr << "错误 (" << RdKafka::err2str(event.err()) << "): " 
                          << event.str() << std::endl;
                break;
            case RdKafka::Event::EVENT_LOG:
                std::cerr << "日志-" << event.severity() << "-" << event.fac().c_str() 
                          << ": " << event.str().c_str() << std::endl;
                break;
            default:
                std::cerr << "事件 " << event.type() << " (" 
                          << RdKafka::err2str(event.err()) << "): " 
                          << event.str() << std::endl;
                break;
        }
    }
};

void process_message(RdKafka::Message* message) {
    switch (message->err()) {
        case RdKafka::ERR_NO_ERROR: {
            std::cout << "读取消息偏移量: " << message->offset() << std::endl;
            
            const RdKafka::Headers *headers = message->headers();
            if (headers) {
                std::vector<RdKafka::Headers::Header> hdr_list = headers->get_all();
                for (const auto& header : hdr_list) {
                    if (header.value() != nullptr) {
                        printf(" 头部: %s = \"%.*s\"\n",
                               header.key().c_str(),
                               (int)header.value_size(), 
                               (const char *)header.value());
                    } else {
                        printf(" 头部: %s = NULL\n", header.key().c_str());
                    }
                }
            }
            
            printf("消息内容: %.*s\n",
                   static_cast<int>(message->len()),
                   static_cast<const char *>(message->payload()));
            break;
        }
        case RdKafka::ERR__PARTITION_EOF:
            std::cout << "分区结束" << std::endl;
            break;
        default:
            std::cerr << "消费失败: " << message->errstr() << std::endl;
            keep_running = 0;
    }
}

int main(int argc, char **argv) {
    std::string brokers = "localhost:9092";
    std::string topic_name = "test-topic";
    std::string error_string;
    
    // 创建配置对象
    RdKafka::Conf *global_conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    RdKafka::Conf *topic_conf = RdKafka::Conf::create(RdKafka::Conf::CONF_TOPIC);
    
    // 设置 broker 列表
    global_conf->set("metadata.broker.list", brokers, error_string);
    
    signal(SIGINT, signal_handler);
    signal(SIGTERM, signal_handler);
    
    // 生产者模式
    DeliveryCallback delivery_cb;
    global_conf->set("dr_cb", &delivery_cb, error_string);
    global_conf->set("default_topic_conf", topic_conf, error_string);
    
    RdKafka::Producer *producer = RdKafka::Producer::create(global_conf, error_string);
    if (!producer) {
        std::cerr << "创建生产者失败: " << error_string << std::endl;
        return 1;
    }
    
    std::cout << "创建生产者: " << producer->name() << std::endl;
    
    // 发送消息
    std::string message_content = "Hello Kafka from C++";
    RdKafka::Headers *msg_headers = RdKafka::Headers::create();
    msg_headers->add("source", "cpp-client");
    
    RdKafka::ErrorCode result = producer->produce(
        topic_name,
        RdKafka::Topic::PARTITION_UA,
        RdKafka::Producer::RK_MSG_COPY,
        const_cast<char*>(message_content.c_str()), 
        message_content.size(),
        nullptr, 0,
        0,
        msg_headers,
        nullptr);
        
    if (result != RdKafka::ERR_NO_ERROR) {
        std::cerr << "发送失败: " << RdKafka::err2str(result) << std::endl;
        delete msg_headers;
    } else {
        std::cout << "消息已发送 (" << message_content.size() << " 字节)" << std::endl;
    }
    
    // 等待消息投递完成
    while (keep_running && producer->outq_len() > 0) {
        producer->poll(1000);
    }
    
    delete producer;
    delete global_conf;
    delete topic_conf;
    
    RdKafka::wait_destroyed(5000);
    return 0;
}

Kafka 环境启动

启动 ZooKeeper 服务:

zookeeper-server-start.sh config/zookeeper.properties

启动 Kafka 服务:

kafka-server-start.sh config/server.properties

Protobuf 数据序列化

使用 Protocol Buffers 定义数据结构可以实现高效的跨语言数据交换。首先定义 .proto 文件:

syntax = "proto3";

import "google/protobuf/timestamp.proto";

package tutorial;

message Person {
  string name = 1;
  int32 id = 2;
  string email = 3;
  
  enum PhoneType {
    MOBILE = 0;
    HOME = 1;
    WORK = 2;
  }
  
  message PhoneNumber {
    string number = 1;
    PhoneType type = 2;
  }
  
  repeated PhoneNumber phones = 4;
  google.protobuf.Timestamp last_updated = 5;
}

message AddressBook {
  repeated Person people = 1;
}

编译生成 C++ 代码:

protoc --cpp_out=. addressbook.proto

使用生成的类进行数据序列化:

#include "addressbook.pb.h"
#include <fstream>
#include <iostream>

int main() {
    tutorial::AddressBook address_book;
    
    // 添加人员信息
    tutorial::Person* person = address_book.add_people();
    person->set_name("张三");
    person->set_id(123);
    person->set_email("zhangsan@example.com");
    
    tutorial::Person::PhoneNumber* phone = person->add_phones();
    phone->set_number("13800138000");
    phone->set_type(tutorial::Person::MOBILE);
    
    // 序列化到字符串
    std::string serialized_data;
    if (!address_book.SerializeToString(&serialized_data)) {
        std::cerr << "序列化失败" << std::endl;
        return -1;
    }
    
    std::cout << "序列化数据长度: " << serialized_data.size() << " 字节" << std::endl;
    
    // 反序列化验证
    tutorial::AddressBook parsed_book;
    if (!parsed_book.ParseFromString(serialized_data)) {
        std::cerr << "反序列化失败" << std::endl;
        return -1;
    }
    
    std::cout << "解析成功,人员数量: " << parsed_book.people_size() << std::endl;
    
    return 0;
}

运行效果

生产者输出:

创建生产者: rdkafka#producer-1
消息已发送 (21 字节)

消费者输出:

读取消息偏移量: 6
 头部: source = "cpp-client"
消息内容: Hello Kafka from C++

相关文章

Linux crontab 详解

1) crontab 是什么cron 是 Linux 的定时任务守护进程;crontab 是用来编辑/查看“按时间周期执行命令”的表(cron table)。常见两类:用户 crontab:每个用户一份(crontab -e 编辑)系统级 crontab / cron.d:可指定执行用户(/etc/crontab、/etc/cron.d/*)2) crontab 时间...

富文本里可以允许的 HTML 属性

一、所有标签默认允许的安全属性(极少)class        (可选)id           (通常建议禁用)title️ 注意:id 容易被滥用做锚点注入,很多系统直接禁用class 允许的话最好只允许固定前缀(如 editor-*)二、a 标签允许属性<a href="" t...

Mac 安装 Node.js 指南

方法一:通过官网安装包(最简单,适合初学者)如果你只是想快速安装并开始使用,这是最直接的方法。访问 Node.js 官网。页面会显示两个版本:LTS (Recommended For Most Users):长期支持版,最稳定。建议选这个。Current:最新特性版,包含最新功能但可能不够稳定。下载 .pkg 安装包并运行。按照安装向导点击“下一步”即可完成。方法二:使用 Homebrew 安装(...

Laravel 事件和监听器创建

在 Laravel 中,使用 Artisan 命令创建 Events(事件) 和 Listeners(监听器) 是非常高效的。你可以通过以下几种方式来实现:1. 手动创建单个 Event如果你只想创建一个事件类,可以使用 make:event 命令:Bashphp artisan make:event UserRegistered执行后,文件将生成在 app/Even...

自定义域名解析神器 dnsmasq

什么是 dnsmasq?dnsmasq 是一个轻量级、功能强大的网络服务工具,专为小型和中等规模网络设计。它是一个综合的网络基础设施解决方案[1]。dnsmasq 能做什么?功能说明应用场景DNS 转发与缓存将 DNS 查询转发到上游服务器(ISP、Google DNS 等),并在本地缓存结果加快 DNS 查询速度,减少外部 DNS 流量本地 DNS解析本地网络设备的主机名,无需编辑&n...

linux screen 用法详情 (nohup 的替代方案)

一、screen 是什么?能干嘛?screen 是一个终端复用器,可以:在一个 SSH 会话中开多个“虚拟终端”SSH 断线后,程序仍然在后台运行随时重新连接到原来的会话特别适合:nohup 的替代方案跑脚本 / 爬虫 / 训练模型运维、远程开发二、安装 screen# CentOS / Rocky / Almayum install -y screen# Debian / Ubuntuapt i...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。