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

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);
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 11:58:42