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
相关产品推荐
相关产品推荐

