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

如何在Google Pub/Sub中为结果添加消息ID(附Python代码优化)

解决方案

你当前的代码存在几个关键问题导致无法正确收集成功/失败的消息ID:

  1. 键名拼写错误:succeded应为succeeded
  2. 处理失败future时调用future.result()会直接抛出异常,导致代码崩溃
  3. 发布失败时不会生成有效消息ID,需要关联原始UUID进行记录

以下是修正后的完整代码:

from typing import List, Callable
import pubsub_v1
from concurrent import futures
import logging

logger = logging.getLogger(__name__)

GOOGLE_CLOUD_PROJECT_ID = "your-project-id"
GOOGLE_CLOUD_TOPIC_ID = "your-topic-id"

def publish_messages_with_error_handler(project_id: str = GOOGLE_CLOUD_PROJECT_ID,
                                        topic_id: str = GOOGLE_CLOUD_TOPIC_ID,
                                        data: List[str] = []) -> dict:
    """Publishes multiple messages to a Pub/Sub topic with an error handler."""

    publisher = pubsub_v1.PublisherClient()
    topic_path = publisher.topic_path(project_id, topic_id)
    publish_futures = []

    # 修正拼写错误:succeded -> succeeded
    result = {
        "succeeded": [],
        "failed": []
    }

    def get_callback(publish_future: pubsub_v1.publisher.futures.Future,
                     data: str) -> Callable[[pubsub_v1.publisher.futures.Future], None]:
        def callback(publish_future: pubsub_v1.publisher.futures.Future) -> None:
            try:
                message_id = publish_future.result(timeout=0)
                logger.info(f"Published message {data} with ID: {message_id}")
            except futures.TimeoutError:
                logger.error(f"Publishing message {data} timed out")
            except Exception as e:
                logger.error(f"Failed to publish message {data}: {str(e)}")

        return callback

    if data:
        for message in data:
            publish_future = publisher.publish(topic_path, message.encode("utf-8"))
            publish_future.add_done_callback(get_callback(publish_future, message))
            publish_futures.append(publish_future)

    futures.wait(publish_futures, return_when=futures.ALL_COMPLETED)

    print(f"Finished publishing messages to {topic_path}.")

    # 正确收集成功/失败记录
    for future, message in zip(publish_futures, data):
        try:
            # 成功时,result()返回Pub/Sub生成的消息ID
            message_id = future.result()
            result["succeeded"].append(message_id)
        except Exception as e:
            # 失败时无有效消息ID,记录原始UUID及错误信息
            result["failed"].append({
                "uuid": message,
                "error": str(e)
            })

    return result

关键修改说明:

  • 拼写修正:将succeded改为标准拼写succeeded,避免后续键名匹配错误
  • 安全的结果收集:通过try-except捕获所有异常,替代原代码中future.exception()的判断方式,避免调用result()时崩溃
  • 关联原始UUID:失败时将对应的原始UUID和错误信息绑定记录,便于后续排查问题(若仅需UUID,可简化为result["failed"].append(message))
  • 线程安全:结果收集逻辑统一放在主线程执行,避免多线程修改字典的线程安全问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 03:10:48