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和客户端在同一区域。
代码优化建议
移除不必要的订阅列表查询
getSubscriptions()仅用于打印订阅名称,生产环境中无实际业务价值,还会增加API调用次数和启动时间,建议删除这部分代码。增加消息确认的错误处理
当前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(); } });配置流式拉取的流量控制
通过flowControl设置并发处理的最大消息数和字节数,避免进程因同时处理过多消息而过载:const subscription = pubSubClient.subscription(subscriptionNameOrId, { flowControl: { maxMessages: 100, // 同时处理的最大消息数 maxBytes: 10 * 1024 * 1024 // 限制同时处理的消息总大小为10MB } });增强日志信息
增加消息ID、发布时间、接收时间等字段,便于排查延迟问题,直观看到消息的实际发布与接收时间差。确保进程持续运行
脚本形式的监听进程可能因无活跃任务而退出,增加信号处理和保活逻辑:// 监听退出信号,优雅关闭订阅 process.on('SIGINT', async () => { console.log('Initiating graceful shutdown...'); await subscription.close(); console.log('Subscription closed, exiting process'); process.exit(0); }); // 保持进程活跃 setInterval(() => {}, 1000);监听订阅关闭事件
增加close事件监听,及时记录订阅意外关闭的情况:subscription.on("close", () => { console.log("Subscription connection closed unexpectedly"); });
内容的提问来源于stack exchange,提问作者Sehej Wahla
相关产品推荐
相关产品推荐

