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

如何获取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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 00:24:01