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

Spring Integration场景下,如何将SFTP获取的CSV InputStream转为POJO

解决方案

你已经用到了OpenCSV的@CsvBindByName注解,直接基于OpenCSV实现转换器即可,ObjectInputStream是用于Java对象序列化的工具,完全不适合CSV转POJO的场景。

1. 引入OpenCSV依赖

如果使用Maven,在pom.xml中添加依赖:

<dependency>
    <groupId>com.opencsv</groupId>
    <artifactId>opencsv</artifactId>
    <version>5.6</version> <!-- 可替换为最新稳定版 -->
</dependency>

2. 实现StreamToEmailInteractionConverter

转换器需要结合拆分器存入Header的表头,将单行CSV字符串转换为EmailInteractionTest对象:

import com.opencsv.bean.CsvToBean;
import com.opencsv.bean.CsvToBeanBuilder;
import com.opencsv.bean.HeaderColumnNameMappingStrategy;
import org.springframework.integration.transformer.GenericTransformer;
import org.springframework.messaging.Message;
import java.io.StringReader;

public class StreamToEmailInteractionConverter implements GenericTransformer<Message<String>, EmailInteractionTest> {

    @Override
    public EmailInteractionTest transform(Message<String> message) {
        // 从Header中获取拆分器提取的CSV表头
        String headerLine = message.getHeaders().get("myHeaders", String.class);
        // 当前消息的Payload是单行CSV数据
        String dataLine = message.getPayload();

        try (StringReader reader = new StringReader(headerLine + "\n" + dataLine)) {
            HeaderColumnNameMappingStrategy<EmailInteractionTest> strategy = new HeaderColumnNameMappingStrategy<>();
            strategy.setType(EmailInteractionTest.class);

            CsvToBean<EmailInteractionTest> csvToBean = new CsvToBeanBuilder<EmailInteractionTest>(reader)
                    .withMappingStrategy(strategy)
                    .build();

            // 仅解析一行数据,直接返回第一个结果
            return csvToBean.parse().get(0);
        } catch (Exception e) {
            throw new RuntimeException("转换CSV行到EmailInteractionTest失败", e);
        }
    }
}

3. 调整readCsvFileFlow配置

因为启用了拆分器的markers(true),会额外生成开始/结束标记消息,需要过滤掉这些非数据消息:

@Bean
public IntegrationFlow readCsvFileFlow(MessageChannel inboundFilesMessageChannel,
                                       QueueChannel kafkaPojoMessageChannel) {
    
    return IntegrationFlow.from(inboundFilesMessageChannel)
                          .split(Files.splitter()
                                      .markers(true)
                                      .charset(StandardCharsets.UTF_8)
                                      .firstLineAsHeader("myHeaders")
                                      .applySequence(true))
                          // 过滤拆分器生成的标记消息,仅处理字符串类型的数据行
                          .filter(payload -> payload instanceof String)
                          // 用自定义转换器将CSV行转为POJO
                          .transform(new StreamToEmailInteractionConverter())
                          .transform(new ObjectToJsonTransformer())
                          .log(LoggingHandler.Level.DEBUG,
                               "DataSftpToKafkaIntegrationFlow",
                               m -> "Payload: " + m.getPayload())
                          .channel(kafkaPojoMessageChannel)
                          .get();
}

关键说明

  • Files.splitter()将InputStream拆分为单行字符串,同时把CSV表头存入myHeaders Header,后续转换需要用这个表头匹配实体类字段。
  • OpenCSV的HeaderColumnNameMappingStrategy会自动根据@CsvBindByName注解匹配表头与实体类字段,无需手动编写映射关系。
  • 必须过滤标记消息,否则转换器会收到FileSplitter.FileMarker类型的Payload,导致转换失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 14:53:30