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
相关产品推荐
相关产品推荐

