Spring Integration:确保同文件CONTENT消息先于EMPTY消息发至SFTP
解决Spring Integration中SFTP文件发送顺序问题
针对你遇到的同一文件的EMPTY消息先于CONTENT消息到达SFTP的问题,这里提供几种直接可行的解决方案:
方案1:利用SFTP发送成功回执触发EMPTY消息(推荐)
这种方式基于Spring Integration的事件驱动模型,只有当CONTENT文件成功发送到SFTP后,才会自动触发EMPTY文件的发送,完全保证顺序,同时避免同步等待带来的性能损耗。
步骤1:调整MyHandler逻辑,只生成CONTENT消息并携带EMPTY元数据
修改MyHandler的@ServiceActivator方法,不再返回两个消息,而是仅生成CONTENT消息,并在消息头中携带对应EMPTY文件的文件名信息:
@ServiceActivator public Message<?> handle(Resource resource) { String baseFileName = resource.getFilename().replace("_CONTENT", ""); return MessageBuilder.withPayload(resource) .setHeader("fileName", baseFileName + "_CONTENT") // 携带EMPTY文件的文件名,用于后续生成空文件 .setHeader("emptyFileName", baseFileName + "_EMPTY") .build(); }
步骤2:配置SFTP发送流和成功回执处理流
拆分流程为两个部分:处理CONTENT发送、监听发送成功事件并处理EMPTY发送:
// CONTENT文件发送流 @Bean public IntegrationFlow sftpContentFlow() { return IntegrationFlows.from("test") .enrichHeaders(blah) .transform(myTransformer) .handle(myHandler) .handle(Sftp.outboundAdapter(sftpSessionFactory) .remoteDirectory("/your/remote/path") .fileName(m -> m.getHeaders().get("fileName")) // 指定CONTENT发送成功后的回执通道 .sendSuccessChannel("sftpSendSuccessChannel")) .get(); } // EMPTY文件发送流(由CONTENT发送成功事件触发) @Bean public IntegrationFlow sftpEmptyFlow() { return IntegrationFlows.from("sftpSendSuccessChannel") // 根据回执中的元数据生成0字节空文件 .transform(replyMsg -> { String emptyFileName = replyMsg.getHeaders().get("emptyFileName", String.class); return new ByteArrayResource(new byte[0]) { @Override public String getFilename() { return emptyFileName; } }; }) .handle(Sftp.outboundAdapter(sftpSessionFactory) .remoteDirectory("/your/remote/path") .fileName(m -> ((Resource) m.getPayload()).getFilename())) .get(); }
方案2:同步等待CONTENT发送完成后再发送EMPTY
如果需要保持原有流程结构,可以通过网关同步调用SFTP发送,等待CONTENT发送成功后再生成EMPTY消息。
步骤1:定义SFTP发送网关
创建一个网关接口,用于同步发送SFTP消息并获取发送结果:
@MessagingGateway public interface SftpSyncGateway { @Gateway(requestChannel = "sftpSyncSendChannel", replyTimeout = 15000) Boolean send(Message<?> message); }
步骤2:配置SFTP同步发送流
@Bean public IntegrationFlow sftpSyncSendFlow() { return IntegrationFlows.from("sftpSyncSendChannel") .handle(Sftp.outboundAdapter(sftpSessionFactory) .remoteDirectory("/your/remote/path") .fileName(m -> m.getHeaders().get("fileName"))) // 发送成功返回true,失败返回false或抛出异常 .transform(m -> true) .get(); }
步骤3:修改MyHandler逻辑,同步发送CONTENT后再生成EMPTY
注入SftpSyncGateway,先发送CONTENT并等待成功,再构建EMPTY消息:
@ServiceActivator public List<Message<?>> handle(Resource resource) { List<Message<?>> result = new ArrayList<>(); String baseFileName = resource.getFilename().replace("_CONTENT", ""); // 构建并发送CONTENT消息 Message<?> contentMsg = MessageBuilder.withPayload(resource) .setHeader("fileName", baseFileName + "_CONTENT") .build(); Boolean contentSentSuccess = sftpSyncGateway.send(contentMsg); if (Boolean.TRUE.equals(contentSentSuccess)) { // 构建EMPTY消息 Resource emptyResource = new ByteArrayResource(new byte[0]) { @Override public String getFilename() { return baseFileName + "_EMPTY"; } }; Message<?> emptyMsg = MessageBuilder.withPayload(emptyResource) .setHeader("fileName", baseFileName + "_EMPTY") .build(); result.add(emptyMsg); } return result; }
步骤4:调整主流程
保持原有主流程结构,将MyHandler返回的EMPTY消息路由到SFTP即可:
@Bean public IntegrationFlow mainFlow() { return IntegrationFlows.from("test") .enrichHeaders(blah) .transform(myTransformer) .handle(myHandler) .route(myRouter) .get(); }
内容的提问来源于stack exchange,提问作者innealtóir
相关产品推荐
相关产品推荐

