Google Pub/Sub同步拉取同批次消息为何分次获取?求解决方案
首先可以肯定地说,你遇到的情况绝对不是Google Pub/Sub的Bug,而是它作为分布式消息系统的设计特性导致的。下面给你拆解原因和对应的解决办法:
为什么会出现分次获取消息的情况?
Pub/Sub的消息是分散存储在多个分布式节点(分片)上的,当你发起同步拉取请求时,服务端会从部分分片查询可用消息,而不是立刻遍历所有分片——这是为了保证系统的低延迟和高可扩展性。尤其是当消息数量很少(比如你这里的4条)时,很可能这些消息被分配到了不同的分片,单次拉取请求只从部分分片拿到了消息,剩下的需要后续拉取才能获取到。
另外,你的代码里调用fetch时传了returnImmediately=true,这个参数的作用是:服务端不会等待凑够你设置的maxMessages数量,有多少消息就立刻返回多少。哪怕当前只有1条可用,也直接返回,不会等待其他分片的消息。这也是导致你分次拿到消息的关键因素之一。
可以采取这些措施优化拉取行为
1. 调整returnImmediately参数为false
把拉取时的returnImmediately改成false,这样服务端会等待一段默认时间(最多5秒),尽量凑够你设置的maxMessages(你这里设了100000,远大于4)之后再返回。这样能极大提高一次拉取到所有消息的概率。
修改你的调用代码:
List<ReceivedMessage> messages = mySyncSubscriber.fetch(100000, false);
2. 增加拉取重试逻辑(补充循环判断)
即使调整了returnImmediately,极端情况下还是可能出现分次拉取的情况。你可以在循环里增加一个“拉取到空消息后再退出”的逻辑,或者固定重试几次,确保不会遗漏消息。比如:
// 记录是否还有未拉取的消息 boolean hasMoreMessages = true; int emptyPullCount = 0; // 最多允许3次空拉取,避免无限循环 while (hasMoreMessages && emptyPullCount < 3) { log("before fetched messages by order"); List<ReceivedMessage> messages = mySyncSubscriber.fetch(100000, false); log("fetched " + messages.size() + " messages by order"); if (messages.isEmpty()) { emptyPullCount++; hasMoreMessages = false; } else { emptyPullCount = 0; messages = messages.stream() .sorted(Comparator.comparingLong(message -> message.getMessage().getPublishTime().getSeconds())) .collect(Collectors.toList()); log("before processMessages"); List<String> ackIds = processMessages(messages); log("after processMessages"); mySyncSubscriber.sendAck(ackIds); } }
3. 考虑使用异步订阅者(可选)
如果你的业务场景允许,Google Pub/Sub的异步订阅者(Subscriber类)会自动管理拉取、重试和消息分发,它会在后台持续拉取消息,并且尽量批量获取,可能更适合你的需求。不过如果必须用同步拉取,前面的两种方法就足够解决问题了。
补充说明
你代码里对消息按发布时间排序的逻辑是没问题的,和分次拉取无关——哪怕分次拿到消息,排序后依然能保证处理顺序,不用担心乱序问题。
内容的提问来源于stack exchange,提问作者Elad Benda

