Spring Boot中JMS Listener遇Avro解码错误后无法读取Solace EMS队列消息问题
问题根因
- 核心问题1:
schema和datumReader是类级别的共享变量,多线程消费场景下,前一个线程处理不兼容的Avro消息时会污染这两个对象的内部状态,后续即使处理正常消息,拿到的也是被破坏的变量,直接抛出异常,这就是删除坏消息后依旧无法消费、必须重启的核心原因。 - 核心问题2:异常处理存在疏漏:读取BytesMessage抛出JMSException后,
receivedMsg会为null,后续直接传入SeekableByteArrayInput会触发空指针相关的数组越界异常;且未捕获的异常直接抛到Spring监听容器,你未配置ErrorHandler,会导致消费线程状态异常。 - 核心问题3:默认的自动消息确认机制下,坏消息处理失败后不会被确认,也没有转入死信队列的逻辑,会导致Solace流持续卡住在异常消息位点,出现日志里的乱序消息忽略提示。
修复方案
1. 调整变量作用域,缓存Schema
将schema、datumReader改为方法局部变量,Schema只在初始化时解析一次,避免重复解析和线程安全问题。
2. 补全异常分支拦截,避免异常外抛
所有错误分支处理完直接终止当前消息的处理逻辑,不要让异常抛到Spring容器层。
3. 配置手动确认和死信逻辑
消费成功后手动确认消息,解码失败的消息直接转入死信队列,避免卡住正常消费流。
4. 配置Spring监听容器的ErrorHandler
兜底处理未捕获的异常,避免消费线程异常退出。
修改后的代码示例:
// 提前缓存Schema,只解析一次 private final Schema customerSchema = new Schema.Parser().parse(ConsumerConfiguration.class.getResourceAsStream("/avro/input/customer.avsc")); @Override public void configureJmsListeners(JmsListenerEndpointRegistrar registrar) { // 给容器配置ErrorHandler,兜底处理异常 DefaultJmsListenerContainerFactory factory = (DefaultJmsListenerContainerFactory) registrar.getContainerFactory(); factory.setErrorHandler(e -> logger.error("JMS监听器未捕获异常", e)); factory.setSessionAcknowledgeMode(Session.CLIENT_ACKNOWLEDGE); // 开启手动确认 for (String queue : inQueues) { SimpleJmsListenerEndpoint endpoint = new SimpleJmsListenerEndpoint(); endpoint.setId("myJmsEndpoint-" + queue); endpoint.setDestination(queue); endpoint.setMessageListener(message -> { if (!(message instanceof BytesMessage)) { logger.warn("收到非二进制格式消息,无法处理"); try { message.acknowledge(); // 非预期格式消息直接确认,避免卡队列 } catch (JMSException e) { logger.error("消息确认失败", e); } return; } BytesMessage bytemsg = (BytesMessage) message; byte[] receivedMsg; try { logger.info("收到BytesMessage,CorrelationId:{}", bytemsg.getJMSCorrelationID()); receivedMsg = new byte[(int) bytemsg.getBodyLength()]; bytemsg.readBytes(receivedMsg); } catch (JMSException ex) { logger.error("读取字节消息失败", ex); // 读取失败的消息可以转入死信队列,此处直接确认避免卡队列 try { message.acknowledge(); } catch (JMSException e) { logger.error("消息确认失败", e); } return; } try (SeekableByteArrayInput inputStream = new SeekableByteArrayInput(receivedMsg)) { BinaryDecoder decoder = DecoderFactory.get().binaryDecoder(inputStream, null); List<Customer> recordList = new ArrayList<>(); // 每次新建datumReader,避免状态污染 SpecificDatumReader<Customer> datumReader = new SpecificDatumReader<>(customerSchema); while (!decoder.isEnd()) { Customer avroRecordRead = datumReader.read(null, decoder); recordList.add(avroRecordRead); } logger.debug("收到Customer数据:{}", recordList); generalTopicWriter.sendAvroCustomerBinaryMessagetoTopic(recordList, queuetoTopicMap.get(queue)); // 处理成功手动确认 message.acknowledge(); } catch (IOException e) { logger.error("Avro消息解码失败,消息格式不兼容", e); // 解码失败的消息转入死信队列,此处示例直接确认,你可以根据需求调整死信逻辑 try { message.acknowledge(); } catch (JMSException ex) { logger.error("消息确认失败", ex); } } catch (JMSException e) { logger.error("消息确认失败", e); } }); registrar.registerEndpoint(endpoint); logger.info("队列{}的监听端点注册完成", queue); } }
内容的提问来源于stack exchange,提问作者Chetana Rathore
相关产品推荐
相关产品推荐

