【问题标题】:RequestHandlerRetryAdvice cannot be made to work with Ftp.outboundGateway in Spring IntegrationRequestHandlerRetryAdvice 不能与 Spring Integration 中的 Ftp.outboundGateway 一起使用
【发布时间】:2018-09-27 02:05:16
【问题描述】:

我的情况类似于this SO question中描述的情况。不同之处在于我不使用WebFlux.outboundGateway,而是使用Ftp.outboundGateway,我在其上调用AbstractRemoteFileOutboundGateway.Command.GETcommand,常见的问题是我无法使用定义的RequestHandlerRetryAdvice。

配置如下(精简到相关部分):

@RestController
@RequestMapping( value = "/somepath" )
public class DownloadController
{
   private DownloadGateway downloadGateway;

   public DownloadController( DownloadGateway downloadGateway )
   {
      this.downloadGateway = downloadGateway;
   }

   @PostMapping( "/downloads" )
   public void download( @RequestParam( "filename" ) String filename )
   {
      Map<String, Object> headers = new HashMap<>();

      downloadGateway.triggerDownload( filename, headers );
   }
}    
@MessagingGateway
public interface DownloadGateway
{
   @Gateway( requestChannel = "downloadFiles.input" )
   void triggerDownload( Object value, Map<String, Object> headers );
}
@Configuration
@EnableIntegration
public class FtpDefinition
{
   private FtpProperties ftpProperties;

   public FtpDefinition( FtpProperties ftpProperties )
   {
      this.ftpProperties = ftpProperties;
   }

   @Bean
   public DirectChannel gatewayDownloadsOutputChannel()
   {
      return new DirectChannel();
   }

   @Bean
   public IntegrationFlow downloadFiles( RemoteFileOutboundGatewaySpec<FTPFile, FtpOutboundGatewaySpec> getRemoteFile )
   {
      return f -> f.handle( getRemoteFile, getRetryAdvice() )
                   .channel( "gatewayDownloadsOutputChannel" );
   }

   private Consumer<GenericEndpointSpec<AbstractRemoteFileOutboundGateway<FTPFile>>> getRetryAdvice()
   {
      return e -> e.advice( ( (Supplier<RequestHandlerRetryAdvice>) () -> {
         RequestHandlerRetryAdvice advice = new RequestHandlerRetryAdvice();
         advice.setRetryTemplate( getRetryTemplate() );
         return advice;
      } ).get() );
   }

   private RetryTemplate getRetryTemplate()
   {
      RetryTemplate result = new RetryTemplate();

      FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
      backOffPolicy.setBackOffPeriod( 5000 );

      result.setBackOffPolicy( backOffPolicy );
      return result;
   }

   @Bean
   public RemoteFileOutboundGatewaySpec<FTPFile, FtpOutboundGatewaySpec> getRemoteFile( SessionFactory sessionFactory )
   {
      return 
         Ftp.outboundGateway( sessionFactory,
                              AbstractRemoteFileOutboundGateway.Command.GET,
                              "payload" )
            .fileExistsMode( FileExistsMode.REPLACE )
            .localDirectoryExpression( "'" + ftpProperties.getLocalDir() + "'" )
            .autoCreateLocalDirectory( true );
   }

   @Bean
   public SessionFactory<FTPFile> ftpSessionFactory()
   {
      DefaultFtpSessionFactory sessionFactory = new DefaultFtpSessionFactory();
      sessionFactory.setHost( ftpProperties.getServers().get( 0 ).getHost() );
      sessionFactory.setPort( ftpProperties.getServers().get( 0 ).getPort() );
      sessionFactory.setUsername( ftpProperties.getServers().get( 0 ).getUser() );
      sessionFactory.setPassword( ftpProperties.getServers().get( 0 ).getPassword() );
      return sessionFactory;
   }
}
@SpringBootApplication
@EnableIntegration
@IntegrationComponentScan
public class FtpTestApplication {

    public static void main(String[] args) {
        SpringApplication.run( FtpTestApplication.class, args );
    }
}
@Configuration
@PropertySource( "classpath:ftp.properties" )
@ConfigurationProperties( prefix = "ftp" )
@Data
public class FtpProperties
{
   @NotNull
   private String localDir;

   @NotNull
   private List<Server> servers;

   @Data
   public static class Server
   {
      @NotNull
      private String host;

      @NotNull
      private int port;

      @NotNull
      private String user;

      @NotNull
      private String password;
   }
}

控制器主要用于测试目的,在实际实现中有一个轮询器。我的FtpProperties 持有一个服务器列表,因为在实际实现中,我使用DelegatingSessionFactory 来根据一些参数选择一个实例。

根据Gary Russell's comment,我预计会重试失败的下载。但是,如果我中断下载服务器端(通过在 FileZilla 实例中发出“Kick user”),我只会立即获得堆栈跟踪并且不会重试:

org.apache.commons.net.ftp.FTPConnectionClosedException: FTP response 421 received.  Server closed connection.
[...]

我还需要上传文件,为此我使用Ftp.outboundAdapter。在这种情况下,使用相同的RetryTemplate,如果我中断上传服务器端,Spring Integration 会再执行两次尝试,每次延迟 5 秒,然后才记录java.net.SocketException: Connection reset,一切都符合预期。

我尝试调试了一下,发现在第一次尝试通过Ftp.outboundAdapter 上传之前,RequestHandlerRetryAdvice.doInvoke() 上的断点被命中。但是当通过Ftp.outboundGateway 下载时,该断点从不命中。

我的配置有问题吗,有人可以让RequestHandlerRetryAdvice 与Ftp.outboundGateway/AbstractRemoteFileOutboundGateway.Command.GET 一起工作吗?

【问题讨论】:

    标签: spring-boot spring-integration spring-integration-dsl


    【解决方案1】:

    抱歉耽搁了;本周我们在 SpringOne 平台。

    问题是由于网关规范是一个 bean ——网关最终在应用建议之前被初始化。

    我这样修改了你的代码...

    @Bean
    public IntegrationFlow downloadFiles(SessionFactory<FTPFile> sessionFactory) {
        return f -> f.handle(getRemoteFile(sessionFactory), getRetryAdvice())
                .channel("gatewayDownloadsOutputChannel");
    }
    
    ...
    
    private RemoteFileOutboundGatewaySpec<FTPFile, FtpOutboundGatewaySpec> getRemoteFile(SessionFactory<FTPFile> sessionFactory) {
        return Ftp.outboundGateway(sessionFactory,
                AbstractRemoteFileOutboundGateway.Command.GET,
                "payload")
                .fileExistsMode(FileExistsMode.REPLACE)
                .localDirectoryExpression("'/tmp'")
                .autoCreateLocalDirectory(true);
    }
    

    ...它奏效了。

    通常最好不要直接处理 Specs,而是将它们内联到流定义中...

    @Bean
    public IntegrationFlow downloadFiles(SessionFactory<FTPFile> sessionFactory) {
        return f -> f.handle(Ftp.outboundGateway(sessionFactory,
                AbstractRemoteFileOutboundGateway.Command.GET,
                "payload")
                .fileExistsMode(FileExistsMode.REPLACE)
                .localDirectoryExpression("'/tmp'")
                .autoCreateLocalDirectory(true), getRetryAdvice())
            .channel("gatewayDownloadsOutputChannel");
    }
    

    【讨论】:

    • 太棒了,你拯救了我的一天,非常感谢!作为 Spring 的新手,尤其是 Spring 集成的新手,我永远不会找到原因和解决方案。
    猜你喜欢
    • 1970-01-01
    • 2018-08-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-12-25
    • 1970-01-01
    • 1970-01-01
    • 2014-01-18
    相关资源
    最近更新 更多