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

Spring Batch如何使用ItemReader读取Unix服务器上的文件?

Spring Batch读取Unix服务器SFTP文件的ItemReader实现方案

你当前手动封装JSch操作SFTP的方式可以满足基础需求,但不符合Spring Batch的批处理规范,无法直接复用框架提供的重试、跳过、事务、分片等能力,以下两种更适配Spring Batch生态的实现方式可以参考:

方案1:使用官方原生SftpItemReader(推荐)

Spring Integration提供了封装完善的SFTP操作组件,配合Spring Batch的RemoteFileItemReader即可实现无侵入的远程文件读取,不需要手动维护连接、遍历文件逻辑,是官方推荐的实现方式。

配置步骤

  1. 引入依赖(Maven示例)
<dependency>
    <groupId>org.springframework.batch</groupId>
    <artifactId>spring-batch-integration</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.integration</groupId>
    <artifactId>spring-integration-sftp</artifactId>
</dependency>
  1. 配置SFTP会话工厂
@Configuration
public class SftpConfig {
    @Value("${sftp.host}")
    private String host;
    @Value("${sftp.port}")
    private int port;
    @Value("${sftp.username}")
    private String username;
    @Value("${sftp.password}")
    private String password;

    @Bean
    public SessionFactory<LsEntry> sftpSessionFactory() {
        DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory();
        factory.setHost(host);
        factory.setPort(port);
        factory.setUser(username);
        factory.setPassword(password);
        factory.setAllowUnknownKeys(false); // 生产环境建议关闭,配置knownHosts
        // 生产环境建议用CachingSessionFactory实现连接池
        return new CachingSessionFactory<>(factory);
    }
}
  1. 配置SftpItemReader
@Bean
public ItemReader<YourDataEntity> sftpItemReader(SessionFactory<LsEntry> sftpSessionFactory) {
    SftpItemReader<YourDataEntity> reader = new SftpItemReader<>();
    reader.setSessionFactory(sftpSessionFactory);
    reader.setRemoteDirectory("/target/unix/directory");
    reader.setFilenamePattern("*" + fileName + "*"); // 匹配你需要的文件规则
    // 配置文件解析规则,比如按行读取csv
    reader.setLineMapper(new DefaultLineMapper<YourDataEntity>() {{
        setLineTokenizer(new DelimitedLineTokenizer());
        setFieldSetMapper(new BeanWrapperFieldSetMapper<YourDataEntity>() {{
            setTargetType(YourDataEntity.class);
        }});
    }});
    return reader;
}

方案2:自定义JSch逻辑的ItemReader

如果不想额外引入Spring Integration依赖,可以把你现有的JSch逻辑封装成符合Spring Batch规范的ItemReader,复用框架的批处理能力:

public class CustomSftpItemReader extends AbstractItemCountingItemStreamItemReader<YourDataEntity> {
    private ChannelSftp channelSftp;
    private Session session;
    private BufferedReader fileReader;
    private List<String> targetFiles;
    private int currentFileIndex = 0;

    // 配置参数通过setter注入
    private String host;
    private int port;
    private String username;
    private String password;
    private String remoteDir;
    private String fileNamePattern;

    @Override
    protected void doOpen() throws Exception {
        // 初始化SFTP连接,对应你原有main方法里的连接逻辑
        JSch jsch = new JSch();
        session = jsch.getSession(username, host, port);
        session.setPassword(password);
        Properties config = new Properties();
        config.put("StrictHostKeyChecking", "no");
        session.setConfig(config);
        session.connect();
        Channel channel = session.openChannel("sftp");
        channel.connect();
        channelSftp = (ChannelSftp) channel;
        channelSftp.cd(remoteDir);
        // 匹配目标文件
        targetFiles = channelSftp.ls(fileNamePattern).stream()
                .map(LsEntry::getFilename)
                .collect(Collectors.toList());
        // 打开第一个文件流
        openNextFile();
    }

    @Override
    protected YourDataEntity doRead() throws Exception {
        String line = fileReader.readLine();
        if (line == null) {
            currentFileIndex++;
            if (currentFileIndex < targetFiles.size()) {
                openNextFile();
                return doRead();
            }
            return null; // 所有文件读取完毕
        }
        // 把行转换为你需要的实体对象
        return convertLineToEntity(line);
    }

    @Override
    protected void doClose() throws Exception {
        // 释放资源
        if (fileReader != null) fileReader.close();
        if (channelSftp != null) channelSftp.disconnect();
        if (session != null) session.disconnect();
    }

    private void openNextFile() throws SftpException, IOException {
        InputStream inputStream = channelSftp.get(targetFiles.get(currentFileIndex));
        fileReader = new BufferedReader(new InputStreamReader(inputStream, StandardCharsets.UTF_8));
    }

    private YourDataEntity convertLineToEntity(String line) {
        // 实现你自己的行解析逻辑
    }
}

实用建议

  • 不要在代码中硬编码连接账号密码,建议放到Spring配置文件中,敏感信息可以用加密组件加密存储,避免泄露
  • 生产环境建议开启SFTP主机密钥校验,不要直接设置StrictHostKeyChecking=no,避免中间人攻击
  • 大文件场景下使用流式读取,不要将整个文件加载到内存,避免OOM问题
  • 多文件处理可以配合MultiResourceItemReader实现文件级的分片、重试逻辑,不需要自己遍历文件
  • 建议使用连接池管理SFTP会话,避免频繁创建销毁连接带来的性能开销

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 04:45:07