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

如何用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;
        }
    }
}

问题解决说明

  1. MonoMap转换错误:通过Mono.fromCallable处理同步文件操作,再用flatMapMany将结果列表转为Flux流式输出,确保每个消息的payload是String类型(JSON字符串),适配RabbitMQ的SimpleMessageConverter要求。
  2. 批量消息拆分:返回Flux<Message<String>>时,Spring Cloud Stream会自动将流中每个元素作为单独消息发送到RabbitMQ,无需额外配置批量拆分开关。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 18:08:13