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

Google My Business Profile PubSub消息监听器延迟数小时接收消息求助

Pub/Sub消息延迟排查及代码优化建议

问题背景

你编写了以下Pub/Sub监听代码用于接收新评论通知,同时配置了对应的通知规则,但监听器接收消息存在数小时延迟:

监听代码

const projectId = 'project-id'
const pubSubClient = new PubSub({ projectId });
const subscriptionNameOrId =
  "subscription-id";

async function pubsubReview() {
  try {
    const [subs] = await pubSubClient.getSubscriptions();
    console.log("Subscriptions:");
    subs.forEach((s) => console.log(s.name));

    const subscription = pubSubClient.subscription(subscriptionNameOrId);

    // Receive callbacks for new messages on the subscription
    subscription.on("message", async (message) => {
      await message.ackWithResponse();
      console.log("Received message:", message.data.toString());
    });

    // Receive callbacks for errors on the subscription
    subscription.on("error", (error) => {
      console.error("Received error:", error);
    });
  } catch (err) {
    console.log("\n\n", err, "\n\n");
  }
}

pubsubReview();

通知配置代码

const updateMask = "notification_types,pubsub_topic";
const projectID = "project-id";
const topicName = "topic-name";
const updateData = {
  notificationTypes: ["NEW_REVIEW", "UPDATED_REVIEW"],
  pubsubTopic: `projects/${projectID}/topics/${topicName}`,
};

延迟问题排查方向

  • 检查订阅消息积压:查看GCP控制台中该订阅的「未确认消息数」指标,若存在大量积压,说明旧消息处理速度跟不上,导致新消息延迟投递。
  • 验证ACK超时设置:如果订阅的ACK超时时间过长,未及时确认的消息会被重新投递,挤占新消息的处理资源,可根据实际消息处理时长调整超时时间。
  • 排查发布端批量配置:若通知发送端开启了批量发送(如设置较大的batchSize或batchDuration),会攒够一定数量或等待特定时间才发送消息,造成延迟,检查通知服务的批量参数配置。
  • 检查服务器资源:监听进程所在服务器的CPU、内存、网络带宽不足,会限制消息的接收和处理速度,导致延迟,查看服务器资源使用率。
  • 确认区域一致性:如果Topic所在GCP区域和监听客户端的区域不一致,跨区域传输会增加延迟,尽量保持Topic和客户端在同一区域。

代码优化建议

  1. 移除不必要的订阅列表查询
    getSubscriptions()仅用于打印订阅名称,生产环境中无实际业务价值,还会增加API调用次数和启动时间,建议删除这部分代码。

  2. 增加消息确认的错误处理
    当前message回调未捕获ackWithResponse()的异常,若确认失败会导致消息无法被正确标记为已处理,进而重复投递。优化后代码:

    subscription.on("message", async (message) => {
      try {
        await message.ackWithResponse();
        console.log("Received message:", {
          messageId: message.id,
          publishTime: message.publishTime.toISOString(),
          receiveTime: new Date().toISOString(),
          data: message.data.toString()
        });
      } catch (ackErr) {
        console.error("Failed to acknowledge message:", ackErr);
        // 确认失败时标记消息为未处理,触发重新投递
        message.nack();
      }
    });
    
  3. 配置流式拉取的流量控制
    通过flowControl设置并发处理的最大消息数和字节数,避免进程因同时处理过多消息而过载:

    const subscription = pubSubClient.subscription(subscriptionNameOrId, {
      flowControl: {
        maxMessages: 100, // 同时处理的最大消息数
        maxBytes: 10 * 1024 * 1024 // 限制同时处理的消息总大小为10MB
      }
    });
    
  4. 增强日志信息
    增加消息ID、发布时间、接收时间等字段,便于排查延迟问题,直观看到消息的实际发布与接收时间差。

  5. 确保进程持续运行
    脚本形式的监听进程可能因无活跃任务而退出,增加信号处理和保活逻辑:

    // 监听退出信号,优雅关闭订阅
    process.on('SIGINT', async () => {
      console.log('Initiating graceful shutdown...');
      await subscription.close();
      console.log('Subscription closed, exiting process');
      process.exit(0);
    });
    
    // 保持进程活跃
    setInterval(() => {}, 1000);
    
  6. 监听订阅关闭事件
    增加close事件监听,及时记录订阅意外关闭的情况:

    subscription.on("close", () => {
      console.log("Subscription connection closed unexpectedly");
    });
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 11:00:55