如何让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.
关键修改说明
- 补全缺失导入:添加
import time、import json、from google.oauth2 import service_account,解决原脚本的依赖问题。 - 设置串行并发:在
subscriber.subscribe()中加入max_concurrency=1,限制同时处理的消息数为1,实现逐条串行执行。 - 消息确认机制:在callback成功处理后调用
message.ack(),确保Pub/Sub移除该消息;失败时调用message.nack(),触发重新投递逻辑。
内容的提问来源于stack exchange,提问作者Moses
相关产品推荐
相关产品推荐

