关键参数:
-
bootstrap.servers: Kafka 集群的连接地址,生产者通过这些地址发现 Kafka 集群。可以配置多个地址,以提高可用性。
-
key.serializer: 将消息的 Key 序列化成字节数组的类。
-
value.serializer: 将消息的 Value 序列化成字节数组的类。
-
acks: 指定了生产者需要多少个 Broker 确认才能认为消息发送成功。常见的值有 0(不等待确认,性能最高,可靠性最低)、1(等待 Leader Broker 确认)、all 或 -1(等待所有同步副本确认,可靠性最高,性能最低)。
-
retries: 当发送失败时,生产者重试发送消息的次数。
-
batch.size: 生产者尝试将多个消息批量发送,减少网络请求,提升效率。 该参数定义了每个批次的大小。
-
linger.ms: 生产者在批量发送消息之前等待的时间。即使批次未满,也会在等待时间结束后发送。
-
buffer.memory: 生产者用于缓存等待发送消息的总内存大小。
(2)创建生产者实例: 在配置完所有必要的参数后,需要根据配置创建一个 KafkaProducer 实例。这个实例将负责与 Kafka 集群建立连接,并执行消息发送的操作。
(3)构建待发送的消息(ProducerRecord): 消息在发送之前需要被封装成 ProducerRecord 对象。ProducerRecord 包含了消息的关键信息,例如目标 Topic,Key(可选),Value,以及可选的分区信息。
关键元素:
-
topic: 消息要发送到的 Kafka Topic 的名称。
-
key: 消息的 Key,用于分区(如果未指定分区)。 Key 相同的所有消息将会被发送到同一个分区。
-
value: 消息的内容。
-
partition: 可选参数,指定消息发送到的分区。 如果没有指定,Kafka 会根据 Key 和配置的分区策略来选择分区。
(4)发送消息:使用 KafkaProducer 实例的 send() 方法发送 ProducerRecord。 send() 方法是异步的,它会立即返回一个 Future 对象。
-
异步发送: 生产者会将消息放入内部缓冲区,并由后台线程批量发送到 Kafka Broker。
-
同步发送 (可选): 可以调用 Future 对象的 get() 方法来同步等待消息发送完成的结果。 这种方式会阻塞当前线程,直到 Broker 返回确认信息。
-
回调函数 (可选): 可以在 send() 方法中传入一个 Callback 对象,当消息发送完成(成功或失败)时,会调用该对象的回调函数。
-
错误处理: 在异步发送中,可以通过 Future 对象或回调函数来处理发送失败的情况。 常见的错误包括网络错误,Topic 不存在,权限不足等。
(5)关闭生产者实例:在程序结束时,必须关闭 KafkaProducer 实例,释放资源,并确保所有待发送的消息都被发送。 close() 方法会阻塞当前线程,直到所有未完成的请求完成。
重要性: 如果不关闭生产者,可能会导致消息丢失,或者资源泄漏。

二、Kafka C++ API 详解
librdkafka 提供了强大的 C++ API 用于与 Kafka 集群进行交互。 包括 RdKafka::Conf, RdKafka::Message, RdKafka::DeliveryReportCb, RdKafka::Event, RdKafka::EventCb, RdKafka::PartitionerCb, RdKafka::Topic, 和 RdKafka::Producer。
2.1、RdKafka::Conf
RdKafka::Conf 用于配置 Kafka 客户端(Producer 或 Consumer)。
ConfType Enum:
|
CONF_GLOBAL |
全局配置,适用于整个客户端实例。 |
|
CONF_TOPIC |
Topic 配置,只适用于特定的 Topic。 |
ConfResult Enum:
|
CONF_UNKNOWN |
未知的配置属性。 |
|
CONF_INVALID |
配置属性值无效。 |
|
CONF_OK |
配置成功。 |
方法:
|
static Conf * create(ConfType type); |
创建配置对象。 type 参数指定配置类型 (CONF_GLOBAL 或 CONF_TOPIC)。 |
|
Conf::ConfResult set(const std::string &name, const std::string &value, std::string &errstr); |
设置配置对象的属性值。 name 是属性名称,value 是属性值(字符串类型),errstr 用于返回错误信息。 |
|
Conf::ConfResult set(const std::string &name, DeliveryReportCb *dr_cb, std::string &errstr); |
设置投递报告回调函数 dr_cb。 |
|
Conf::ConfResult set(const std::string &name, EventCb *event_cb, std::string &errstr); |
设置事件回调函数 event_cb。 |
|
Conf::ConfResult set(const std::string &name, const Conf *topic_conf, std::string &errstr); |
设置用于自动订阅 Topic 的默认 Topic 配置。 |
|
Conf::ConfResult set(const std::string &name, PartitionerCb *partitioner_cb, std::string &errstr); |
设置自定义分区器回调函数 partitioner_cb。 注意: 必须是 CONF_TOPIC 类型的配置对象。 |
|
Conf::ConfResult set(const std::string &name, PartitionerKeyPointerCb *partitioner_kp_cb,std::string &errstr); |
设置自定义分区器KeyPointer回调函数 partitioner_kp_cb。 |
|
Conf::ConfResult set(const std::string &name, SocketCb *socket_cb, std::string &errstr); |
设置 Socket 回调函数 socket_cb。 |
|
Conf::ConfResult set(const std::string &name, OpenCb *open_cb, std::string &errstr); |
设置 Open 回调函数 open_cb。 |
|
Conf::ConfResult set(const std::string &name, RebalanceCb *rebalance_cb, std::string &errstr); |
设置 Rebalance 回调函数 rebalance_cb(用于 Consumer Group)。 |
|
Conf::ConfResult set(const std::string &name, OffsetCommitCb *offset_commit_cb, std::string &errstr); |
设置 Offset 提交回调函数 offset_commit_cb。 |
|
Conf::ConfResult get(const std::string &name, std::string &value) const; |
查询指定属性 name 的配置值,并将结果存储在 value 中。 |
2.2、RdKafka::Message
RdKafka::Message 表示一条消费或生产的消息,或者是一个事件(例如错误)。
|
std::string errstr() const; |
如果消息是错误事件,返回错误字符串;否则返回空字符串。 |
|
ErrorCode err() const; |
如果消息是错误事件,返回错误代码;否则返回 ERR_NO_ERROR (0)。 |
|
Topic * topic() const; |
返回消息的 Topic 对象。 如果 Topic 对象不是通过 RdKafka::Topic::create() 创建的,则使用 topic_name()。 |
|
std::string topic_name() const; |
返回消息的 Topic 名称。 |
|
int32_t partition() const; |
如果分区可用,返回分区号。 |
|
void * payload() const; |
返回消息数据(payload)的指针。 |
|
size_t len() const; |
返回消息数据的长度。 |
|
const std::string * key() const; |
返回字符串类型的消息 Key 指针。如果 Key 不存在,可能返回 nullptr。使用前需要判空。 |
|
const void * key_pointer() const; |
返回void类型的消息 Key 指针。如果 Key 不存在,可能返回 nullptr。使用前需要判空。 |
|
size_t key_len() const; |
返回消息 key的二进制长度。 |
|
int64_t offset () const; |
返回消息或错误的位移 (Offset)。 |
|
void * msg_opaque() const; |
返回通过 RdKafka::Producer::produce() 提供的 msg_opaque 指针。 |
|
virtual MessageTimestamp timestamp() const = 0; |
返回消息的时间戳。 |
|
virtual int64_t latency() const = 0; |
返回生产消息的微秒级时间延迟(在 produce 函数内部)。如果延迟不可用,返回 -1。 |
|
virtual struct rd_kafka_message_s *c_ptr () = 0; |
返回底层数据结构的 C 句柄 (rd_kafka_message_t)。 不推荐直接使用 C API,除非 C++ API 没有提供相应功能。 |
|
virtual Status status () const = 0; |
返回消息在 Topic Log 中的持久化状态。 |
|
virtual RdKafka::Headers *headers () = 0; |
返回消息头。 |
|
virtual RdKafka::Headers *headers (RdKafka::ErrorCode *err) = 0; |
返回消息头,错误信息会输出到err。 |
2.3、RdKafka::DeliveryReportCb
投递报告回调函数。 每当使用 RdKafka::Producer::produce() 发送的消息被成功传递或遇到永久性错误(或重试次数耗尽)时,该回调函数会被调用。


