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')"
解决方案
- 将入站通道适配器定义为PublishSubscribeChannel:
@Bean public MessageChannel streamChannel() { return new PublishSubscribeChannel(); }
- 添加优先级为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; }
- 额外:若需移动处理失败的文件(来自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
相关产品推荐
相关产品推荐

