Pulsar 也会重复消费?

简介: 排查了一个问题:在使用 Pulsar 消费时,发生了同一条消息反复消费的情况。

排查


当他告诉我这个现象的时候我就持怀疑态度,根据之前使用的经验 Pulsar 在官方文档以及 API 中都解释过:


网络异常,图片无法展示
|

网络异常,图片无法展示
|


只有当设置了消费的 ackTimeout 并超时消费时才会重复投递消息,默认情况下是关闭的,查看代码也确实没有开启。


那会不会是调用了 negativeAcknowledge() 方法呢(调用该方法也会触发重新投递),因为我们使了一个第三方库 github.com/majusko/pul… 只有当抛出异常时才会调用该方法。


查阅代码之后也没有地方抛出异常,甚至整个过程中都没看到异常产生;这就有点诡异了。


复现


为了捋清楚整个事情的来龙去脉,详细了解了他的使用流程;


其实也就是业务出现了 bug,他在消息消费时 debug 然后进行单步调试,当走完一次调试后,没多久马上又收到了同样的消息。


但奇怪的是也不是每次 debug 后都能重复消费,我们都说如果一个 bug 能 100% 完全复现,那基本上就解决一大半了。


所以我们排查的第一步就是完全复现这个问题。


为了排除掉是 IDEA 的问题(虽然极大概率不太可能)既然是 debug 的时候产生的问题,那其实转换到代码也就是 sleep 嘛,所以我们打算在消费逻辑里直接 sleep 一段时间看能否复现。


经过测试,sleep 几秒到几十秒都无法复现,最后索性 sleep 一分钟,神奇的事情发生了,每次都成功复现!


既然能成功复现那就好说了,因为我自己的业务代码也有使用到 Pulsar 的地方,为了方便调试就准备在自己的项目里再复现一次。


结果诡异的事情再次发生,我这里又不能复现了。


虽然这才是符合预期的,但这就没法调了呀。


本着相信现代科学的前提,我们俩唯一的区别就是项目不一样了,为此我对比了两边的代码。


@PulsarConsumer(
            topic = xx,
            clazz = Xx.class,
            subscriptionType = SubscriptionType.Shared
    )
    public void consume(Data msg) {
        log.info("consume msg:{}", msg.getOrderId());
        Lock lock = redisLockRegistry.obtain(msg.getOrderId());
        if (lock.tryLock()) {
            try {
                orderService.do(msg.getOrderId());
            } catch (Exception e) {
                log.error("consumer msg:{} err:", msg.toString(), e);
            } finally {
                lock.unlock();
            }
        }
    }


结果不出所料,同事那边的代码加了锁;一个基于 Redis 的分布式锁,这时我一拍大腿不会是解锁的时候超时了导致抛了异常吧。


为了验证这个问题,在能复现的基础上我在框架的 Pulsar 消费处打了断点:


网络异常,图片无法展示
|


网络异常,图片无法展示
|


果然破案了,异常提示已经非常清楚了:加锁已经过了超时时间。


进入异常后直接 negative 消息,同时异常也被吃掉了,所以之前没有发现。


网络异常,图片无法展示
|


查阅了 RedisLockRegistry 的源码,默认超时时间正好是一分钟,所以之前我们 sleep 几十秒也无法复现这个问题。


总结


事后我向同事了解了下为啥这里要加锁,因为我看下来完全没有加锁的必要;结果他是因为从别人那里复制的代码才加上的,压根没想那么多。


所以这事也能得出一些教训:


  • ctrl C/V 虽然方便,但也得充分考虑自己的业务场景。


  • 使用一些第三方 API 时,需要充分了解其作用、参数。


相关文章
|
Java 存储 jvm-sandbox
海量流量下,淘宝如何进行稳定的流量回放?
随着业务的不断发展, 整个淘系的服务端已经有数千个应用,在淘宝已经有非常大的应用数量和变更次数的基础上, 对流量回放也有更高的要求。那么在不断尝试流量的录制与回放的过程中,我们遇到了什么问题?那么在不断尝试的过程中,我们遇到了什么问题?我们由从中得到了什么启示?流量录制回放又能给我们带来多少收益?
11099 1
|
9月前
|
人工智能 Rust 运维
这个神器让你白嫖ClaudeOpus 4.5,Gemini 3!还能接Claude Code等任意平台
加我进AI讨论学习群,公众号右下角“联系方式”文末有老金的 开源知识库地址·全免费
11114 22
|
9月前
|
人工智能 API Android开发
送给GLM Coding Plan用户和开源社区的“AI手机”
智谱推出“AI手机”新体验,通过Claude Code输入提示词,即可自动部署开源Agent模型AutoGLM。三步操作,轻松拥有专属AI设备,享受技术平权。倡导开源生态与AI协同,推动人人可用的AGI未来。
975 2
|
3月前
|
数据采集 人工智能 自然语言处理
跨源数据整合为什么不一定要先建大中台
跨源整合不再只是一个工程问题,而是企业是否具备智能分析与数据执行能力的基础问题。
|
8月前
|
人工智能 安全 JavaScript
Claude Code子代理实战:10个即用模板分享
Claude Code单次泛化指令易失效?作者提出“子代理”理念:为AI分配专属角色(如重构专家、测试员、安全审查员),每代理专注一事、规则明确、输出可控。10个实战模板覆盖开发全链路,让AI协作更接近真实工程团队——专注比全能更可靠。
2392 0
Claude Code子代理实战:10个即用模板分享
|
XML JSON Java
Java 反射:从原理到实战的全面解析与应用指南
本文深度解析Java反射机制,从原理到实战应用全覆盖。首先讲解反射的概念与核心原理,包括类加载过程和`Class`对象的作用;接着详细分析反射的核心API用法,如`Class`、`Constructor`、`Method`和`Field`的操作方法;最后通过动态代理和注解驱动配置解析等实战场景,帮助读者掌握反射技术的实际应用。内容翔实,适合希望深入理解Java反射机制的开发者。
1100 13
Future原理解析
介绍了Java多线程中Future类的原理
Future原理解析
|
存储 供应链 安全
区块链技术在选举中的应用:透明与安全的新时代
区块链技术在选举中的应用:透明与安全的新时代
871 16
|
监控 IDE Java
XXL-JOB任务调度详解
XXL-JOB任务调度详解
2708 0
|
Dubbo Java 应用服务中间件
5分钟学会 gRPC(中)
我猜测大部分长期使用 Java 的开发者应该较少会接触 gRPC,毕竟在 Java 圈子里大部分使用的还是 Dubbo/SpringClound 这两类服务框架。 我也是近段时间有机会从零开始重构业务才接触到 gRPC 的,当时选择 gRPC 时也有几个原因

热门文章

最新文章