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

Apache Pulsar新订阅ACK时抛出recycled already异常求助

问题:Apache Pulsar 确认消息时抛出 "recycled already" 异常

我们的Apache Pulsar已运行数周,新增了一个无分区的持久化队列。突然出现如下异常,仅该新订阅存在此问题,其余队列均无异常。所有消费者的实现方式一致,异常发生在执行消息确认操作consumer.acknowledgeAsync(msg);时。

异常信息

java.lang.IllegalStateException: recycled already
at io.netty.util.Recycler$WeakOrderQueue.transfer(Recycler.java:442)
at io.netty.util.Recycler$Stack.scavengeSome(Recycler.java:587)
at io.netty.util.Recycler$Stack.scavenge(Recycler.java:562)
at io.netty.util.Recycler$Stack.pop(Recycler.java:535)
at io.netty.util.Recycler.get(Recycler.java:162)
at org.apache.pulsar.common.util.collections.ConcurrentBitSetRecyclable.create(ConcurrentBitSetRecyclable.java:51)
at org.apache.pulsar.client.impl.PersistentAcknowledgmentsGroupingTracker.lambda$doIndividualBatchAckAsync$8(PersistentAcknowledgmentsGroupingTracker.java:328)
at java.base/java.util.concurrent.ConcurrentHashMap.computeIfAbsent(ConcurrentHashMap.java:1705)
at org.apache.pulsar.client.impl.PersistentAcknowledgmentsGroupingTracker.doIndividualBatchAckAsync(PersistentAcknowledgmentsGroupingTracker.java:232)
at org.apache.pulsar.client.impl.PersistentAcknowledgmentsGroupingTracker.doIndividualBatchAck(PersistentAcknowledgmentsGroupingTracker.java:294)
at org.apache.pulsar.client.impl.PersistentAcknowledgmentsGroupingTracker.doIndividualBatchAck(PersistentAcknowledgmentsGroupingTracker.java:287)
at org.apache.pulsar.client.impl.PersistentAcknowledgmentsGroupingTracker.lambda$addAcknowledgment$4(PersistentAcknowledgmentsGroupingTracker.java:235)
at org.apache.pulsar.client.impl.PersistentAcknowledgmentsGroupingTracker.addIndividualAcknowledgment(PersistentAcknowledgmentsGroupingTracker.java:220)
at org.apache.pulsar.client.impl.PersistentAcknowledgmentsGroupingTracker.addAcknowledgment(PersistentAcknowledgmentsGroupingTracker.java:232)
at org.apache.pulsar.client.impl.PersistentAcknowledgmentsGroupingTracker.addAcknowledgment(PersistentAcknowledgmentsGroupingTracker.java:196)
at org.apache.pulsar.client.impl.ConsumerImpl.doAcknowledge(ConsumerImpl.java:560)
at org.apache.pulsar.client.impl.ConsumerBase.doAcknowledgeWithTxn(ConsumerBase.java:688)
at org.apache.pulsar.client.impl.ConsumerBase.acknowledgeAsync(ConsumerBase.java:634)
at org.apache.pulsar.client.impl.ConsumerBase.acknowledgeAsync(ConsumerBase.java:619)
at org.apache.pulsar.client.impl.ConsumerBase.acknowledgeAsync(ConsumerBase.java:521)
at de.seepex.service.lastactive.PulsarLastActiveConsumer.messageListener(PulsarLastActiveConsumer.java:82)
at org.apache.pulsar.client.impl.ConsumerBase.callMessageListener(ConsumerBase.java:1157)

消费者代码

private void messageListener(Consumer<byte[]> consumer, Message<byte[]> msg) {
    try {
        byte[] payload = msg.getData();

        if (payload == null || payload.length == 0) {
            consumer.acknowledgeAsync(msg);
            return;
        }

        final LastActiveUpdate lastActiveUpdate = objectMapper.readValue(payload, LastActiveUpdate.class);
        lastActiveProcessor.process(lastActiveUpdate);

        consumer.acknowledgeAsync(msg);

    } catch (Exception e) {
        LOG.error("failed to process", e);
        consumer.negativeAcknowledge(msg);
    }
}

异常原因分析

这个异常来自Netty的Recycler组件,说明Pulsar客户端复用的ConcurrentBitSetRecyclable对象被重复回收,本质是并发场景下的对象复用竞态冲突。结合你的场景,可能的触发点:

  • 批量确认的并发冲突:Pulsar默认开启ACK分组批量确认,当这个无分区队列的消息量极大时,ACK分组的并发处理逻辑可能出现竞态,导致同一个可回收BitSet对象被多次回收。
  • 客户端版本BUG:部分Pulsar客户端版本(如2.8.x-2.9.x的部分版本)存在ACK分组逻辑中对象回收的并发BUG,升级到2.10.0及以上稳定版本可解决。
  • 无分区队列的ACK压力集中:无分区队列的所有消息都集中在少量线程处理,ACK操作的并发压力比多分区队列更集中,更容易触发这种边缘场景的回收冲突。

临时解决方案

  • 禁用ACK分组:创建Consumer时配置ackGroupingTimeMillis=0,改为单条消息确认,避免BitSet复用的冲突。
  • 排查重复确认:确保每个消息只被调用一次acknowledgeAsync或negativeAcknowledge(当前代码逻辑无明显重复,但需确认是否有其他隐式调用)。
  • 升级Pulsar客户端到最新稳定版本,修复已知的并发回收BUG。

内容的提问来源于stack exchange,提问作者Alex Tbk

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 06:34:52