启用消息排序键后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()确认发布成功),但仍出现重复,可能的原因有这些:
发布端重试机制误触发
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 )订阅端消息确认问题
重复消息也可能是订阅者未正确确认消息导致的。如果订阅者处理消息超时,或未调用ack()确认,Pub/Sub会认为消息未被处理并重新投递。需要检查订阅端代码,确保处理完消息后及时调用ack(),并根据业务需求调整ack超时时间(默认10分钟)。排序键流积压问题
启用排序键后,每个排序键对应一个有序消息流,如果同一个排序键的消息发布过于频繁,可能导致客户端内部队列积压进而触发重试。可以分散排序键粒度,比如在原有排序键基础上添加时间戳或批次号,降低单个排序键的流压力。
关于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

