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

如何让Google Pub/Sub Python订阅脚本按顺序执行而非并行

问题:Google Pub/Sub Python客户端并行处理消息,需改为串行执行

当前使用Google Pub/Sub Python客户端拉取消息时,脚本默认采用多线程并行处理,导致消息未按预期逐条串行执行(5秒等待逻辑未生效,所有消息先被接收再批量处理)。

预期执行流程

Message from Pub/Sub: Message 1
<< Script should wait for 5 Seconds >>
Success
After 5 Seconds
Message from Pub/Sub: Message 2
<< Script should wait for 5 Seconds >>
Success
After 5 Seconds
Message from Pub/Sub: Message 3
<< Script should wait for 5 Seconds >>
Success
After 5 Seconds

实际执行流程

Message from Pub/Sub: Message 1
Message from Pub/Sub: Message 2
Message from Pub/Sub: Message 3
<<Script is waiting for 5 seconds>>
Success
After 5 Seconds
Success
After 5 Seconds
Success
After 5 Seconds

解决方案:限制并发数为1实现串行处理

Google Pub/Sub的subscribe方法默认会启动多个线程并行调用callback,只需设置max_concurrency=1即可强制串行处理。同时补全原脚本缺失的导入模块,并确保处理完成后手动确认消息(避免重复投递)。

修改后的完整代码:

from google.cloud import pubsub_v1
from concurrent.futures import TimeoutError
import time
import json
from google.oauth2 import service_account


def time_sleep():
    time.sleep(5)
    print("Success")

# Manage the Message and Decode it
def callback(message: pubsub_v1.subscriber.message.Message) -> None:
    try:
        # Get the message data
        message_data = message.data.decode("utf-8")

        # Parse the message data as JSON
        json_data = json.loads(message_data)
        print("Message from Pub/Sub: ", json_data)
        time_sleep()
        print("After 5 seconds")
        # 处理完成后确认消息,避免重复投递
        message.ack()
    except Exception as e:
        # Handle any exceptions that may occur
        print(f"Error processing message: {e}")
        # 处理失败时拒绝消息,让Pub/Sub重新投递
        message.nack()

# 替换为你的实际配置
secret_data = {"google_auth": "你的服务账号JSON内容"}
project_id = "你的项目ID"
subscription_id = "你的订阅ID"
timeout = 300  # 按需设置超时时间,单位秒

# Create credentials using the JSON private key
credentials = service_account.Credentials.from_service_account_info(
    secret_data["google_auth"]
)

subscriber = pubsub_v1.SubscriberClient(credentials=credentials)
# The `subscription_path` method creates a fully qualified identifier
# in the form `projects/{project_id}/subscriptions/{subscription_id}`
subscription_path = subscriber.subscription_path(project_id, subscription_id)

# 设置max_concurrency=1,强制串行处理消息
streaming_pull_future = subscriber.subscribe(
    subscription_path, 
    callback=callback,
    max_concurrency=1
)
print(f"Listening for messages on {subscription_path}..\n")

# Wrap subscriber in a 'with' block to automatically call close() when done.
with subscriber:
    try:
        # When `timeout` is not set, result() will block indefinitely,
        # unless an exception is encountered first.
        streaming_pull_future.result(timeout=timeout)
    except TimeoutError:
        streaming_pull_future.cancel()  # Trigger the shutdown.
        streaming_pull_future.result()  # Block until the shutdown is complete.

关键修改说明

  1. 补全缺失导入:添加import time、import json、from google.oauth2 import service_account,解决原脚本的依赖问题。
  2. 设置串行并发:在subscriber.subscribe()中加入max_concurrency=1,限制同时处理的消息数为1,实现逐条串行执行。
  3. 消息确认机制:在callback成功处理后调用message.ack(),确保Pub/Sub移除该消息;失败时调用message.nack(),触发重新投递逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 14:53:13