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

K8环境下Pubsub模拟器出现消息重复发布及数据损坏问题求助

K8环境下Pubsub模拟器出现消息重复发布及数据损坏问题求助

各位大佬好,我遇到一个挺棘手的问题,折腾了好一阵都没搞定,来这儿求助大家!

先说说我的场景:我有个Go写的服务,会创建K8 Job启动Python Pod处理任务,处理完成后用Pubsub把结果发回给Go服务。目前用的是Pubsub模拟器,所有Pod都在同一个K8命名空间里,属于开发环境。

问题现象

Python端我打印了publish()返回的future.result(),日志显示只发了6条消息:

Published msg: 4
Published msg: 7
Published msg: 10
Published msg: 13
Published msg: 16
Published msg: 19

但Go服务这边的消费日志却收到了12条消息,而且只有msgId和Python端匹配的那些消息数据是完整的,不匹配的消息要么缺失Proto字段,要么自定义属性丢失,数据直接损坏了:

Consumed msg_id: 4, publish_time: 2023-08-10 13:29:02.365 +0000 UTC
Consumed msg_id: 5, publish_time: 2023-08-10 13:29:02.638 +0000 UTC
Consumed msg_id: 7, publish_time: 2023-08-10 13:29:03.167 +0000 UTC
Consumed msg_id: 8, publish_time: 2023-08-10 13:29:03.312 +0000 UTC
Consumed msg_id: 10, publish_time: 2023-08-10 13:29:03.946 +0000 UTC
Consumed msg_id: 11, publish_time: 2023-08-10 13:29:04.107 +0000 UTC
Consumed msg_id: 13, publish_time: 2023-08-10 13:29:04.674 +0000 UTC
Consumed msg_id: 14, publish_time: 2023-08-10 13:29:04.827 +0000 UTC
Consumed msg_id: 16, publish_time: 2023-08-10 13:29:05.408 +0000 UTC
Consumed msg_id: 17, publish_time: 2023-08-10 13:29:05.61 +0000 UTC
Consumed msg_id: 19, publish_time: 2023-08-10 13:29:06.145 +0000 UTC
Consumed msg_id: 20, publish_time: 2023-08-10 13:29:06.294 +0000 UTC

代码相关细节

Python的核心发布逻辑大概是这样(去掉了业务代码):

for page in self._config.pages:
    try:
        data = self._get_results(page)
        if not data:
            continue
        self._publish(data, page, success=True)
    except Exception as e:
        self._publish(success=False)

Publisher Client的初始化配置:

self.publisher: Client = pubsub_v1.PublisherClient(
    publisher_options=pubsub_v1.types.PublisherOptions(enable_message_ordering=True),
    batch_settings=pubsub_v1.types.BatchSettings(
        max_messages=100,
        max_bytes=1000000,
        max_latency=0.01,
    )
)

后来我修改了重试配置(原本用默认参数):

retry_options = retry.Retry(
    initial=100,  # 默认是0.1,我调大了
    maximum=60.0,
    multiplier=1.45,
    predicate=retry.if_exception_type(
        core_exceptions.Aborted,
        core_exceptions.Cancelled,
        core_exceptions.DeadlineExceeded,
        core_exceptions.InternalServerError,
        core_exceptions.ResourceExhausted,
        core_exceptions.ServiceUnavailable,
        core_exceptions.Unknown,
    ),
    deadline=600.0)

publish_future = self.publisher.publish(topic_path, data, retry=retry_options, **attributes)
if callback is not None:
    publish_future.add_done_callback(callback)
    return publish_future
else:
    result = publish_future.result()
    self._logger.log_info(f"Published msg: {result}")

我确认应用层没有手动触发重试发布逻辑,而且Python代码单独在本地跑(用docker-compose起Pubsub模拟器)是完全正常的,发5条收5条,只有放到K8环境里才会出问题。

用的依赖版本:

  • Python 3.11.1
  • grpcio 1.51.1(相关grpc衍生包都是同版本)
  • google-cloud-pubsub 2.16.0
  • google-api-core 2.11.1
  • google-auth 2.22.0

已尝试的排查方案

  • 测试单条消息:依旧出现重复问题
  • 减小每页数据量(怀疑消息大小导致异常):哪怕只发1条数据,问题还是存在
  • 调整重试的initial值(怀疑初始RPC超时引发重试):调大后没有改善
  • 移除所有重试异常(取消自动重试):问题依旧
  • 把Python Pod的CPU从100m调到200m:没有明显变化

更新发现

我试了发送空数据,Go端居然只收到6条正常消息!这说明消息数据大小是导致延迟进而引发重复的原因,但我不知道该怎么针对性解决?是继续调大初始超时时间?还是给Pod分配更多CPU资源?

有没有大佬遇到过类似问题,或者能给我指个调试方向的?感激不尽!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.21 15:48:11