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

Spring Integration发送Kafka消息后移动SFTP文件失败排查与解决

问题描述

基于Spring Integration实现从SFTP服务器读取文件流转发送至Kafka,期望成功发送消息后将远程文件移动到同服务器的其他目录,但文件完全未移动且无报错。

尝试用以下代码将upload/test.xml移动到upload/processed/test.xml:

@Bean
@ServiceActivator(inputChannel = "success")
public MessageHandler handler() {
  return new SftpOutboundGateway(sftpSessionFactory(), "mv", "");
}

已在消息头中设置file_renameTo=upload/processed/test.xml,同时想了解是否可通过类似advice.setOnSuccessExpressionString("@template.copy(headers['file_remoteDirectory']+'/'+headers['file_remoteFile'])");的方式实现文件移动。

消息内容如下:

"GenericMessage [payload=Note(to=Toves, from=Jani, heading=Reminder, body=Don't forget me this weekend!!!!!), headers={file_remoteHostPort=localhost:2222, file_remoteFileInfo={"directory":false,"filename":"test.xml","link":false,"modified":1674550080000,"permissions":"rw-r--r--","remoteDirectory":"upload","size":122}, kafka_messageKey=test.xml, file_remoteDirectory=upload, kafka_recordMetadata=test-0@298, file_renameTo=upload/processed/test.xml, id=708d04c4-5abc-9f45-e83b-1fea7ffa5e8d, closeableResource=org.springframework.integration.file.remote.session.CachingSessionFactory$CachedSession@2e0aa05, file_remoteFile=test.xml, timestamp=1674556856288}]"

调试发现报错发生在以下代码处:

private String obtainRemoteFilePath(Message<?> requestMessage) {
Error here----> String remoteFilePath = this.fileNameProcessor.processMessage(requestMessage);
        Assert.state(remoteFilePath != null,
                () -> "The 'fileNameProcessor' evaluated to null 'remoteFilePath' from message: " + requestMessage);
        return remoteFilePath;
}

错误信息:

"class com.demo.sftp.models.Note cannot be cast to class java.lang.String (com.demo.sftp.models.Note is in unnamed module of loader org.springframework.boot.devtools.restart.classloader.RestartClassLoader @387f9ed2; java.lang.String is in module java.base of loader 'bootstrap')"

解决方案
  1. 将入站通道适配器定义为PublishSubscribeChannel:
@Bean
public MessageChannel streamChannel() {
  return new PublishSubscribeChannel();
}
  1. 添加优先级为LOWEST_PRECEDENCE的OutboundGateway:
@Bean
@Order(Ordered.LOWEST_PRECEDENCE)
@ServiceActivator(inputChannel = "streamChannel")
public MessageHandler moveFile() {
  SftpOutboundGateway sftpOutboundGateway = new SftpOutboundGateway(sftpSessionFactory(), Command.MV.getCommand(),
      "headers['file_remoteDirectory'] + '/' + headers['file_remoteFile']");
  sftpOutboundGateway
      .setRenameExpressionString(
          "headers['file_remoteDirectory'] + '/processed/' +headers['timestamp'] + '-' + headers['file_remoteFile']");
  sftpOutboundGateway.setRequiresReply(false);
  sftpOutboundGateway.setUseTemporaryFileName(true);
  sftpOutboundGateway.setOutputChannelName("nullChannel");
  sftpOutboundGateway.setOrder(Ordered.LOWEST_PRECEDENCE);
  sftpOutboundGateway.setAsync(true);
  return sftpOutboundGateway;
}
  1. 额外:若需移动处理失败的文件(来自errorChannel),需修改setRenameExpressionString以适配ErrorMessage结构:
@Bean
@Order(Ordered.LOWEST_PRECEDENCE)
@ServiceActivator(inputChannel = IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME)
public MessageHandler moveErrorFile() {
  SftpOutboundGateway sftpOutboundGateway = new SftpOutboundGateway(sftpSessionFactory(), Command.MV.getCommand(),
      "payload['failedMessage']['headers']['file_remoteDirectory'] + '/' + payload['failedMessage']['headers']['file_remoteFile']");
  sftpOutboundGateway
      .setRenameExpressionString(
          "payload['failedMessage']['headers']['file_remoteDirectory'] + '/error/' + payload['failedMessage']['headers']['timestamp'] + '-' + payload['failedMessage']['headers']['file_remoteFile']");
  sftpOutboundGateway.setRequiresReply(false);
  sftpOutboundGateway.setUseTemporaryFileName(true);
  sftpOutboundGateway.setOutputChannelName("nullChannel");
  sftpOutboundGateway.setOrder(Ordered.HIGHEST_PRECEDENCE);
  sftpOutboundGateway.setAsync(true);
  return sftpOutboundGateway;
}

内容的提问来源于stack exchange,提问作者Wicharn Rueangkhajorn

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 06:56:55