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

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,无法实现去重效果。

修复方案

  1. 移除重复的状态管理逻辑:删除handleEvent中的processedMessages.add(eventId)和对应的finally清理代码,统一在messageHandler中管理事件状态。

  2. 统一超时配置:

    • 将Pub/Sub订阅的确认超时设置为至少120秒(大于消息处理+清理的最大耗时)。
    • 调整processedMessages的清理延迟为确认超时的1.5倍(如180秒),避免消息重投时eventId已被删除。
  3. 添加JSON解析异常处理:

    let parsedMessage;
    try {
      parsedMessage = JSON.parse(message.data);
    } catch (parseError) {
      console.error(`Failed to parse message ${message.id}:`, parseError);
      await message.ackWithResponse();
      return;
    }
    
  4. 优化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(); // 标记消息为未处理,触发重新投递
    }
    
  5. 初始化全局去重集合:在代码顶部添加全局声明:

    const processedMessages = new Set();
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 07:02:07