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

使用Flask处理Pub/Sub推送订阅长任务的ACK问题及解决方案

问题场景与故障现象

我使用EventArc生成事件,通过Cloud Run处理,以Cloud Pub/Sub作为传输层。代码结构如下:

from flask import Flask, request
import time

app = Flask(__name__)

@app.route('/', methods=['POST'])
def index():
    # 提取Pub/Sub消息
    envelope = request.get_json()
    message = envelope['message']

    try:
        # 处理消息
        # ...
        time.sleep(200)

        # 返回200确认消息
        return '', 200
    except Exception as e:
        # 记录异常
        # ...

        # 返回500让消息重试
        return '', 500

if __name__ == '__main__':
    app.run(port=8080, debug=True)

故障现象:当消息处理时间超过1分钟时,即使处理任务正常执行并返回200 OK,Pub/Sub仍会反复发送相同消息,说明消息未被成功确认(ACK)。

补充:此前已将Pub/Sub的ACK超时设置为10分钟(处理时间在1-3分钟),重试间隔最小值设为9分50秒。

消息无法ACK的原因

核心问题是Cloud Run的默认请求超时限制为1分钟:
即使你配置了Pub/Sub的ACK超时为10分钟,Cloud Run会在请求到达1分钟时强制断开连接,此时Pub/Sub无法收到你的200响应,会判定消息处理失败,进而触发重试。

直接解决ACK失败的方案
  1. 调整Cloud Run服务的请求超时时间
    在Cloud Run控制台或通过gcloud命令,将服务的请求超时设置为大于你的最大处理时间(比如4分钟),确保处理过程能完整执行并返回200响应。
  2. 保持Pub/Sub的ACK超时配置
    继续维持10分钟的ACK超时,确保其大于Cloud Run的请求超时,避免Pub/Sub在Cloud Run处理完成前就重新投递消息。
实现「立即ACK+异步处理」的方案

如果想彻底规避请求超时的影响,可以采用接收到消息后立即ACK,再异步执行处理逻辑的方案,具体实现如下:

技术选型

使用Cloud Tasks作为异步任务队列(比Cloud Run内部后台线程更可靠,避免实例缩容导致任务中断),将消息内容提交到队列后立即返回200给Pub/Sub,再由Cloud Tasks触发Cloud Run的处理接口执行任务。

代码示例

from flask import Flask, request
import json
from google.cloud import tasks_v2

app = Flask(__name__)
tasks_client = tasks_v2.CloudTasksClient()

# 替换为你的项目信息
PROJECT_ID = "你的项目ID"
REGION = "你的服务区域"
QUEUE_NAME = "你的Cloud Tasks队列名"
QUEUE_PATH = tasks_client.queue_path(PROJECT_ID, REGION, QUEUE_NAME)
# 替换为你的Cloud Run处理接口地址
PROCESS_TASK_URL = "https://你的Cloud Run服务域名/process-task"

@app.route('/', methods=['POST'])
def receive_message():
    envelope = request.get_json()
    message = envelope['message']

    try:
        # 构建异步任务
        task_payload = {
            "http_request": {
                "http_method": tasks_v2.HttpMethod.POST,
                "url": PROCESS_TASK_URL,
                "headers": {"Content-Type": "application/json"},
                "body": json.dumps(message).encode()
            }
        }
        # 提交任务到Cloud Tasks队列
        tasks_client.create_task(request={"parent": QUEUE_PATH, "task": task_payload})
        # 立即返回200,确认Pub/Sub消息
        return '', 200
    except Exception as e:
        # 任务提交失败时返回500,让Pub/Sub重试
        return '', 500

# 异步处理任务的接口
@app.route('/process-task', methods=['POST'])
def process_task():
    task_data = request.get_json()
    # 这里编写你的消息处理逻辑
    # 例如:time.sleep(200)
    # ...
    return '', 200

if __name__ == '__main__':
    app.run(port=8080, debug=True)

注意事项

  • 需为Cloud Run服务账号配置Cloud Tasks的创建权限,确保能正常提交任务。
  • 处理接口需实现幂等性:因为如果任务提交成功但Cloud Run返回200前出现异常,Pub/Sub会重试消息,导致重复提交任务,需保证重复执行不会产生副作用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 12:10:05