You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在C++中结合Confluent Schema Registry实现Kafka Protobuf序列化生产者

C++集成Confluent Schema Registry实现Protobuf序列化的Kafka生产者

要在C中实现集成Confluent Schema Registry的Protobuf序列化Kafka生产者,你需要依赖**Confluent C Kafka客户端**(包含Schema Registry支持)和Protobuf库,以下是具体实现步骤:

1. 准备依赖

  • 安装Confluent C++ Kafka客户端:确保安装的版本包含Schema Registry模块(可通过confluent-kafka-cpp包获取,或源码编译时开启Schema Registry支持)
  • 安装Protobuf编译器(protoc)和C++库:用于编译自定义Protobuf消息定义

2. 定义并编译Protobuf消息

先创建Protobuf消息文件(示例:test.proto):

syntax = "proto3";

package example;

message Meal {
  string name = 1;
  int32 calories = 2;
}

使用protoc编译为C++代码:

protoc --cpp_out=. test.proto

编译后生成test.pb.h和test.pb.cc,需引入到项目中。

3. 编写生产者代码

以下是完整的生产者实现示例,包含Schema Registry集成与Protobuf序列化:

#include <iostream>
#include <string>
#include <map>
#include <chrono>
#include <librdkafka/rdkafkacpp.h>
#include <confluent/schema_registry/SchemaRegistryClient.h>
#include <confluent/schema_registry/ProtobufSerializer.h>

int main() {
    std::string errstr;

    // 1. 初始化Schema Registry客户端配置
    std::map<std::string, std::string> sr_config = {
        {"schema.registry.url", "http://localhost:8081"}
    };
    auto sr_client = std::make_shared<confluent::schema_registry::SchemaRegistryClient>(sr_config);

    // 2. 创建Protobuf序列化器(绑定自定义消息类型)
    auto protobuf_serializer = std::make_shared<confluent::schema_registry::ProtobufSerializer<example::Meal>>(sr_client);

    // 3. 配置Kafka生产者参数
    std::map<std::string, std::string> producer_config = {
        {"bootstrap.servers", "localhost:9092"},
        {"acks", "all"}
    };

    // 4. 创建Kafka生产者实例
    RdKafka::Producer* producer = RdKafka::Producer::create(
        RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL),
        producer_config,
        errstr
    );
    if (!producer) {
        std::cerr << "创建生产者失败: " << errstr << std::endl;
        return 1;
    }

    // 5. 构造并序列化Protobuf消息
    example::Meal meal;
    meal.set_name("hamburger");
    meal.set_calories(500);

    std::string serialized_value;
    protobuf_serializer->serialize(meal, "TEST_PROTO", serialized_value);

    // 6. 发送消息到Kafka
    std::string key = "meal_1";
    RdKafka::ErrorCode resp = producer->produce(
        RdKafka::Topic::create(producer, "TEST_PROTO", nullptr, errstr),
        RdKafka::Topic::PARTITION_UA,
        RdKafka::Producer::RK_MSG_COPY,
        const_cast<char*>(serialized_value.data()), serialized_value.size(),
        const_cast<char*>(key.data()), key.size(),
        nullptr
    );

    if (resp != RdKafka::ERR_NO_ERROR) {
        std::cerr << "发送消息失败: " << RdKafka::err2str(resp) << std::endl;
    } else {
        std::cout << "消息发送成功" << std::endl;
    }

    // 等待消息发送完成
    producer->flush(std::chrono::seconds(5));

    delete producer;
    return 0;
}

关键说明

  • Schema Registry客户端:负责与Schema Registry服务交互,自动处理Protobuf Schema的注册、版本管理与获取
  • Protobuf序列化器:ProtobufSerializer<T>会将Protobuf消息实例序列化为符合Confluent规范的字节流(包含Schema ID等元数据)
  • 生产者配置:无需指定内置value序列化器,直接通过自定义序列化器处理消息体

编译注意事项

编译时需链接对应库:

g++ -std=c++17 kafka_protobuf_producer.cpp test.pb.cc -o kafka_protobuf_producer -lrdkafka++ -lprotobuf -lconfluent-schema-registry

内容的提问来源于stack exchange,提问作者Alchemist

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 19:15:11