Google Marketplace SaaS集成中Pub/Sub消息无法确认问题排查
Google Marketplace SaaS集成:重复事件通知与消息确认失败问题排查
问题描述
我正在开展Google Marketplace SaaS集成工作,运行集成测试时,持续接收到同一条事件通知,且消息无法被确认。已遵循Google SaaS集成相关文档开展工作。
相关代码片段
const { PubSub } = require('@google-cloud/pubsub'); const { authClient, getAccessToken } = require('../../services/googleAuthClient'); const axios = require('axios'); // project id const projectId = 'your_project_id'; // create pubsub instance const pubsub = new PubSub({ projectId, // project id authClient // service account credentials }); const subscriptionName = 'your_subscription_name'; async function listenForMessages(subscriptionName) { // Create subscription with adjusted flow control and acknowledgment deadline const subscription = pubsub.subscription(subscriptionName); // Create message handler const messageHandler = async (message) => { console.log(`Received message ${message.id}`); console.log(`Data: ${message.data}`); console.log(`Attributes: ${JSON.stringify(message.attributes)}`); // Parse the message data const parsedMessage = JSON.parse(message.data); const eventType = parsedMessage.eventType; const { eventId } = parsedMessage; // Debounce check to avoid processing the same message multiple times if (processedMessages.has(eventId)) { console.log(`Skipping already processed message with eventId: ${eventId}`); try { await message.ackWithResponse(); console.log(`Ack for message ${message.id} successful.`); } catch (e) { console.log(`Ack for message ${message.id} failed with error: ${e.errorCode}`); } return; } try { // Mark the event as processed processedMessages.add(eventId); // Handle the event based on eventType await handleEvent(eventType, parsedMessage); // Acknowledge the message immediately after processing try { await message.ackWithResponse(); console.log(`Ack for message ${message.id} successful.`); } catch (e) { console.log(`Ack for message ${message.id} failed with error: ${e.errorCode}`); } } catch (error) { console.error(`Failed to process message ${message.id}:`, error); // Optionally nack the message if processing failed and you want to retry // message.nack(); } finally { // Clean up the processed message after a delay setTimeout(() => { processedMessages.delete(eventId); }, 60000); // 60 seconds, adjust as needed } }; // Start listening for messages subscription.on('message', messageHandler); } // Call the function to start listening for messages indefinitely listenForMessages(subscriptionName); async function handleEvent(eventType, parsedMessage) { console.log('eventType:', eventType); console.log('parsedMessage:', parsedMessage); const { eventId } = parsedMessage; // Debounce check // if (processedMessages.has(eventId)) { // console.log(`Skipping already processed message with eventId: ${eventId}`); // return; // } try { // Mark the event as processed processedMessages.add(eventId); // Switch-case to handle different event types switch (eventType) { // case 'ACCOUNT_ACTIVE': // await handleAccountActive(parsedMessage); // break; case 'ACCOUNT_DELETED': await handleAccountDeleted(parsedMessage); break; case 'ENTITLEMENT_CREATION_REQUESTED': await handleEntitlementCreationRequested(parsedMessage); break; case 'ENTITLEMENT_OFFER_ACCEPTED': await handleEntitlementOfferAccepted(parsedMessage); break; case 'ENTITLEMENT_ACTIVE': await handleEntitlementActive(parsedMessage); break; case 'ENTITLEMENT_PLAN_CHANGE_REQUESTED': await handleEntitlementPlanChangeRequested(parsedMessage); break; case 'ENTITLEMENT_PLAN_CHANGED': await handleEntitlementPlanChanged(parsedMessage); break; case 'ENTITLEMENT_PLAN_CHANGE_CANCELLED': await handleEntitlementPlanChangeCancelled(parsedMessage); break; case 'ENTITLEMENT_PENDING_CANCELLATION': await handleEntitlementPendingCancellation(parsedMessage); break; case 'ENTITLEMENT_CANCELLATION_REVERTED': await handleEntitlementCancellationReverted(parsedMessage); break; case 'ENTITLEMENT_CANCELLED': await handleEntitlementCancelled(parsedMessage); break; case 'ENTITLEMENT_CANCELLING': await handleEntitlementCancelling(parsedMessage); break; case 'ENTITLEMENT_RENEWED': await handleEntitlementRenewed(parsedMessage); break; case 'ENTITLEMENT_OFFER_ENDED': await handleEntitlementOfferEnded(parsedMessage); break; case 'ENTITLEMENT_DELETED': await handleEntitlementDeleted(parsedMessage); break; default: console.log(`Unknown event type: ${eventType}`); } } catch (error) { console.error('Error handling event:', error); } finally { // Remove the eventId from the processed set after a delay setTimeout(() => { processedMessages.delete(eventId); }, 60000); // 60 seconds, adjust as needed } }
问题根源分析
- 重复管理事件处理状态:
listenForMessages中已将eventId加入processedMessages,但handleEvent里又重复执行添加和清理操作,导致状态混乱,可能引发ack逻辑失效。 - 超时配置不匹配:
processedMessages的清理延迟设为60秒,若Pub/Sub订阅的确认超时小于这个值,消息会在未被确认时重新投递,造成重复接收。 - JSON解析未做异常处理:如果消息数据不是合法JSON,
JSON.parse会抛出错误,直接进入catch块但未执行ack或nack,导致消息超时重投。 - ack失败无重试机制:
ackWithResponse失败后仅打印日志,未做重试或标记消息为未处理,消息会因未确认而重复投递。 - 全局变量未初始化:代码中未显式声明全局的
processedMessages,每次运行都会创建新Set,无法实现去重效果。
修复方案
移除重复的状态管理逻辑:删除
handleEvent中的processedMessages.add(eventId)和对应的finally清理代码,统一在messageHandler中管理事件状态。统一超时配置:
- 将Pub/Sub订阅的确认超时设置为至少120秒(大于消息处理+清理的最大耗时)。
- 调整
processedMessages的清理延迟为确认超时的1.5倍(如180秒),避免消息重投时eventId已被删除。
添加JSON解析异常处理:
let parsedMessage; try { parsedMessage = JSON.parse(message.data); } catch (parseError) { console.error(`Failed to parse message ${message.id}:`, parseError); await message.ackWithResponse(); return; }优化ack失败的重试逻辑:
const maxAckRetries = 3; let ackSuccess = false; for (let i = 0; i < maxAckRetries; i++) { try { await message.ackWithResponse(); console.log(`Ack for message ${message.id} successful.`); ackSuccess = true; break; } catch (e) { console.log(`Ack attempt ${i+1} failed for message ${message.id}:`, e.errorCode); await new Promise(resolve => setTimeout(resolve, 1000 * (i+1))); } } if (!ackSuccess) { console.error(`All ack attempts failed for message ${message.id}`); message.nack(); // 标记消息为未处理,触发重新投递 }初始化全局去重集合:在代码顶部添加全局声明:
const processedMessages = new Set();
内容的提问来源于stack exchange,提问作者Gaurav Pant
相关产品推荐
相关产品推荐

