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

如何避免Cloud Run返回错误时PubSub出现无限循环?

问题:Pub/Sub 推送订阅重复投递非200响应的消息

我正在构建一套文件处理架构:Cloud Functions脚本下载文件后交由Cloud Run应用处理。原本Cloud Functions直接调用Cloud Run接口,但需要等待处理完成,所以引入Pub/Sub解耦:Cloud Functions发布消息,Cloud Run通过推送订阅接收。但遇到问题:当Cloud Run返回非200状态码(比如400)时,Pub/Sub会无限重复投递这条消息,除非手动从主题清除。

期望行为:只要Cloud Run给出响应(无论状态码是200、400还是500),就让Pub/Sub确认该消息(ACK),不再重复投递。后续我会通过日志和告警发现问题,修复后再重新发送消息。请问这个需求是否可行,还是我的设计有问题?

当前配置代码

Terraform 创建Pub/Sub主题和订阅

resource "google_pubsub_topic" "schema_validator_trigger" {
  name   = "oney-schema-validator-trigger"
  labels = var.tags
}

resource "google_pubsub_subscription" "schema_validator_subscription" {
  name  = "oney_schema_validator_subscription"
  topic = google_pubsub_topic.schema_validator_trigger.name

  ack_deadline_seconds = 600

  expiration_policy {
    ttl = ""
  }

  push_config {
    push_endpoint = "https://my-cloudrun-url.a.run.app/my/endpoint"
    oidc_token {
      service_account_email = google_service_account.downloader_sa.email
    }
    attributes = {
      x-goog-version = "v1"
    }
  }
}

Cloud Run 消息处理代码

from flask import Flask, request

app = Flask(__name__)

@app.route("/my/endpoint", methods=["POST"])
def my_endpoint():
    # Check received PubSub message
    envelope = request.get_json()

    expected_params = ["bucketname", "filename"]
    msg, status_code = utils.requests.validate_pubsub_message(envelope, expected_params)
    if status_code != 200:
        utils.logging.log_message(msg, severity="ERROR")
        return msg, status_code

    # 后续处理逻辑
    ...

解决方案:确保所有响应都触发Pub/Sub ACK

你的需求完全可行,问题出在Pub/Sub推送订阅的默认行为:只有当端点返回2xx状态码时,Pub/Sub才会自动ACK消息;返回非2xx状态码时,会将消息重新放入待投递队列,直到达到重试次数上限(默认无限重试,除非配置死信策略)。

要实现“只要有响应就ACK”的目标,只需调整Cloud Run的处理逻辑:

1. 统一返回2xx状态码,携带错误信息

无论消息验证失败还是处理异常,都让端点返回200状态码,同时在响应体中携带错误详情用于日志记录。这样Pub/Sub会判定投递成功,停止重复发送。

修改后的代码示例:

from flask import Flask, request

app = Flask(__name__)

@app.route("/my/endpoint", methods=["POST"])
def my_endpoint():
    try:
        envelope = request.get_json()

        expected_params = ["bucketname", "filename"]
        msg, status_code = utils.requests.validate_pubsub_message(envelope, expected_params)
        if status_code != 200:
            utils.logging.log_message(msg, severity="ERROR")
            # 返回200状态码,携带错误信息
            return {"error": msg}, 200

        # 后续处理逻辑
        ...
        return {"status": "success"}, 200
    except Exception as e:
        # 捕获所有未预期异常,记录日志后返回200
        error_msg = f"Unexpected error: {str(e)}"
        utils.logging.log_message(error_msg, severity="CRITICAL")
        return {"error": error_msg}, 200

2. 可选:配置死信队列兜底

如果担心极端情况(比如Cloud Run完全无响应)导致消息无限重试,可以给订阅配置死信队列,将多次投递失败的消息转发到专用主题,避免占用资源:

修改Terraform订阅配置:

resource "google_pubsub_subscription" "schema_validator_subscription" {
  name  = "oney_schema_validator_subscription"
  topic = google_pubsub_topic.schema_validator_trigger.name

  ack_deadline_seconds = 600

  expiration_policy {
    ttl = ""
  }

  # 新增死信队列配置
  dead_letter_policy {
    dead_letter_topic = google_pubsub_topic.schema_validator_dlq.name
    max_delivery_attempts = 5 # 最多重试5次后转入死信队列
  }

  push_config {
    push_endpoint = "https://my-cloudrun-url.a.run.app/my/endpoint"
    oidc_token {
      service_account_email = google_service_account.downloader_sa.email
    }
    attributes = {
      x-goog-version = "v1"
    }
  }
}

# 新增死信主题
resource "google_pubsub_topic" "schema_validator_dlq" {
  name   = "oney-schema-validator-dlq"
  labels = var.tags
}

设计合理性说明

你的解耦思路是正确的:用Pub/Sub实现Cloud Functions和Cloud Run的异步通信,避免阻塞Cloud Functions。只要确保Cloud Run始终返回2xx状态码,就能让Pub/Sub停止重试,同时通过日志和告警追踪错误,后续可手动从死信队列(或原始主题)重新发送消息,完全符合你的需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 09:52:45