C++ 使用 Protobuf 实现跨语言 Kafka 数据交换
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++