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

