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

启用消息排序键后Google Cloud Pub/Sub Python发布者出现大量重复消息问题

启用消息排序键后Google Cloud Pub/Sub Python发布者出现大量重复消息问题

我来梳理下问题场景并分享排查思路:

我们团队在使用Google Cloud Pub/Sub结合Python的pubsub_v1.PublisherClient时,已经在订阅端和发布端都启用了消息排序键,但启用后频繁遇到重复消息,重复率高达30%-50%,而且这些重复消息的message ID完全一致。另外还有个小疑问:我们从未手动指定过消息ID,但接收的消息里同时出现了message_id和messageId两个字段,搞不清原因。

先贴出核心代码供参考:

from google.cloud import pubsub_v1
import dataclasses
import json
import typing


class GCPubsubPublisher:
    def __init__(self, publisher_client: pubsub_v1.PublisherClient):
        self.publisher_client = publisher_client

    def serialize_data(self, data: typing.Any) -> bytes:
        if data is None:
            raise ValueError("data cannot be None")
        try:
            if dataclasses.is_dataclass(data):
                data = dataclasses.asdict(data)
            if not isinstance(data, str):
                data = json.dumps(data)

            return data.encode("utf-8")  # type: ignore[no-any-return]
        except Exception as exc:
            raise Exception(f"{data} cannot be serialized, Error: {exc}") from exc

    def publish(
        self,
        topic: str,
        data: typing.Any,
        async_wait_for_result: bool,
        ordering_key: str = "",
        **kwargs: typing.Dict[str, typing.Any],
    ) -> str | None:
        serialized_data = self.serialize_data(data)
        future = self.publisher_client.publish(
            topic=topic, data=serialized_data, ordering_key=ordering_key, **kwargs
        )

        # Calling future.result() waits asynchronously until the message has been published successfully
        # and returns the generated message_id
        return None if not async_wait_for_result else future.result()

    @staticmethod
    def build() -> "GCPubsubPublisher":
        publisher_options = pubsub_v1.types.PublisherOptions(
            enable_message_ordering=True, timeout=300
        )
        publisher = pubsub_v1.PublisherClient(publisher_options=publisher_options)
        return GCPubsubPublisher(publisher)

    @staticmethod
    def get_topic_path(project_name: str, topic_name: str) -> str:
        return pubsub_v1.PublisherClient.topic_path(project_name, topic_name)


if __name__ == "__main__":
    pubsub_client = GCPubsubPublisher.build()

    ordering_key1 = "50030469_bill_of_sales"
    form_payload1 = {
        "application_id": "50030469",
        "form_path": "gs://form_path/bill_of_sales.pdf",
        "timestamp": "2025-01-31T19:24:22+0000",
        "form_hash": "9191b4f4ea85d3a2be0f6c2f1cf513a5a40e301bc6306fdca49ee8c88144b696",
    }

    future = pubsub_client.publish(
        topic="projects/project_name/topics/form_submission",
        data=form_payload1,
        async_wait_for_result=True,
        ordering_key=ordering_key1,
    )

    ordering_key2 = "50030469_csc"
    form_payload2 = {
        "application_id": "50030469",
        "form_path": "gs://form_path/csc.pdf",
        "timestamp": "2025-01-31T19:24:22+0000",
        "form_hash": "9191b4f4ea85d3a2be0f6c2f1cf513a5a40e301bc6306fdca49ee8c88144b696",
    }

    future = pubsub_client.publish(
        topic="projects/project_name/topics/form_submission",
        data=form_payload2,
        async_wait_for_result=True,
        ordering_key=ordering_key2,
    )

关于重复消息的排查方向

虽然你设置了async_wait_for_result=True(等待future.result()确认发布成功),但仍出现重复,可能的原因有这些:

  1. 发布端重试机制误触发
    Pub/Sub Python客户端默认会对发布失败进行重试,即使调用了future.result(),如果网络波动导致服务端确认延迟,客户端可能误以为发布失败触发重试,而实际上第一次发布已经成功。可以自定义重试策略,避免不必要的重试:

    from google.api_core import retry
    
    # 自定义重试策略:仅重试特定错误,限制最大重试次数
    custom_retry = retry.Retry(
        retry.if_exception_type(pubsub_v1.exceptions.PubsubError),
        max_attempts=3,
    )
    
    # 在publish时传入自定义重试策略
    future = self.publisher_client.publish(
        topic=topic, data=serialized_data, ordering_key=ordering_key, retry=custom_retry, **kwargs
    )
    
  2. 订阅端消息确认问题
    重复消息也可能是订阅者未正确确认消息导致的。如果订阅者处理消息超时,或未调用ack()确认,Pub/Sub会认为消息未被处理并重新投递。需要检查订阅端代码,确保处理完消息后及时调用ack(),并根据业务需求调整ack超时时间(默认10分钟)。

  3. 排序键流积压问题
    启用排序键后,每个排序键对应一个有序消息流,如果同一个排序键的消息发布过于频繁,可能导致客户端内部队列积压进而触发重试。可以分散排序键粒度,比如在原有排序键基础上添加时间戳或批次号,降低单个排序键的流压力。

关于message_id和messageId的疑问

这是正常现象,无需担心:

  • message_id是Python SDK为符合PEP8命名规范做的字段映射;
  • messageId是Pub/Sub REST API原生的字段名;
    两者指向同一个消息ID,只是不同层面的命名方式,并非代码问题。

额外排查建议

建议在发布端添加详细日志,记录每次发布的参数、返回的message ID以及future.result()的执行情况,明确是客户端重复调用publish还是服务端重复投递:

import logging
logging.basicConfig(level=logging.INFO)

# 在publish方法中添加日志
def publish(...):
    serialized_data = self.serialize_data(data)
    logging.info(f"准备发布消息,排序键:{ordering_key},数据哈希:{hash(serialized_data)}")
    future = self.publisher_client.publish(...)
    if async_wait_for_result:
        try:
            msg_id = future.result()
            logging.info(f"消息发布成功,ID:{msg_id}")
            return msg_id
        except Exception as e:
            logging.error(f"消息发布失败,错误:{str(e)}")
            raise
    return None

如果日志中每个message ID仅出现一次,但订阅端仍收到重复,大概率是服务端投递问题,可联系GCP技术支持协助排查。

备注:内容来源于stack exchange,提问作者bioniclebeastmaster

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 08:53:09