Spring Integration并行调用报错:接收线程已收到回复
Spring Integration高并发场景下ConcurrentModificationException与Reply重复接收问题解决
问题背景
我搭建了一个Spring Integration流,通过网关从控制器调用该流,微服务中4个API均使用此配置。单请求调用功能正常,但运行大量并行的自动化测试时,出现错误提示Reply message received but the receiving thread has already received a reply,同时伴随ConcurrentModificationException异常。
异常栈信息
pool-19-thread-1 ERROR com.fasterxml.jackson.databind.JsonMappingException: (was java.util.ConcurrentModificationException) (through reference chain: org.apache.logging.log4j.core.layout.AbstractJacksonLayout$LogEventWithAdditionalFields["logEvent"]->org.apache.logging.log4j.core.impl.Log4jLogEvent["message"]) com.fasterxml.jackson.databind.JsonMappingException: (was java.util.ConcurrentModificationException) (through reference chain: org.apache.logging.log4j.core.layout.AbstractJacksonLayout$LogEventWithAdditionalFields["logEvent"]->org.apache.logging.log4j.core.impl.Log4jLogEvent["message"]) at com.fasterxml.jackson.databind.JsonMappingException.wrapWithPath(JsonMappingException.java:397) at com.fasterxml.jackson.databind.JsonMappingException.wrapWithPath(JsonMappingException.java:356) at com.fasterxml.jackson.databind.ser.std.StdSerializer.wrapAndThrow(StdSerializer.java:316) at com.fasterxml.jackson.databind.ser.std.BeanSerializerBase.serializeFieldsFiltered(BeanSerializerBase.java:815) at com.fasterxml.jackson.databind.ser.impl.UnwrappingBeanSerializer.serialize(UnwrappingBeanSerializer.java:132) at com.fasterxml.jackson.databind.ser.impl.UnwrappingBeanPropertyWriter.serializeAsField(UnwrappingBeanPropertyWriter.java:127) at com.fasterxml.jackson.databind.ser.std.BeanSerializerBase.serializeFields(BeanSerializerBase.java:755) at com.fasterxml.jackson.databind.ser.BeanSerializer.serialize(BeanSerializer.java:178) at com.fasterxml.jackson.databind.ser.DefaultSerializerProvider._serialize(DefaultSerializerProvider.java:480) at com.fasterxml.jackson.databind.ser.DefaultSerializerProvider.serializeValue(DefaultSerializerProvider.java:319) at com.fasterxml.jackson.databind.ObjectWriter$Prefetch.serialize(ObjectWriter.java:1516) at com.fasterxml.jackson.databind.ObjectWriter._writeValueAndClose(ObjectWriter.java:1217) at com.fasterxml.jackson.databind.ObjectWriter.writeValue(ObjectWriter.java:1059) at org.apache.logging.log4j.core.layout.AbstractJacksonLayout.toSerializable(AbstractJacksonLayout.java:344) at org.apache.logging.log4j.core.layout.JsonLayout.toSerializable(JsonLayout.java:292) at org.apache.logging.log4j.core.layout.AbstractJacksonLayout.toSerializable(AbstractJacksonLayout.java:292) at org.apache.logging.log4j.core.layout.JsonLayout.toSerializable(JsonLayout.java:70) at org.apache.logging.log4j.core.layout.AbstractJacksonLayout.toSerializable(AbstractJacksonLayout.java:52) at org.apache.logging.log4j.core.layout.AbstractStringLayout.toByteArray(AbstractStringLayout.java:282) at org.apache.logging.log4j.core.layout.AbstractLayout.encode(AbstractLayout.java:209) at org.apache.logging.log4j.core.layout.AbstractLayout.encode(AbstractLayout.java:37) at org.apache.logging.log4j.core.appender.AbstractOutputStreamAppender.directEncodeEvent(AbstractOutputStreamAppender.java:197) at org.apache.logging.log4j.core.appender.AbstractOutputStreamAppender.tryAppend(AbstractOutputStreamAppender.java:190) at org.apache.logging.log4j.core.appender.AbstractOutputStreamAppender.append(AbstractOutputStreamAppender.java:181) at org.apache.logging.log4j.core.config.AppenderControl.tryCallAppender(AppenderControl.java:161) at org.apache.logging.log4j.core.config.AppenderControl.callAppender0(AppenderControl.java:134) at org.apache.logging.log4j.core.config.AppenderControl.callAppenderPreventRecursion(AppenderControl.java:125) at org.apache.logging.log4j.core.config.AppenderControl.callAppender(AppenderControl.java:89) at org.apache.logging.log4j.core.config.LoggerConfig.callAppenders(LoggerConfig.java:542) at org.apache.logging.log4j.core.config.LoggerConfig.processLogEvent(LoggerConfig.java:500) at org.apache.logging.log4j.core.config.LoggerConfig.log(LoggerConfig.java:483) at org.apache.logging.log4j.core.config.LoggerConfig.log(LoggerConfig.java:417) at org.apache.logging.log4j.core.config.AwaitCompletionReliabilityStrategy.log(AwaitCompletionReliabilityStrategy.java:82) at org.apache.logging.log4j.core.Logger.log(Logger.java:161) at org.apache.logging.log4j.spi.AbstractLogger.tryLogMessage(AbstractLogger.java:2205) at org.apache.logging.log4j.spi.AbstractLogger.logMessageTrackRecursion(AbstractLogger.java:2159) at org.apache.logging.log4j.spi.AbstractLogger.logMessageSafely(AbstractLogger.java:2142) at org.apache.logging.log4j.spi.AbstractLogger.logMessage(AbstractLogger.java:1994) at org.apache.logging.log4j.spi.AbstractLogger.logIfEnabled(AbstractLogger.java:1852) at org.apache.commons.logging.LogAdapter$Log4jLog.log(LogAdapter.java:270) at org.apache.commons.logging.LogAdapter$Log4jLog.info(LogAdapter.java:230) at org.springframework.core.log.LogAccessor.info(LogAccessor.java:292) at org.springframework.integration.handler.LoggingHandler.handleMessageInternal(LoggingHandler.java:191) at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:56) at org.springframework.integration.channel.FixedSubscriberChannel.send(FixedSubscriberChannel.java:77) at org.springframework.integration.channel.interceptor.WireTap.preSend(WireTap.java:166) at org.springframework.integration.channel.AbstractMessageChannel$ChannelInterceptorList.preSend(AbstractMessageChannel.java:469) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:309) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:272) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47) at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109) at org.springframework.integration.handler.AbstractMessageProducingHandler.sendOutput(AbstractMessageProducingHandler.java:457) at org.springframework.integration.handler.AbstractMessageProducingHandler.doProduceOutput(AbstractMessageProducingHandler.java:325) at org.springframework.integration.handler.AbstractMessageProducingHandler.produceOutput(AbstractMessageProducingHandler.java:268) at org.springframework.integration.handler.AbstractMessageProducingHandler.sendOutputs(AbstractMessageProducingHandler.java:232) at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:142) at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:56) at org.springframework.integration.dispatcher.AbstractDispatcher.tryOptimizedDispatch(AbstractDispatcher.java:115) at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:133) at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106) at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:317) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:272) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47) at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109) at org.springframework.integration.handler.AbstractMessageProducingHandler.sendOutput(AbstractMessageProducingHandler.java:457) at org.springframework.integration.handler.AbstractMessageProducingHandler.doProduceOutput(AbstractMessageProducingHandler.java:325) at org.springframework.integration.handler.AbstractMessageProducingHandler.produceOutput(AbstractMessageProducingHandler.java:268) at org.springframework.integration.handler.AbstractMessageProducingHandler.sendOutputs(AbstractMessageProducingHandler.java:232) at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:142) at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:56) at org.springframework.integration.dispatcher.AbstractDispatcher.tryOptimizedDispatch(AbstractDispatcher.java:115) at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:133) at org.springframework.integration.dispatcher.UnicastingDispatcher.access$000(UnicastingDispatcher.java:55) at org.springframework.integration.dispatcher.UnicastingDispatcher$1.run(UnicastingDispatcher.java:114) at org.springframework.integration.util.ErrorHandlingTaskExecutor.lambda$execute$0(ErrorHandlingTaskExecutor.java:57) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.util.ConcurrentModificationException at java.base/java.util.LinkedHashMap$LinkedHashIterator.nextNode(LinkedHashMap.java:719) at java.base/java.util.LinkedHashMap$LinkedEntryIterator.next(LinkedHashMap.java:751) at java.base/java.util.LinkedHashMap$LinkedEntryIterator.next(LinkedHashMap.java:749) at java.base/java.util.AbstractMap.toString(AbstractMap.java:551) at java.base/java.lang.String.valueOf(String.java:2951) at java.base/java.lang.StringBuilder.append(StringBuilder.java:172) at java.base/java.util.AbstractMap.toString(AbstractMap.java:556) at java.base/java.lang.String.valueOf(String.java:2951) at java.base/java.lang.StringBuilder.append(StringBuilder.java:172) at java.base/java.util.AbstractMap.toString(AbstractMap.java:556) at org.springframework.messaging.MessageHeaders.toString(MessageHeaders.java:348) at java.base/java.lang.String.valueOf(String.java:2951) at java.base/java.lang.StringBuilder.append(StringBuilder.java:172) at org.springframework.messaging.support.GenericMessage.toString(GenericMessage.java:120) at java.base/java.lang.String.valueOf(String.java:2951) at java.base/java.util.Objects.toString(Objects.java:160) at org.springframework.integration.handler.LoggingHandler.createLogMessage(LoggingHandler.java:209) at org.springframework.integration.handler.LoggingHandler.lambda$handleMessageInternal$0(LoggingHandler.java:179) at org.springframework.core.log.LogMessage$SupplierMessage.buildString(LogMessage.java:155) at org.springframework.core.log.LogMessage.toString(LogMessage.java:70) at java.base/java.lang.String.valueOf(String.java:2951) at org.apache.logging.log4j.message.ObjectMessage.getFormattedMessage(ObjectMessage.java:55) at org.apache.logging.log4j.core.jackson.MessageSerializer.serialize(MessageSerializer.java:44) at org.apache.logging.log4j.core.jackson.MessageSerializer.serialize(MessageSerializer.java:33) at com.fasterxml.jackson.databind.ser.BeanPropertyWriter.serializeAsField(BeanPropertyWriter.java:728) at com.fasterxml.jackson.databind.ser.impl.SimpleBeanPropertyFilter.serializeAsField(SimpleBeanPropertyFilter.java:208) at com.fasterxml.jackson.databind.ser.std.BeanSerializerBase.serializeFieldsFiltered(BeanSerializerBase.java:807) ... 79 more
原Spring Integration配置
@Bean public IntegrationFlow myFlow() { return flow -> flow.handle(validatorService, "validateRequest") .channel(c -> c.executor(Executors.newCachedThreadPool())) .log("Parallel Processing started") .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(savingRequestToTheDB()) .recipientFlow(flow1()) .recipientFlow(flow2()) .recipientFlow(conditionToCallFlow3, flow3()) .recipientFlow( conditionToCallFlow4, flow4()), gatherer -> gatherer.outputProcessor(gatherAndProcess())) .log("Parallel Processing ended") .gateway(getSomethng(), f -> f.errorChannel("cdErrorChannel")) .split() .channel(c -> c.executor(Executors.newCachedThreadPool())) .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(prepareResponse()) .recipientFlow(conditionToCallFlow5, flow5()), gatherer -> gatherer.outputProcessor(gatherLionResponse())) .to(savingLionResponse()); }
gatherLionResponse方法
private MessageGroupProcessor gatherLionResponse() { return lionsService::gatherLionResponse; } public Object gatherLionResponse(MessageGroup group) { String lionResponse; MessageHeaders headers; String JStage ; String JName ; long dbID; Optional<Message<?>> lionResponseMessage = group.getMessages().stream() .filter(m -> m.getHeaders().get("sequenceNumber", Integer.class) == 1) .findFirst(); if (lionResponseMessage.isPresent()) { Message<?> message = lionResponseMessage.get(); lionResponse= message.getPayload().toString(); headers = message.getHeaders(); } logger.atInfo().log("Lion-Response-Body : " + lionResponse); return MessageBuilder.withPayload(lionResponse) .copyHeaders(headers) .removeHeaders( "Date", CommonConstants.D_FLAG, CommonConstants.CREATED_AT, CommonConstants.RE_QUANT); }
尝试修改后的配置
@Bean public IntegrationFlow myFlow() { return flow -> flow.handle(validatorService, "validateRequest") .channel(c -> c.executor(Executors.newCachedThreadPool())) .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(savingLoanRequestToTheDB()) .recipientFlow(flow1()) .recipientFlow(flow2()) .recipientFlow(conditionToCallFlow3, flow3()) .recipientFlow( conditionToCallFlow4, flow4()), gatherer -> gatherer.outputProcessor(gatherAndProcess())) .gateway(getSomethng()) .publishSubscribeChannel( pubSub1 -> pubSub1 .subscribe( f -> f.publishSubscribeChannel( Executors.newCachedThreadPool(), pubSub2 -> pubSub2 .subscribe(flow5()) .subscribe(prepareResponse()))) .subscribe(f -> f.to(saveResponse()))); }
问题分析与修复方案
1. ConcurrentModificationException根源与修复
从异常栈可见,问题出在MessageHeaders.toString()序列化时触发了LinkedHashMap的并发修改。MessageHeaders内部用LinkedHashMap存储,本身线程安全,但如果流程中直接修改原消息的headers(比如gatherLionResponse中的removeHeaders操作),同时日志线程在序列化该headers,就会触发异常。
修复代码:
public Object gatherLionResponse(MessageGroup group) { String lionResponse = null; MessageHeaders headers = null; Optional<Message<?>> lionResponseMessage = group.getMessages().stream() .filter(m -> m.getHeaders().get("sequenceNumber", Integer.class) == 1) .findFirst(); if (lionResponseMessage.isPresent()) { Message<?> message = lionResponseMessage.get(); lionResponse = message.getPayload().toString(); // 创建headers副本,避免修改原headers导致并发冲突 Map<String, Object> headerMap = new HashMap<>(message.getHeaders()); headerMap.remove("Date"); headerMap.remove(CommonConstants.D_FLAG); headerMap.remove(CommonConstants.CREATED_AT); headerMap.remove(CommonConstants.RE_QUANT); headers = new MessageHeaders(headerMap);
相关产品推荐
相关产品推荐

