如何获取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
相关产品推荐
相关产品推荐

