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

迁移至Jakarta API客户端后ActiveMQ Artemis频繁报AMQ224016错误

ActiveMQ Artemis AMQ224016错误排查求助

问题现象

将ActiveMQ Artemis从2.20.0升级到2.32.0版本后,搭配对应客户端(2.20或2.32版本客户端连接2.32服务端),在Spring Boot从2.7迁移至3.x(同步升级jakarta.jms-api从2.0.3到3.1.0并切换Jakarta包导入)后,出现AMQ224016错误。

每条已处理消息都会触发一次堆栈跟踪输出,导致服务器日志每日超过500MB,而此前仅约100KB。

旧版本组合(Java 8 + artemis-jms-client:2.19.1、Java 17 + Spring Boot 2.7 + artemis-jms-client:2.20.0)均无此问题。测试发现,仅切换Spring Boot版本到3.x并使用Jakarta包时,服务端才会出现异常,监听实现未做任何修改。

尝试过不同Artemis实例配置:默认单实例、2主2从集群、默认/自定义队列,所有场景下旧JMS API正常,新JMS API均触发错误。已知Artemis从2.23.0开始支持JMS 3.1.0(JakartaEE 10),求排查方向或兼容性问题分析。

复现条件

  • ActiveMQ Artemis 2.32.0默认配置实例(客户端通过TCP 61616连接)
  • Spring Boot 3.x + Jakarta JMS 3.x API的客户端监听应用

服务端日志堆栈示例

2024-03-18 12:08:46,069 ERROR [org.apache.activemq.artemis.core.server] AMQ224016: Caught exception
org.apache.activemq.artemis.api.core.ActiveMQIllegalStateException: AMQ229027: Could not find reference on consumer ID=0, messageId = 25612055904 queue = example.queue
    at org.apache.activemq.artemis.core.server.impl.ServerConsumerImpl.individualAcknowledge(ServerConsumerImpl.java:1013) ~[artemis-server-2.32.0.jar:2.32.0]
    at org.apache.activemq.artemis.core.server.impl.ServerSessionImpl.individualAcknowledge(ServerSessionImpl.java:1314) ~[artemis-server-2.32.0.jar:2.32.0]
    at org.apache.activemq.artemis.core.protocol.core.ServerSessionPacketHandler.slowPacketHandler(ServerSessionPacketHandler.java:618) ~[artemis-server-2.32.0.jar:2.32.0]
    at org.apache.activemq.artemis.core.protocol.core.ServerSessionPacketHandler.onMessagePacket(ServerSessionPacketHandler.java:319) ~[artemis-server-2.32.0.jar:2.32.0]
    at org.apache.activemq.artemis.utils.actors.Actor.doTask(Actor.java:32) ~[artemis-commons-2.32.0.jar:2.32.0]
    at org.apache.activemq.artemis.utils.actors.ProcessorBase.executePendingTasks(ProcessorBase.java:68) ~[artemis-commons-2.32.0.jar:2.32.0]
    at org.apache.activemq.artemis.utils.actors.OrderedExecutor.doTask(OrderedExecutor.java:57) ~[artemis-commons-2.32.0.jar:2.32.0]
    at org.apache.activemq.artemis.utils.actors.OrderedExecutor.doTask(OrderedExecutor.java:32) ~[artemis-commons-2.32.0.jar:2.32.0]
    at org.apache.activemq.artemis.utils.actors.ProcessorBase.executePendingTasks(ProcessorBase.java:68) ~[artemis-commons-2.32.0.jar:2.32.0]
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) [?:?]
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) [?:?]
    at org.apache.activemq.artemis.utils.ActiveMQThreadFactory$1.run(ActiveMQThreadFactory.java:118) [artemis-commons-2.32.0.jar:2.32.0]

代码实现

ExampleMqListener.kt

@Component
class ExampleMqListener : MqListener<ExampleMessage> {
    override fun handleMessage(message: ExampleMessage): Boolean {
        return runBlocking {
            LoggerFactory.getLogger(ExampleMqListener::class.java).info("Incoming ExampleMessage: $message")
            true
        }
    }
}

data class ExampleMessage(val data: String) : MqMessage

class MqListenerImpl<T : MqMessage>(private val messageConverter: MessageConverter, private val handler: MqListener<T>) : MessageListener {
    override fun onMessage(jmsMessage: Message) {
        val message = messageConverter.fromMessage(jmsMessage) as? T ?: error("Unexpected message type")
        runCatching {
            handler.handleMessage(message)
        }.getOrElse {
            LoggerFactory.getLogger(MqListenerImpl::class.java).error("Unable to handle message", it)
        }
        jmsMessage.acknowledge()
    }
}

interface MqListener<T> {
    fun handleMessage(message: T): Boolean
}

MqListenerConfig.kt

@Configuration
class MqListenerConfig(
    override val listenerContainerFactory: DefaultJmsListenerContainerFactory,
    override val jmsListenerEndpointRegistry: JmsListenerEndpointRegistry,
    override val messageConverter: MessageConverter,
    private val mqListeners: List<MqListener<out MqMessage>>,
) : JmsListenerConfig {
    @EventListener(ApplicationReadyEvent::class)
    fun startListeners() {
        mqListeners.forEach {
            listenToQueue(it)
        }
    }

    private fun <T : MqMessage> listenToQueue(handler: MqListener<T>) {
        handler.messageType().let {
            listen(
                listenerConfig = ListenerConfig(it.simpleName, it.simpleName, "example.queue", "1-2"),
                handler = handler,
                messageClass = it
            )
        }
    }

    private fun <T : MqMessage> MqListener<T>.messageType(): Class<T> = (javaClass.genericInterfaces[0] as ParameterizedType).actualTypeArguments[0] as Class<T>
}

@JsonIgnoreProperties(value = ["type"])
interface MqMessage : Serializable {
    val type: String
        get() = javaClass.simpleName
}

class ListenerConfig(
    val messageType: String,
    val id: String,
    val destination: String,
    val concurrency: String
)

interface JmsListenerConfig {
    val listenerContainerFactory: DefaultJmsListenerContainerFactory
    val jmsListenerEndpointRegistry: JmsListenerEndpointRegistry
    val messageConverter: MessageConverter
}

fun <T : MqMessage> JmsListenerConfig.listen(
    listenerConfig: ListenerConfig,
    handler: MqListener<T>,
    messageClass: Class<T>
) = SimpleJmsListenerEndpoint().apply {
    id = listenerConfig.id
    destination = listenerConfig.destination
    concurrency = listenerConfig.concurrency
    selector = "_type = '${messageClass.canonicalName}'"
    messageListener = MqListenerImpl(messageConverter, handler)
}.let {
    jmsListenerEndpointRegistry.registerListenerContainer(it, listenerContainerFactory, true)
}

MqConnectionConfig.kt

@Configuration
@ConfigurationProperties(prefix = "mq")
@EnableJms
class MqConnectionConfig {
    var front: MqConfig = MqConfig()

    @Bean
    fun backConnectionFactory(): SingleConnectionFactory =
        SingleConnectionFactory(createFactory(front)).apply {
            setReconnectOnException(true)
        }

    private fun createFactory(config: MqConfig) =
        ActiveMQConnectionFactory(config.url).apply {
            user = config.user
            password = config.password
        }

    @Bean
    fun jmsListenerContainerFactory(): DefaultJmsListenerContainerFactory = DefaultJmsListenerContainerFactory().apply {
        setConnectionFactory(backConnectionFactory())
        setMessageConverter(jacksonJmsMessageConverter())
        setSessionAcknowledgeMode(ActiveMQJMSConstants.INDIVIDUAL_ACKNOWLEDGE)
    }

    @Bean
    fun jmsListenerEndpointRegistry(): JmsListenerEndpointRegistry = JmsListenerEndpointRegistry()

    @Bean
    fun jacksonJmsMessageConverter() = MappingJackson2MessageConverter().apply {
        setTargetType(MessageType.TEXT)
        setTypeIdPropertyName("_type")
        setObjectMapper(jacksonObjectMapper().apply {
            configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false)
            registerModule(JavaTimeModule())
            disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS)
            disable(DeserializationFeature.ADJUST_DATES_TO_CONTEXT_TIME_ZONE)
        })
    }
}

class MqConfig {
    var url: String = "tcp://localhost:61616"
    var user: String? = null
    var password: String? = null
}

依赖配置(build.gradle.kts)

implementation("org.springframework.boot:spring-boot-starter-artemis")
implementation("com.fasterxml.jackson.module:jackson-module-kotlin")
implementation("com.fasterxml.jackson.datatype:jackson-datatype-jsr310")
implementation("org.jetbrains.kotlin:kotlin-reflect")
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-reactor")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 16:17:03