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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 17:06:32