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

GCP PubSub拉取模式订阅有未消费消息但拉取不到消息问题咨询

问题排查&修复方案

1 首先修复代码中的明确错误

  • 变量引用错误:MessageReceiverImpl中解析得到的消息对象是logReqMsg,但调用logSave.addLogMsg()时传入了未定义的battleLogMsg,会直接抛出空指针异常,即使拉取到消息也会处理失败。
  • 缺少异常兜底逻辑:switch判断编码类型时没有default分支,如果消息没有携带googclient_schemaencoding属性,encoding为null,会直接跳过解析逻辑,后续打印日志时调用JsonFormat.printer().print(null)同样会抛出空指针。且当前代码无论处理是否成功都会调用consumer.ack()确认消息,出错的消息会直接被丢弃,不会留存到未消费队列中。
  • 修复示例:
@Override
public void receiveMessage(PubsubMessage message, AckReplyConsumer consumer) {
    // 先打日志确认是否收到消息
    LOGGER.info("Received message id: " + message.getMessageId());
    ByteString data = message.getData();
    String encoding = message.getAttributesMap().get("googclient_schemaencoding");
    Req.LogReq logReqMsg = null;
    try {
        if (encoding == null) {
            // 兜底默认编码,或者打错误日志后nack
            LOGGER.warning("Message has no googclient_schemaencoding attribute, message id: " + message.getMessageId());
            consumer.nack();
            return;
        }
        switch (encoding) {
            case "BINARY":
                logReqMsg = Req.LogReq.parseFrom(data);
                break;
            case "JSON":
                Req.LogReq.Builder msgBuilder = Req.LogReq.newBuilder();
                JsonFormat.parser().merge(data.toStringUtf8(), msgBuilder);
                logReqMsg = msgBuilder.build();
                break;
            default:
                LOGGER.warning("Unsupported encoding: " + encoding + ", message id: " + message.getMessageId());
                consumer.nack();
                return;
        }
        LOGGER.info(JsonFormat.printer().omittingInsignificantWhitespace().print(logReqMsg));
        // 修正变量引用
        logSave.addLogMsg(logReqMsg);
        consumer.ack();
    } catch (Exception e) {
        LOGGER.error("Process message failed, message id: " + message.getMessageId(), e);
        // 处理失败时nack,让消息重新投递
        consumer.nack();
    }
}

2 检查订阅器存活逻辑

你使用的是异步Subscriber,调用startAsync().awaitRunning()只会等待订阅器启动完成,不会阻塞主线程。如果你的主程序启动订阅器后没有加阻塞逻辑(比如CountDownLatch.await()、死循环等待关闭信号等),程序启动后会直接退出,订阅器也会被销毁,自然无法拉取到消息。

3 权限&网络排查

  • 确认K8s Pod使用的服务账号拥有roles/pubsub.subscriber角色,具备目标订阅的拉取、确认权限。
  • 在Pod内执行curl https://pubsub.googleapis.com/$discovery/rest?version=v1,确认网络可以连通Pub/Sub服务,没有被防火墙、网络策略拦截。

4 配置校验

  • 确认代码中使用的projectId、subscriptionId和你在GCP控制台看到的有7条未消费消息的订阅ID完全一致,没有拼写、大小写错误。
  • 确认订阅的拉取模式配置正确,没有被误设为推送模式。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 09:15:03