如何获取librdkafka生产者发送队列中已生产未发送的消息数量
librdkafka生产者查询待发送消息数量解决方案
完全可以实现该需求,librdkafka原生提供了对应接口直接获取已生产未发送的消息数量。
核心API说明
你可以直接调用rd_kafka_outq_len()接口,传入生产者实例指针,返回值就是当前已调用生产接口但还未从生产者侧发出的消息数量,完全匹配你对“已生产未发送队列”的统计需求。
适配你的场景的流控方案
你已经关闭了ack机制(acks=0),该场景下消息只要被发送到网络层就会从待统计队列中移除,不会等待broker响应,统计值完全符合你“只要消息从生产者侧发出就算完成”的逻辑。
具体流控逻辑可以按以下方式实现:
- 每次调用生产接口前,先调用
rd_kafka_outq_len()获取当前待发送队列长度 - 若长度超过你预设的阈值,暂停生产,间隔一小段时间后再重试,避免队列堆积
除了自行判断队列长度外,你也可以直接配置librdkafka的queue.buffering.max.messages参数设置内置缓冲队列的最大消息数,当队列满时调用rd_kafka_produce()会直接返回RD_KAFKA_RESP_ERR__QUEUE_FULL错误,你捕获该错误做流控即可。
示例代码片段
// 自定义待发送队列最大阈值 #define MAX_PENDING_MSG_NUM 5000 while (生产消息逻辑循环) { // 获取当前待发送队列长度 int pending_cnt = rd_kafka_outq_len(producer_handle); if (pending_cnt >= MAX_PENDING_MSG_NUM) { // 队列已满,等待10ms后重试 usleep(10 * 1000); continue; } // 队列未满,执行消息生产逻辑 int ret = rd_kafka_produce( topic_handle, partition, RD_KAFKA_MSG_F_COPY, msg_payload, payload_len, msg_key, key_len, NULL ); // 后续可增加生产返回值处理逻辑 }
内容的提问来源于stack exchange,提问作者Pras
相关产品推荐
相关产品推荐

