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

Spring Integration中SFTP流式CSV转对象再转JSON的问题

问题描述

通过Sftp.inboundStreamingAdapter获取CSV格式的InputStream,传入消息通道后,期望将流转为MyObject对象再转成JSON。当前集成流用Files.splitter()拆分每行,自定义StreamToMyObjectConverter(基于ObjectInputStream)转换时触发报错,堆栈显示无法处理FileMarker类型消息,推测和CSV表头及拆分方式有关,求可行解决方案。

集成流代码

@Bean
public IntegrationFlow readCsvFileFlow(MessageChannel inboundFilesMessageChannel,
                                       QueueChannel kafkaPojoMessageChannel) {

    return IntegrationFlow.from(inboundFilesMessageChannel)
                          .split(Files.splitter())
                          .transform(new StreamToMyObject()) // TODO: Turn InputStream to MyObject Object
                          .transform(new ObjectToJsonTransformer())
                          .log(LoggingHandler.Level.DEBUG,
                               "AcousticEngageDataSftpToKafkaIntegrationFlow",
                               m -> "Payload: " + m.getPayload())
                          .channel(kafkaPojoMessageChannel)
                          .get();
}

自定义转换器代码

public class StreamToMyObjectConverter implements GenericTransformer<InputStream, MyObject> {
    @Override
    public MyObject transform(InputStream inputStream) {
        try (ObjectInputStream ois = new ObjectInputStream(inputStream)){
            return (MyObject) ois.readObject();
        }
        catch (IOException | ClassNotFoundException e) {
            throw new RuntimeException(e);
        }

    }
}

MyObject类(jsonschema2pojo生成)

@JsonInclude(Include.NON_NULL)
@JsonPropertyOrder({"email", "RECIPIENT_ID", "ENCODED_RECIPIENT_ID", "contactId", "code", "messageId", "userAgent", "messageName", "mailingTemplateId", "subjectLine", "docType", "reportId", "sendType", "bounceType", "urlDescription", "clickUrl", "optOutDetails", "messageGroupId", "programId", "timestamp", "originatedFrom", "eventId", "externalSystemName", "externalSystemReferenceId", "trackingCode"})
public class MyObject {
    @JsonProperty("email")
    private String email;
    @JsonProperty("RECIPIENT_ID")
    private String recipientId;
    @JsonProperty("ENCODED_RECIPIENT_ID")
    private String encodedRecipientId;
    @JsonProperty("contactId")
    private String contactId;
    @JsonProperty("code")
    private String code;
    @JsonProperty("messageId")
    private String messageId;
    @JsonProperty("userAgent")
    private String userAgent;
    @JsonProperty("messageName")
    private String messageName;
    @JsonProperty("mailingTemplateId")
    private String mailingTemplateId;
    @JsonProperty("subjectLine")
    private String subjectLine;
    @JsonProperty("docType")
    private String docType;
    @JsonProperty("reportId")
    private String reportId;
    @JsonProperty("sendType")
    private String sendType;
    @JsonProperty("bounceType")
    private String bounceType;
    @JsonProperty("urlDescription")
    private String urlDescription;
    @JsonProperty("clickUrl")
    private String clickUrl;
    @JsonProperty("optOutDetails")
    private String optOutDetails;
    @JsonProperty("messageGroupId")
    private String messageGroupId;
    @JsonProperty("programId")
    private String programId;
    @JsonProperty("timestamp")
    private String timestamp;
    @JsonProperty("originatedFrom")
    private String originatedFrom;
    @JsonProperty("eventId")
    private String eventId;
    @JsonProperty("externalSystemName")
    private String externalSystemName;
    @JsonProperty("externalSystemReferenceId")
    private String externalSystemReferenceId;
    @JsonProperty("trackingCode")
    private String trackingCode;
}

报错堆栈

Caused by: org.springframework.messaging.MessageHandlingException: error occurred during processing message in 'MethodInvokingMessageProcessor' [org.springframework.integration.handler.MethodInvokingMessageProcessor@5f2f9a4d], failedMessage=GenericMessage [payload=FileMarker [filePath=/downloaddummy_acoustic.csv, mark=START], headers={file_remoteHostPort=transfer-campaign-us-2.goacoustic.com:22, file_remoteFileInfo={"directory":false,"filename":"dummy_acoustic.csv","link":false,"modified":1686691008000,"permissions":"rw-r-----","remoteDirectory":"/download","size":1250}, file_remoteDirectory=/download, id=ca7bd14f-6967-0040-4f56-e63c715ae1c5, file_marker=START, closeableResource=org.springframework.integration.file.remote.session.CachingSessionFactory$CachedSession@4e9ca373, file_remoteFile=dummy_acoustic.csv, timestamp=1686758001364}]
    at org.springframework.integration.support.utils.IntegrationUtils.wrapInHandlingExceptionIfNecessary(IntegrationUtils.java:191)
    at org.springframework.integration.handler.MethodInvokingMessageProcessor.processMessage(MethodInvokingMessageProcessor.java:117)
    at org.springframework.integration.transformer.AbstractMessageProcessingTransformer.transform(AbstractMessageProcessingTransformer.java:115)
    at org.springframework.integration.transformer.MessageTransformingHandler.handleRequestMessage(MessageTransformingHandler.java:119)
    ... 42 more
Caused by: org.springframework.expression.spel.SpelEvaluationException: EL1004E: Method call: Method transform(org.springframework.integration.file.splitter.FileSplitter$FileMarker) cannot be found on type com.thrivent.enterprisemarketingchannelactivation.engage.converter.StreamToEmailInteractionConverter
    at org.springframework.expression.spel.ast.MethodReference.findAccessorForMethod(MethodReference.java:225)
    at org.springframework.expression.spel.ast.MethodReference.getValueInternal(MethodReference.java:135)
    at org.springframework.expression.spel.ast.MethodReference$MethodValueRef.getValue(MethodReference.java:380)
    at org.springframework.expression.spel.ast.CompoundExpression.getValueInternal(CompoundExpression.java:93)
    at org.springframework.expression.spel.ast.SpelNodeImpl.getTypedValue(SpelNodeImpl.java:119)
    at org.springframework.expression.spel.standard.SpelExpression.getValue(SpelExpression.java:376)
    at org.springframework.integration.util.AbstractExpressionEvaluator.evaluateExpression(AbstractExpressionEvaluator.java:169)
    at org.springframework.integration.util.AbstractExpressionEvaluator.evaluateExpression(AbstractExpressionEvaluator.java:154)
    at org.springframework.integration.handler.support.MessagingMethodInvokerHelper.invokeExpression(MessagingMethodInvokerHelper.java:611)
    at org.springframework.integration.handler.support.MessagingMethodInvokerHelper.fallbackToInvokeExpression(MessagingMethodInvokerHelper.java:604)
    at org.springframework.integration.handler.support.MessagingMethodInvokerHelper.processInvokeExceptionAndFallbackToExpressionIfAny(MessagingMethodInvokerHelper.java:590)
    at org.springframework.integration.handler.support.MessagingMethodInvokerHelper.invokeHandlerMethod(MessagingMethodInvokerHelper.java:561)
    at org.springframework.integration.handler.support.MessagingMethodInvokerHelper.processInternal(MessagingMethodInvokerHelper.java:476)
    at org.springframework.integration.handler.support.MessagingMethodInvokerHelper.process(MessagingMethodInvokerHelper.java:354)
    at org.springframework.integration.handler.MethodInvokingMessageProcessor.processMessage(MethodInvokingMessageProcessor.java:114)
    ... 44 more
解决方案

1. 处理FileMarker消息,过滤或禁用标记

Files.splitter()默认会发送FileMarker消息(START/END标记),这些消息不是CSV行数据,而你的转换器仅处理InputStream类型,导致类型不匹配报错。两种处理方式:

方式一:禁用标记消息

创建splitter时关闭标记生成:

.split(Files.splitter().markers(false))

方式二:过滤标记消息

在split后添加过滤器,跳过FileMarker类型的消息:

.split(Files.splitter())
// 过滤FileMarker类型消息
.filter(payload -> !(payload instanceof FileSplitter.FileMarker))

2. 修复CSV转换逻辑,替换ObjectInputStream用法

ObjectInputStream用于读取Java序列化对象,完全不适合解析CSV文本。需改用CSV解析库(如OpenCSV、Apache Commons CSV)来解析每行文本。

示例:基于OpenCSV实现转换器

首先添加OpenCSV依赖(Maven):

<dependency>
    <groupId>com.opencsv</groupId>
    <artifactId>opencsv</artifactId>
    <version>5.6</version>
</dependency>

重写转换器,接收String类型的CSV行并映射到MyObject:

public class StreamToMyObjectConverter implements GenericTransformer<String, MyObject> {
    // 表头顺序需与CSV列顺序一致,对应MyObject的JsonPropertyOrder
    private static final String[] CSV_HEADERS = {
        "email", "RECIPIENT_ID", "ENCODED_RECIPIENT_ID", "contactId", "code", 
        "messageId", "userAgent", "messageName", "mailingTemplateId", "subjectLine", 
        "docType", "reportId", "sendType", "bounceType", "urlDescription", 
        "clickUrl", "optOutDetails", "messageGroupId", "programId", "timestamp", 
        "originatedFrom", "eventId", "externalSystemName", "externalSystemReferenceId", 
        "trackingCode"
    };

    @Override
    public MyObject transform(String csvLine) {
        try (CSVReader reader = new CSVReader(new StringReader(csvLine))) {
            HeaderColumnNameMappingStrategy<MyObject> strategy = new HeaderColumnNameMappingStrategy<>();
            strategy.setType(MyObject.class);
            strategy.setColumnMapping(CSV_HEADERS);

            CSVToBean<MyObject> csvToBean = new CSVToBeanBuilder<MyObject>(reader)
                    .withMappingStrategy(strategy)
                    .build();

            return csvToBean.parse().get(0);
        } catch (IOException e) {
            throw new RuntimeException("解析CSV行失败: " + csvLine, e);
        }
    }
}

3. 跳过CSV表头行

如果CSV第一行是表头,需在流转中跳过:

.split(Files.splitter().markers(false))
// 过滤表头行,匹配表头开头特征
.filter(payload -> {
    if (payload instanceof String line) {
        return !line.startsWith("email,RECIPIENT_ID");
    }
    return false;
})

或者在CSV解析时配置跳过表头:

CSVToBean<MyObject> csvToBean = new CSVToBeanBuilder<MyObject>(reader)
        .withMappingStrategy(strategy)
        .withSkipLines(1) // 跳过第一行表头
        .build();

最终完整集成流示例

@Bean
public IntegrationFlow readCsvFileFlow(MessageChannel inboundFilesMessageChannel,
                                       QueueChannel kafkaPojoMessageChannel) {

    return IntegrationFlow.from(inboundFilesMessageChannel)
                          .split(Files.splitter().markers(false))
                          .filter(payload -> {
                              if (payload instanceof String line) {
                                  return !line.startsWith("email,RECIPIENT_ID");
                              }
                              return false;
                          })
                          .transform(new StreamToMyObjectConverter())
                          .transform(new ObjectToJsonTransformer())
                          .log(LoggingHandler.Level.DEBUG,
                               "AcousticEngageDataSftpToKafkaIntegrationFlow",
                               m -> "Payload: " + m.getPayload())
                          .channel(kafkaPojoMessageChannel)
                          .get();
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 00:24:56