如何为Spring Integration SFTP Outbound Gateway添加重试机制(附Java示例)
Spring Integration SFTP Outbound Gateway 上传重试实现(Java配置)
核心思路是通过RequestHandlerRetryAdvice配置重试策略,再通过@ServiceActivator的adviceChain属性将其绑定到SFTP网关,实现上传失败后的自动重试。
1. 重试策略配置类
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.handler.advice.RequestHandlerRetryAdvice; import org.springframework.retry.RetryCallback; import org.springframework.retry.RetryContext; import org.springframework.retry.RetryListener; import org.springframework.retry.backoff.FixedBackOffPolicy; import org.springframework.retry.policy.SimpleRetryPolicy; import org.springframework.retry.support.RetryTemplate; @Configuration public class RetryConfig { @Bean public RequestHandlerRetryAdvice sftpUploadRetryAdvice() { RequestHandlerRetryAdvice retryAdvice = new RequestHandlerRetryAdvice(); retryAdvice.setRetryTemplate(sftpRetryTemplate()); // 添加重试监听,记录重试过程日志 retryAdvice.setRetryListeners(new RetryListener() { @Override public <T, E extends Throwable> boolean open(RetryContext context, RetryCallback<T, E> callback) { return true; } @Override public <T, E extends Throwable> void close(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) { if (throwable != null) { System.err.println("SFTP上传重试最终失败,异常:" + throwable.getMessage()); } else { System.out.println("SFTP上传重试成功"); } } @Override public <T, E extends Throwable> void onError(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) { int retryCount = context.getRetryCount(); System.err.println("SFTP上传第" + retryCount + "次失败,异常:" + throwable.getMessage()); } }); return retryAdvice; } private RetryTemplate sftpRetryTemplate() { RetryTemplate retryTemplate = new RetryTemplate(); // 重试策略:最多重试3次,指定触发重试的异常类型 SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); retryPolicy.setMaxAttempts(3); retryPolicy.setRetryableExceptions( java.io.IOException.class, org.springframework.integration.sftp.session.SftpSessionException.class ); retryTemplate.setRetryPolicy(retryPolicy); // 退避策略:每次重试间隔1秒 FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy(); backOffPolicy.setBackOffPeriod(1000); retryTemplate.setBackOffPolicy(backOffPolicy); return retryTemplate; } }
2. SFTP网关配置类
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.annotation.ServiceActivator; import org.springframework.integration.sftp.outbound.SftpOutboundGateway; import org.springframework.integration.sftp.session.DefaultSftpSessionFactory; import org.springframework.messaging.MessageHandler; @Configuration public class SftpConfig { @Bean public DefaultSftpSessionFactory sftpSessionFactory() { DefaultSftpSessionFactory sessionFactory = new DefaultSftpSessionFactory(); sessionFactory.setHost("your-sftp-host"); sessionFactory.setPort(22); sessionFactory.setUser("your-username"); sessionFactory.setPassword("your-password"); // 允许未知主机密钥(生产环境建议配置可信密钥) sessionFactory.setAllowUnknownKeys(true); return sessionFactory; } @Bean @ServiceActivator(inputChannel = "sftpUploadChannel", adviceChain = "sftpUploadRetryAdvice") public MessageHandler sftpOutboundGateway() { // "put"表示上传操作,第二个参数是远程目录路径(支持SpEL表达式) SftpOutboundGateway gateway = new SftpOutboundGateway(sftpSessionFactory(), "put", "/remote/upload/path"); // 设置文件存在时的处理策略:覆盖已存在文件 gateway.setFileExistsMode(org.springframework.integration.file.FileExistsMode.REPLACE); return gateway; } }
关键说明
@ServiceActivator的adviceChain属性直接指定重试增强的bean名称,即可将重试逻辑绑定到SFTP网关。- 可根据业务需求调整重试次数、退避间隔,以及需要触发重试的异常类型。
- 重试监听用于记录重试过程,便于排查上传失败问题。
内容的提问来源于stack exchange,提问作者Choff
相关产品推荐
相关产品推荐

