如何用Spring Cloud Stream+RabbitMQ实现SFTP供应商拆分器(批量生产者)?
解决方案
核心思路
基于Spring Cloud Stream的函数式编程模型,结合SFTP Supplier完成文件读取,通过Apache Commons Compress处理解压,解析定长文件生成JSON后,以Flux流式输出单个消息,适配RabbitMQ绑定器的要求,无需切换到Spring Integration。
步骤1:依赖配置
添加必要的Maven/Gradle依赖:
spring-cloud-starter-stream-rabbit:RabbitMQ绑定器spring-cloud-starter-stream-sftp:SFTP数据源支持org.apache.commons:commons-compress:处理gzip/tar解压com.fasterxml.jackson.core:jackson-databind:JSON序列化
步骤2:应用配置(application.yml)
spring: cloud: stream: bindings: sftpSupplier-out-0: destination: your-target-queue # RabbitMQ队列名称 producer: batch-mode: false # 禁用批量模式,确保单消息发送 sftp: suppliers: sftpSupplier: remote-dir: /sftp/remote/path # SFTP待读取文件目录 local-dir: /tmp/sftp-local # 本地临时目录(可选) auto-create-local-dir: true poller: fixed-delay: 30000 # 轮询间隔(30秒) session: host: your-sftp-host port: 22 username: sftp-username password: sftp-password # 若使用密钥认证,替换为以下配置 # private-key: file:/path/to/private/key # passphrase: key-passphrase
步骤3:核心业务代码
import org.apache.commons.compress.archivers.tar.TarArchiveEntry; import org.apache.commons.compress.archivers.tar.TarArchiveInputStream; import org.apache.commons.compress.compressors.gzip.GzipCompressorInputStream; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.sftp.session.SftpFileInfo; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import com.fasterxml.jackson.databind.ObjectMapper; import java.io.*; import java.util.HashMap; import java.util.Map; import java.util.function.Function; @Configuration public class SftpFileProcessingConfig { private final ObjectMapper objectMapper = new ObjectMapper(); @Bean public Function<Flux<Message<SftpFileInfo>>, Flux<Message<String>>> sftpSupplier() { return sftpMessageFlux -> sftpMessageFlux.flatMap(this::processSingleSftpFile); } private Flux<Message<String>> processSingleSftpFile(Message<SftpFileInfo> sftpFileMsg) { SftpFileInfo fileInfo = sftpFileMsg.getPayload(); return Mono.fromCallable(() -> { // 读取SFTP文件流并解压 try (InputStream sftpIn = fileInfo.getSession().readRaw(fileInfo.getRemoteDirectory() + "/" + fileInfo.getFilename()); GzipCompressorInputStream gzipIn = new GzipCompressorInputStream(sftpIn); TarArchiveInputStream tarIn = new TarArchiveInputStream(gzipIn)) { String layoutContent = null; String dataContent = null; TarArchiveEntry entry; // 遍历tar包中的文件,提取layout和data内容 while ((entry = tarIn.getNextTarEntry()) != null) { if (!entry.isDirectory()) { ByteArrayOutputStream baos = new ByteArrayOutputStream(); byte[] buffer = new byte[1024]; int len; while ((len = tarIn.read(buffer)) != -1) { baos.write(buffer, 0, len); } if (entry.getName().endsWith("layout.txt")) { layoutContent = baos.toString("UTF-8"); } else if (entry.getName().endsWith("data.txt")) { dataContent = baos.toString("UTF-8"); } } } if (layoutContent == null || dataContent == null) { throw new IllegalArgumentException("压缩包中未找到layout或data文件"); } // 解析layout规则,转换data为JSON列表 Map<String, FieldRule> fieldRules = parseLayout(layoutContent); return dataContent.lines() .map(line -> buildJsonRecord(line, fieldRules)) .toList(); } catch (IOException e) { throw new RuntimeException("处理SFTP文件失败", e); } }) .flatMapMany(Flux::fromIterable) // 将JSON列表转为流式消息 .map(jsonStr -> MessageBuilder.withPayload(jsonStr).build()); } private Map<String, FieldRule> parseLayout(String layoutContent) { Map<String, FieldRule> rules = new HashMap<>(); // 假设layout每行格式:字段名 起始索引 长度 数据类型(示例:id 0 10 String) for (String line : layoutContent.split("\n")) { line = line.trim(); if (line.isEmpty()) continue; String[] parts = line.split("\\s+"); rules.put(parts[0], new FieldRule(Integer.parseInt(parts[1]), Integer.parseInt(parts[2]), parts[3])); } return rules; } private String buildJsonRecord(String dataLine, Map<String, FieldRule> fieldRules) { Map<String, Object> record = new HashMap<>(); fieldRules.forEach((fieldName, rule) -> { String rawValue = dataLine.substring(rule.startIndex, rule.startIndex + rule.length).trim(); switch (rule.dataType) { case "Integer": record.put(fieldName, Integer.parseInt(rawValue)); break; case "Long": record.put(fieldName, Long.parseLong(rawValue)); break; default: record.put(fieldName, rawValue); } }); try { return objectMapper.writeValueAsString(record); } catch (IOException e) { throw new RuntimeException("转换记录为JSON失败", e); } } // 内部类:存储字段解析规则 private static class FieldRule { int startIndex; int length; String dataType; FieldRule(int startIndex, int length, String dataType) { this.startIndex = startIndex; this.length = length; this.dataType = dataType; } } }
问题解决说明
- MonoMap转换错误:通过
Mono.fromCallable处理同步文件操作,再用flatMapMany将结果列表转为Flux流式输出,确保每个消息的payload是String类型(JSON字符串),适配RabbitMQ的SimpleMessageConverter要求。 - 批量消息拆分:返回
Flux<Message<String>>时,Spring Cloud Stream会自动将流中每个元素作为单独消息发送到RabbitMQ,无需额外配置批量拆分开关。
内容的提问来源于stack exchange,提问作者rmarianni
相关产品推荐
相关产品推荐

