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

Node.js中Google PubSub订阅者接收首条消息后无法拉取后续消息

问题描述

Node.js应用通过Google PubSub发布消息,订阅者仅能拉取第一条消息,后续无法接收。流量低(每秒约1条),必须重启PM2服务器才能接收第二条消息。

发布者代码
const postNotificationsService = async (payload) => {
  try {

    const dataBuffer = Buffer.from(JSON.stringify(payload));
    console.log("Connected to Producer");
    console.log("Successfully posted to web notification Queue");
    let publishRes = await pubsub
      .topic(topicName)
      .publishMessage({ data: dataBuffer });

    console.log(publishRes)
    return { message: "Success" };
  } catch (e) {
    console.log(e)
    throw new Error(e.message);
  }
};
订阅者代码
const pullSubscriptionsMessage = () => {
  // const subscription = pubsub.subscription(subscriptionName);

  const subscriberOptions = {
    flowControl: {
      maxMessages: maxInProgress,
    },
  };

  // References an existing subscription.
  // Note that flow control settings are not persistent across subscribers.
  const subscription = pubsub.subscription(
    subscriptionName,
    subscriberOptions
  );

  console.log(
    `Subscriber to subscription ${subscription.name} is ready to receive messages at a controlled volume of ${maxInProgress} messages.`
  );

  // Create an event handler to handle messages
  let messageCount = 0;
  const messageHandler = async message => {
    console.log(`Received message ${message.id}:`);
    // console.log(`Data: ${message.data}`);
    // console.log(`tAttributes: ${message.attributes}`);
    messageCount += 1;
    let messagePayload = JSON.parse(message.data.toString())
    await mailerService.sendEmail(messagePayload);

    await message.ack();
  };

  

  // Listen for new messages/errors until timeout is hit
  subscription.on('message', messageHandler);
  setTimeout(() => {
    subscription.close();
    console.log(`${messageCount} message(s) received.`);
  }, timeout * 1000);

}

pullSubscriptionsMessage();
运行输出
0|serverHTTP  | Connected to Producer
0|serverHTTP  | Successfully posted to web notification Queue
0|serverHTTP  | Received message 6339231272672371:
0|serverHTTP  | 6339231272672371
0|serverHTTP  | 1 message(s) received.
0|serverHTTP  | Connected to Producer
0|serverHTTP  | Successfully posted to web notification Queue
0|serverHTTP  | 6339223654064343
问题原因及解决方案

核心问题

订阅者代码中的setTimeout会在指定时间后调用subscription.close(),直接关闭了订阅连接。从运行输出可见,第一条消息处理完成后打印了1 message(s) received.,说明订阅已被终止,后续发布的消息自然无法被拉取。

修复方案

移除自动关闭订阅的setTimeout逻辑,让订阅保持长连接以持续接收消息:

const pullSubscriptionsMessage = () => {
  const subscriberOptions = {
    flowControl: {
      maxMessages: maxInProgress,
    },
  };

  const subscription = pubsub.subscription(
    subscriptionName,
    subscriberOptions
  );

  console.log(
    `Subscriber to subscription ${subscription.name} is ready to receive messages at a controlled volume of ${maxInProgress} messages.`
  );

  let messageCount = 0;
  const messageHandler = async message => {
    console.log(`Received message ${message.id}:`);
    messageCount += 1;
    let messagePayload = JSON.parse(message.data.toString())
    await mailerService.sendEmail(messagePayload);

    await message.ack();
    console.log(`${messageCount} message(s) received so far.`);
  };

  subscription.on('message', messageHandler);
  
  // 可选:添加错误监听,避免异常导致订阅静默失效
  subscription.on('error', (err) => {
    console.error(`Subscriber error: ${err.message}`);
  });

  // 可选:监听进程信号实现优雅关闭
  process.on('SIGINT', async () => {
    console.log('Closing subscription...');
    await subscription.close();
    console.log('Subscription closed.');
    process.exit(0);
  });
}

pullSubscriptionsMessage();

额外注意事项

  • 确保maxInProgress值设置合理,避免同时处理过多消息导致资源耗尽
  • 检查mailerService.sendEmail是否存在未捕获异常,若抛出错误会导致message.ack()不执行,消息会被重新投递
  • 若需要定期重启订阅,可通过进程管理工具(如PM2)配置重启策略,而非在代码中强制关闭

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 09:25:37