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

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))

问题根因分析

  1. 释放策略逻辑错误:g.size() > 25意味着只有当组内消息数超过25时才触发释放,导致前25条消息需要等待第26条消息进入才会被释放,后续不足25的消息则无法触发释放条件。
  2. 关联策略缺失:两次拆分后,消息的默认correlationId来自原始拆分上下文,导致聚合器无法将所有文件消息归入同一个全局组,部分消息被分散到不同组中无法累积到25条。
  3. 序列头干扰:默认情况下,聚合器会根据拆分后消息的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 03:02:01