Spring Integration聚合释放策略异常:仅释放前25个文件后续被忽略
目录文件处理的聚合释放策略问题
我正在尝试为目录中的文件处理配置正确的释放策略。在对目录和文件进行两次*拆分(split)后,执行聚合(aggregate)*并设置释放策略:每满25个文件就释放提交。目前前25个文件能正常提交,但之后的所有文件都被完全忽略,集成流程直接结束,剩余文件未被释放。
相关IntegrationFlow代码
@Bean public IntegrationFlow QueueJamsPostJobsByFile() { return IntegrationFlows.from(IntegrationNamesEnum.QUEUE_JAMS_POST_JOBS_BY_FILE.toString()) .handle((p, h) -> { try { /* need to get all the files in this folder */ JsonObject obj = JsonParser.parseString(p.toString()).getAsJsonObject(); fileFilter = obj.get("fileFilter").getAsString(); return obj.get("dirName").getAsString(); } catch (Exception e) { QUEUE_JAMS_POST_JOBS_BY_FILE.info("Exception Queueing File(s) to JAMS " + e.getLocalizedMessage()); responseBean = new EsbResponsePayloadBean(); responseBean.setReturnHttpStatusCode("500"); responseBean.setReturnPayload("Error: " + e.getLocalizedMessage()); responseBean.setReturnStatusCode("-1"); responseBean.setReturnStatusDesc("Exception Queueing File(s) to JAMS " + e.getLocalizedMessage()); return MessageBuilder.withPayload(responseBean.getReturnPayload()) .setHeader("rtn_http_status_code", responseBean.getReturnHttpStatusCode()) .setHeader("rtn_status_code", responseBean.getReturnStatusCode()) .setHeader("x-error-message", "Exception Queueing File(s) to JAMS " + e.getLocalizedMessage()) .setHeader(org.springframework.integration.http.HttpHeaders.STATUS_CODE, responseBean.getReturnHttpStatusCode()) .build(); } }) .gateway(IntegrationNamesEnum.GET_SUBDIRECTORIES_INTGRTN.toString()) /*split for each file. each one will gateway into jams submit integration independently*/ .split() .handle((p, h) -> { return MessageBuilder.withPayload(p) .setHeader("x-esb-file-filter", fileFilter) .setHeader("x-esb-directory-path", p.toString()) .build(); }) .gateway(IntegrationNamesEnum.GET_FILES_INTGRTN.toString()) .split() .filter(Message.class, m -> m.getHeaders().get("x-error-message") == null || m.getHeaders().get("x-error-message").toString().isEmpty(), e -> e.discardChannel("discardChannel")) .handle((p, h) -> { try { String Directory=h.get("x-esb-directory-path").toString(); File f = new File(Directory+p.toString()); String file = f.getAbsolutePath(); f=null; return file; } catch (Exception e) { QUEUE_JAMS_POST_JOBS_BY_FILE.info("Exception creating JAMS Payload" + e.getLocalizedMessage()); responseBean = new EsbResponsePayloadBean(); responseBean.setReturnHttpStatusCode("500"); responseBean.setReturnPayload("Error: " + e.getLocalizedMessage()); responseBean.setReturnStatusCode("-1"); responseBean.setReturnStatusDesc("Exception creating JAMS Payload " + e.getLocalizedMessage()); return MessageBuilder.withPayload(responseBean.getReturnPayload()) .setHeader("rtn_http_status_code", responseBean.getReturnHttpStatusCode()) .setHeader("rtn_status_code", responseBean.getReturnStatusCode()) .setHeader(org.springframework.integration.http.HttpHeaders.STATUS_CODE, responseBean.getReturnHttpStatusCode()) .build(); } }) //Release Strategy Here: .aggregate(a -> a.releaseStrategy(g -> g.size() > 25) .expireTimeout(30000) .groupTimeout(30000) .sendPartialResultOnExpiry(true) .expireGroupsUponTimeout(true)) .handle((p, h) -> { FileListMessageBean fileListMessageBean = new FileListMessageBean(); List<FileListMessageBean.FileName> filesList = new ArrayList<FileListMessageBean.FileName>(); List<String> fileList = (List<String>) p; for(String fileName : fileList) { FileListMessageBean.FileName payloadFile = fileListMessageBean.new FileName(); payloadFile.fileName = fileName; FileProperty fp = fileListMessageBean.new FileProperty(); fp.filePropertyName = "DirectoryToWritePayloads"; fp.filePropertyValue = FilenameUtils.getFullPath(fileName); List<FileListMessageBean.FileProperty> fileProps = new ArrayList<FileProperty>(); fileProps.add(fp); payloadFile.fileProperties = fileProps; filesList.add(payloadFile); } fileListMessageBean.setFiles(filesList); JAMSScheduleJobRequestPayload structureBean = new JAMSScheduleJobRequestPayload(); JAMSScheduleJobRequestPayload.Root structureBeanRoot = structureBean.new Root(); structureBeanRoot.setName("\AlliantData\Installs\Prep\Esb_Post_UID_File"); structureBeanRoot.setFolderName("POST API"); structureBeanRoot.setOverrideName("Esb_Post_UID_File"); List<JAMSScheduleJobRequestPayload.Parameter> parameterList = new ArrayList<>(); JAMSScheduleJobRequestPayload.Parameter ESB_s_Payload = structureBean.new Parameter(); ESB_s_Payload.setParamName("ESB_s_Payload"); ESB_s_Payload.setParamValue(GsonUtil.gson.toJson(fileListMessageBean, FileListMessageBean.class)); parameterList.add(ESB_s_Payload); JAMSScheduleJobRequestPayload.Parameter ESB_s_Topic_Header = structureBean.new Parameter(); ESB_s_Topic_Header.setParamName("ESB_s_Topic_Header"); ESB_s_Topic_Header.setParamValue(IntegrationNamesEnum.POST_UID2_API.toString()); parameterList.add(ESB_s_Topic_Header); structureBeanRoot.setParameters(parameterList); String json = GsonUtil.gsonbuilder.disableHtmlEscaping().create().toJson(structureBeanRoot); return json; }) .handle((p, h) -> { QUEUE_JAMS_POST_JOBS_BY_FILE.info("Payload: ", p.toString()); return p; }) .gateway(IntegrationNamesEnum.JAMS_SUBMIT_JOB_INTGRTN.getChannelName()) /* aggregate the fileListMessageBean files */ .aggregate() .handle((p, h) -> { QUEUE_JAMS_POST_JOBS_BY_FILE.info("Payload: ", p.toString()); return p; }) .handle((p, h) -> { QUEUE_JAMS_POST_JOBS_BY_FILE.info("Process Completed", h, Level.DEBUG); return p; }) .log(Level.DEBUG, m -> "**** Payload: " + m.getPayload()) .log(Level.DEBUG, m -> "**** Headers: " + m.getHeaders()) .logAndReply(); }
当前聚合配置
.aggregate(a -> a.releaseStrategy(g -> g.size() > 25) .expireTimeout(30000) .groupTimeout(30000) .sendPartialResultOnExpiry(true) .expireGroupsUponTimeout(true))
问题根因分析
- 释放策略逻辑错误:
g.size() > 25意味着只有当组内消息数超过25时才触发释放,导致前25条消息需要等待第26条消息进入才会被释放,后续不足25的消息则无法触发释放条件。 - 关联策略缺失:两次拆分后,消息的默认
correlationId来自原始拆分上下文,导致聚合器无法将所有文件消息归入同一个全局组,部分消息被分散到不同组中无法累积到25条。 - 序列头干扰:默认情况下,聚合器会根据拆分后消息的
sequenceNumber和sequenceSize判断组是否完成,当所有拆分消息处理完毕后,聚合器会认为组已完成,不会触发超时释放剩余消息。
修复方案
调整聚合器配置,解决上述三个核心问题:
完整修正后的聚合配置
.aggregate(a -> a.releaseStrategy(g -> g.size() >= 25) .correlationStrategy(m -> "global-file-group") .ignoreSequenceHeaders(true) .groupTimeout(30000) .sendPartialResultOnExpiry(true) .expireGroupsUponTimeout(true) .expireTimeout(30000))
关键配置说明
- 修正释放策略:将
g.size() >25改为g.size() >=25,确保组内消息数达到25时立即释放。 - 全局关联策略:通过
correlationStrategy(m -> "global-file-group")强制所有文件消息进入同一个聚合组,保证消息累积的连续性。 - 忽略序列头:
ignoreSequenceHeaders(true)禁用基于拆分序列的组完成判断,避免聚合器提前结束组生命周期,确保剩余消息能通过超时触发释放。 - 超时配置:保留
groupTimeout和sendPartialResultOnExpiry,确保当剩余消息不足25时,能在30秒超时后自动释放。
额外检查项
- 后续的无配置
aggregate()会使用默认策略,需根据实际需求补充关联和释放配置,避免消息异常聚合。 - 检查
discardChannel的处理逻辑,确认被过滤的消息未丢失。
内容的提问来源于stack exchange,提问作者Jim Dannucci
相关产品推荐
相关产品推荐

