【问题标题】:Spring Integration - Handling stale sftp sessionsSpring Integration - 处理陈旧的 sftp 会话
【发布时间】:2018-03-01 21:26:27
【问题描述】:

我已经实现了以下场景:

  1. 一个queueChannel以字节[]的形式保存消息
  2. MessageHandler,轮询队列通道并通过 sftp 上传文件
  3. 一个 Transformer,监听 errorChannel 并将从失败消息中提取的有效负载发送回 queueChannel(被认为是处理失败消息的错误处理程序,因此不会丢失任何内容)

如果 sftp 服务器在线,一切正常。

如果 sftp 服务器关闭,则作为转换器到达的错误消息是:

org.springframework.messaging.MessagingException: Failed to obtain pooled item; nested exception is java.lang.IllegalStateException: failed to create SFTP Session

转换器对此无能为力,因为有效负载的 failedMessage 为 null 并且本身会引发异常。转换器丢失了消息。

如何配置我的流程以使转换器获得正确的消息以及未成功上传文件的相应负载?

我的配置:

  @Bean
  public MessageChannel toSftpChannel() {
    final QueueChannel channel = new QueueChannel();
    channel.setLoggingEnabled(true);
    return new QueueChannel();
  }

  @Bean
  public MessageChannel toSplitter() {
    return new PublishSubscribeChannel();
  }

  @Bean
  @ServiceActivator(inputChannel = "toSftpChannel", poller = @Poller(fixedDelay = "10000", maxMessagesPerPoll = "1"))
  public MessageHandler handler() {
    final SftpMessageHandler handler = new SftpMessageHandler(sftpSessionFactory());
    handler.setRemoteDirectoryExpression(new LiteralExpression(sftpRemoteDirectory));
    handler.setFileNameGenerator(message -> {
      if (message.getPayload() instanceof byte[]) {
        return (String) message.getHeaders().get("name");
      } else {
        throw new IllegalArgumentException("byte[] expected in Payload!");
      }
    });
    return handler;
  }

  @Bean
  public SessionFactory<LsEntry> sftpSessionFactory() {
    final DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true);

    final Properties jschProps = new Properties();
    jschProps.put("StrictHostKeyChecking", "no");
    jschProps.put("PreferredAuthentications", "publickey,password");
    factory.setSessionConfig(jschProps);

    factory.setHost(sftpHost);
    factory.setPort(sftpPort);
    factory.setUser(sftpUser);
    if (sftpPrivateKey != null) {
      factory.setPrivateKey(sftpPrivateKey);
      factory.setPrivateKeyPassphrase(sftpPrivateKeyPassphrase);
    } else {
      factory.setPassword(sftpPasword);
    }
    factory.setAllowUnknownKeys(true);
    return new CachingSessionFactory<>(factory);
  }

  @Bean
  @Splitter(inputChannel = "toSplitter")
  public DmsDocumentMessageSplitter splitter() {
    final DmsDocumentMessageSplitter splitter = new DmsDocumentMessageSplitter();
    splitter.setOutputChannelName("toSftpChannel");
    return splitter;
  }

  @Transformer(inputChannel = "errorChannel", outputChannel = "toSftpChannel")
  public Message<?> errorChannelHandler(ErrorMessage errorMessage) throws RuntimeException {

    Message<?> failedMessage = ((MessagingException) errorMessage.getPayload())
      .getFailedMessage();
    return MessageBuilder.withPayload(failedMessage)
                         .copyHeadersIfAbsent(failedMessage.getHeaders())
                         .build();
  }

  @MessagingGateway 
  public interface UploadGateway {

    @Gateway(requestChannel = "toSplitter")
    void upload(@Payload List<byte[]> payload, @Header("header") DmsDocumentUploadRequestHeader header);
  }

谢谢..

更新

@Bean(PollerMetadata.DEFAULT_POLLER)
@Transactional(propagation = Propagation.REQUIRED, isolation = Isolation.READ_COMMITTED)
  PollerMetadata poller() {
    return Pollers
      .fixedRate(5000)
      .maxMessagesPerPoll(1)
      .receiveTimeout(500)
      .taskExecutor(taskExecutor())
      .transactionSynchronizationFactory(transactionSynchronizationFactory())
      .get();
  }

  @Bean
  @ServiceActivator(inputChannel = "toMessageStore", poller = @Poller(PollerMetadata.DEFAULT_POLLER))
  public BridgeHandler bridge() {
    BridgeHandler bridgeHandler = new BridgeHandler();
    bridgeHandler.setOutputChannelName("toSftpChannel");
    return bridgeHandler;
  }

【问题讨论】:

  • 抱歉,我以为您在轮询入站通道适配器。打开 DEBUG 日志记录并遵循消息流。如果您仍然无法弄清楚,请将日志发布到某个地方。
  • 如果您担心消息丢失,您可能也不应该使用QueueChannel。
  • @gary 你好加里,感谢您的回复...我主要担心上传失败会导致文件丢失..据我了解,我必须以某种方式/某处将文件排队才能能够重试上传......或者有没有更好的方法来处理这些用例?这是日志link
  • 请看我的回答。

标签: spring-integration jsch spring-integration-sftp


【解决方案1】:

null failedMessage 是一个错误;转载INT-4421。

我不建议在这种情况下使用QueueChannel。如果您使用直接渠道,您可以配置retry advice 以尝试重新投递。当重试次数用尽时(如果这样配置),异常将被抛回调用线程。

将建议添加到SftpMessageHandler 的adviceChain 属性。

编辑

您可以通过在可轮询通道和 sftp 适配器之间插入桥接来解决“丢失”失败消息:

@Bean
@ServiceActivator(inputChannel = "toSftpChannel", poller = @Poller(fixedDelay = "5000", maxMessagesPerPoll = "1"))
public BridgeHandler bridge() {
    BridgeHandler bridgeHandler = new BridgeHandler();
    bridgeHandler.setOutputChannelName("toRealSftpChannel");
    return bridgeHandler;
}

@Bean
@ServiceActivator(inputChannel = "toRealSftpChannel")
public MessageHandler handler() {
    final SftpMessageHandler handler = new SftpMessageHandler(sftpSessionFactory());
    handler.setRemoteDirectoryExpression(new LiteralExpression("foo"));
    handler.setFileNameGenerator(message -> {
        if (message.getPayload() instanceof byte[]) {
            return (String) message.getHeaders().get("name");
        }
        else {
            throw new IllegalArgumentException("byte[] expected in Payload!");
        }
    });
    return handler;
}

【讨论】:

  • 您可以使用QueueChannel,只要您还使用事务性ChannelMessageStore,例如JdbcChannelMessageStore。在这种情况下,当重试用尽后抛出异常时,事务将回滚,消息将保留在存储中。
  • @ArtemBilan:您能否提供一些示例,如何使用带有 Java 注释的 RedisChannelPriorityMessageStore 设置 QueueChannel?
  • @RokPurkeljc 有一个带有 mongo 商店的 java 配置示例 in the documentation。另请参阅我对null failedMessage 的修改。
  • Redis 不是事务性的,所以不保证在异常的情况下不会丢失消息
  • 尝试使用 Oracle 数据源设置 JdbcChannelMessageStore。如何在我的自定义(使用 PollerMetadata)定义的带有 Java 注释的轮询器上启用 @Transactional(请参阅更新)?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-04-04
  • 2011-03-27
  • 1970-01-01
  • 2012-05-31
  • 1970-01-01
相关资源
最近更新 更多