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
相关产品推荐
相关产品推荐

