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

如何获取Apache Pulsar共享订阅下未消费的队列任务?

Apache Pulsar共享订阅获取待处理任务/数量的可行性与实现方案

可行性结论

完全可行。Apache Pulsar通过Admin API和客户端API提供了获取共享订阅下待处理消息(未被消费者确认的消息)数量及具体内容的能力,足以支撑分布式任务队列的监控需求。

实现方案

1. 获取待处理任务数量

方式一:使用Pulsar Admin CLI

通过命令行工具直接查询指定主题和订阅的统计信息,提取msgBacklog字段值即为待处理消息数:

pulsar-admin subscriptions stats persistent://tenant/namespace/topic-name shared-sub-name

返回的JSON结果中,msgBacklog字段对应未确认消息总量。

方式二:使用Admin Client SDK(以Java为例)

通过编程方式调用Admin API获取统计数据:

PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl("http://pulsar-broker:8080").build();
SubscriptionStats stats = admin.subscriptions().getStats("persistent://tenant/namespace/topic-name", "shared-sub-name");
long backlogCount = stats.getMsgBacklog();

方式三:客户端本地统计(仅单消费者)

单个消费者实例可通过自身API获取本地未确认的消息数,但仅适用于单消费者场景,分布式集群下需汇总所有消费者数据:

Consumer<?> consumer = pulsarClient.newConsumer()
    .topic("persistent://tenant/namespace/topic-name")
    .subscriptionName("shared-sub-name")
    .subscriptionType(SubscriptionType.Shared)
    .subscribe();
ConsumerStats consumerStats = consumer.getStats();
long localBacklog = consumerStats.getUnackedMessages();

2. 获取具体待处理任务内容

共享订阅的待处理消息分散在多个cursor(每个消费者对应一个cursor)中,需遍历所有cursor获取消息:

方式一:Admin CLI批量查看

先获取订阅的所有cursor列表,再逐个cursor peek消息:

# 获取所有cursor ID
pulsar-admin subscriptions get-cursors persistent://tenant/namespace/topic-name shared-sub-name

# 针对每个cursor peek消息(示例获取前10条)
pulsar-admin topics peek persistent://tenant/namespace/topic-name --subscription shared-sub-name --cursor cursor-id-1 --num-messages 10

peek操作不会移动cursor位置,不会影响正常消费流程。

方式二:Admin Client SDK编程获取(以Java为例)

// 获取订阅的所有cursor
List<String> cursors = admin.subscriptions().getCursors("persistent://tenant/namespace/topic-name", "shared-sub-name");

// 遍历每个cursor获取消息
for (String cursor : cursors) {
    List<Message<byte[]>> messages = admin.topics().peekMessages(
        "persistent://tenant/namespace/topic-name",
        "shared-sub-name",
        cursor,
        10 // 单次获取的消息数量
    );
    // 处理消息内容
    for (Message<byte[]> msg : messages) {
        System.out.println("待处理任务内容:" + new String(msg.getValue()));
    }
}

注意事项

  • 共享订阅的待处理消息总量是所有cursor的backlog之和,统计时需确保汇总所有cursor数据。
  • 调用Admin API查询时,避免过于频繁,建议根据业务需求设置合理的查询间隔,防止给Broker带来性能压力。
  • 批量获取待处理消息时,需控制单次获取的消息数量,避免内存溢出,可通过分页方式逐步获取。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 21:52:44