SFTP拉取CSV转JSON发Kafka时出现类型转换异常求助
问题
尝试通过SFTP连接服务器获取所有.csv文件,拆分为单行对象并转换为JSON后发送至Kafka Topic,但出现类型转换异常:
Can't convert value of class org.apache.sshd.sftp.client.impl.SftpInputStreamAsync to class org.apache.kafka.common.serialization.StringSerializer specified in value.serializer
完整堆栈跟踪信息:
2023-06-13T14:15:44.196-05:00 ERROR 71399 --- [ scheduling-1] o.s.integration.handler.LoggingHandler : org.springframework.messaging.MessageHandlingException: error occurred in message handler [bean 'publishToKafkaFlow.kafka:outbound-channel-adapter#0' for component 'EngageKafkaProducer'; defined in: 'class path resource [com/thrivent/enterprisemarketingchannelactivation/engage/config/AcousticEngageDataSftpToKafkaIntegrationFlow.class]'; from source: 'bean method publishToKafkaFlow'], failedMessage=GenericMessage [payload=SftpInputStreamAsync[ClientSessionImpl[cp.api@engage.thrivent@transfer-campaign-us-2.goacoustic.com/54.88.73.172:22]][/download/email_metadata_UNICA_CUSTID_Recipient-Event-Bulk-Export_6748469_Email_Jun-08-2023-20-10-00_20-21.csv], headers={file_remoteHostPort=transfer-campaign-us-2.goacoustic.com:22, file_remoteFileInfo={"directory":false,"filename":"email_metadata_UNICA_CUSTID_Recipient-Event-Bulk-Export_6748469_Email_Jun-08-2023-20-10-00_20-21.csv","link":false,"modified":1686258599000,"permissions":"rw-rw-r--","remoteDirectory":"/download","size":7358023}, sequenceNumber=1, file_remoteDirectory=/download, sequenceSize=1, correlationId=79dc50e4-5bb3-8ff6-9edc-1c9ec8f9e88f, id=fa20238f-c641-a407-f8c9-8d318238b4c2, closeableResource=org.springframework.integration.file.remote.session.CachingSessionFactory$CachedSession@64b375da, file_remoteFile=email_metadata_UNICA_CUSTID_Recipient-Event-Bulk-Export_6748469_Email_Jun-08-2023-20-10-00_20-21.csv, timestamp=1686683743081}] at org.springframework.integration.support.utils.IntegrationUtils.wrapInHandlingExceptionIfNecessary(IntegrationUtils.java:191) at org.springframework.integration.handler.AbstractMessageHandler.doHandleMessage(AbstractMessageHandler.java:108) at org.springframework.integration.handler.AbstractMessageHandler.handleWithMetrics(AbstractMessageHandler.java:90) at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:70) at org.springframework.integration.endpoint.PollingConsumer.handleMessage(PollingConsumer.java:158) at org.springframework.integration.endpoint.AbstractPollingEndpoint.messageReceived(AbstractPollingEndpoint.java:474) at org.springframework.integration.endpoint.AbstractPollingEndpoint.doPoll(AbstractPollingEndpoint.java:460) at org.springframework.integration.endpoint.AbstractPollingEndpoint.pollForMessage(AbstractPollingEndpoint.java:412) at org.springframework.integration.endpoint.AbstractPollingEndpoint.lambda$createPoller$4(AbstractPollingEndpoint.java:348) at org.springframework.integration.util.ErrorHandlingTaskExecutor.lambda$execute$0(ErrorHandlingTaskExecutor.java:57) at org.springframework.core.task.SyncTaskExecutor.execute(SyncTaskExecutor.java:50) at org.springframework.integration.util.ErrorHandlingTaskExecutor.execute(ErrorHandlingTaskExecutor.java:55) at org.springframework.integration.endpoint.AbstractPollingEndpoint.lambda$createPoller$5(AbstractPollingEndpoint.java:341) at org.springframework.scheduling.support.DelegatingErrorHandlingRunnable.run(DelegatingErrorHandlingRunnable.java:54) at org.springframework.scheduling.concurrent.ReschedulingRunnable.run(ReschedulingRunnable.java:96) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run$$$capture(FutureTask.java:264) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java) at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:833) Caused by: org.apache.kafka.common.errors.SerializationException: Can't convert value of class org.apache.sshd.sftp.client.impl.SftpInputStreamAsync to class org.apache.kafka.common.serialization.StringSerializer specified in value.serializer at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:1008) at org.apache.kafka.clients.producer.KafkaProducer.send(KafkaProducer.java:952) at org.springframework.kafka.core.DefaultKafkaProducerFactory$CloseSafeProducer.send(DefaultKafkaProducerFactory.java:1022) at org.springframework.kafka.core.KafkaTemplate.doSend(KafkaTemplate.java:783) at org.springframework.kafka.core.KafkaTemplate.observeSend(KafkaTemplate.java:754) at org.springframework.kafka.core.KafkaTemplate.send(KafkaTemplate.java:564) at org.springframework.integration.kafka.outbound.KafkaProducerMessageHandler.handleRequestMessage(KafkaProducerMessageHandler.java:528) at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:136) at org.springframework.integration.handler.AbstractMessageHandler.doHandleMessage(AbstractMessageHandler.java:105) ... 20 more Caused by: java.lang.ClassCastException: class org.apache.sshd.sftp.client.impl.SftpInputStreamAsync cannot be cast to class java.lang.String (org.apache.sshd.sftp.client.impl.SftpInputStreamAsync is in unnamed module of loader 'app'; java.lang.String is in module java.base of loader 'bootstrap') at org.apache.kafka.common.serialization.StringSerializer.serialize(StringSerializer.java:29) at org.apache.kafka.common.serialization.Serializer.serialize(Serializer.java:62) at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:1005) ... 28 more
相关Integration Flow代码:
@Configuration @RequiredArgsConstructor @Slf4j public class AcousticEngageDataSftpToKafkaIntegrationFlow { private static final String KAFKA_TOPIC = "my-topic"; @Bean public QueueChannel inboundFilesMessageChannel() { return MessageChannels.queue().get(); } @Bean public QueueChannel kafkaPojoMessageChannel() { return MessageChannels.queue().get(); } @Bean public MessageChannel kafkaProducerErrorRecordChannel() { return MessageChannels.direct().get(); } @Bean public IntegrationFlow sftpFileTransferFlow(SessionFactory<SftpClient.DirEntry> engageSftpSessionFactory, IntegrationFlowProperties engageProperties, MessageChannel inboundFilesMessageChannel) { return IntegrationFlow .from(Sftp.inboundStreamingAdapter(new RemoteFileTemplate<>(engageSftpSessionFactory)) .filter(new SftpRegexPatternFileListFilter(engageProperties.getRemoteFilePattern())) .remoteDirectory(engageProperties.getRemoteDirectory()), e -> e.id("sftpInboundAdapter") .autoStartup(true) .poller(Pollers.fixedRate(5000))) .log(LoggingHandler.Level.DEBUG, "AcousticEngageDataSftpToKafkaIntegrationFlow", "headers['file_remoteDirectory'] + + T(java.io.File).separator + headers['file_remoteFile']") .channel(inboundFilesMessageChannel) .get(); } @Bean public IntegrationFlow readCsvFileFlow(MessageChannel inboundFilesMessageChannel, QueueChannel kafkaPojoMessageChannel) { return IntegrationFlow.from(inboundFilesMessageChannel) .split(s -> s.delimiters("\n")) .log(LoggingHandler.Level.DEBUG, "AcousticEngageDataSftpToKafkaIntegrationFlow", m -> "Payload: " + m.getPayload()) .channel(kafkaPojoMessageChannel) .get(); } @Bean public IntegrationFlow publishToKafkaFlow(KafkaTemplate<String, String> kafkaTemplate, MessageChannel kafkaProducerErrorRecordChannel, QueueChannel kafkaPojoMessageChannel) { return IntegrationFlow.from(kafkaPojoMessageChannel) .log(LoggingHandler.Level.DEBUG, "AcousticEngageDataSftpToKafkaIntegrationFlow", e -> "Payload: " + e.getPayload()) .handle(Kafka.outboundChannelAdapter(kafkaTemplate) .topic(KAFKA_TOPIC), e -> e.id("EngageKafkaProducer")) .routeByException(r -> r .channelMapping(KafkaProducerException.class, kafkaProducerErrorRecordChannel) .defaultOutputChannel("errorChannel")) .get(); } @Bean @Transformer(inputChannel="inboundFilesMessageChannel", outputChannel="kafkaPojoMessageChannel") ObjectToJsonTransformer objectToJsonTransformer() { return new ObjectToJsonTransformer(); } @Bean public IntegrationFlow logErrorsInErrorQueue() { return IntegrationFlow.from("errorChannel") .wireTap(f -> f.handle(m -> log.error("Error occurred in AcousticEngageDataSftpToKafkaIntegrationFlow: {} ", m))) .channel("kafkaProducerErrorRecordChannel") .get(); } @Bean public IntegrationFlow kafkaProducerErrorFlow( final MessageChannel kafkaProducerErrorRecordChannel) { return IntegrationFlow.from(kafkaProducerErrorRecordChannel) .handle("KafkaExceptionHandler", "kafkaProducerErrorChannel") .get(); } }
已尝试不同转换器及自定义转换器,键值序列化器采用Spring配置,认为问题出在转换器或序列化环节,只需先实现数据发送到Topic,后续再适配JSON Schema。
解决方案
问题根源
Sftp.inboundStreamingAdapter输出的消息payload是SftpInputStreamAsync流对象,而非文件内容字符串,直接拆分或转JSON会导致类型不匹配。- 独立配置的
ObjectToJsonTransformer与readCsvFileFlow形成并行流程,流对象未被读取就直接流向Kafka,触发序列化错误。
修复步骤
1. 转换SFTP流为字符串
在获取SFTP文件流后,添加StreamTransformer将流转换为UTF-8字符串,为后续处理提供可操作的文本内容。
2. 整合流程顺序
将拆分、JSON转换逻辑串联到同一流程中,避免并行流程冲突,确保数据按“流→字符串→单行→JSON→Kafka”的顺序处理。
3. 移除冲突的独立Transformer
删除原有的ObjectToJsonTransformer Bean,将JSON转换逻辑嵌入到处理流程中。
修改后的代码
@Configuration @RequiredArgsConstructor @Slf4j public class AcousticEngageDataSftpToKafkaIntegrationFlow { private static final String KAFKA_TOPIC = "my-topic"; @Bean public QueueChannel inboundFilesMessageChannel() { return MessageChannels.queue().get(); } @Bean public QueueChannel kafkaPojoMessageChannel() { return MessageChannels.queue().get(); } @Bean public MessageChannel kafkaProducerErrorRecordChannel() { return MessageChannels.direct().get(); } @Bean public IntegrationFlow sftpFileTransferFlow(SessionFactory<SftpClient.DirEntry> engageSftpSessionFactory, IntegrationFlowProperties engageProperties, MessageChannel inboundFilesMessageChannel) { return IntegrationFlow .from(Sftp.inboundStreamingAdapter(new RemoteFileTemplate<>(engageSftpSessionFactory)) .filter(new SftpRegexPatternFileListFilter(engageProperties.getRemoteFilePattern())) .remoteDirectory(engageProperties.getRemoteDirectory()), e -> e.id("sftpInboundAdapter") .autoStartup(true) .poller(Pollers.fixedRate(5000))) .log(LoggingHandler.Level.DEBUG, "AcousticEngageDataSftpToKafkaIntegrationFlow", "headers['file_remoteDirectory'] + + T(java.io.File).separator + headers['file_remoteFile']") .transform(new StreamTransformer()) // 将SFTP流转换为字符串 .channel(inboundFilesMessageChannel) .get(); } @Bean public IntegrationFlow processCsvAndConvertToJsonFlow(MessageChannel inboundFilesMessageChannel, QueueChannel kafkaPojoMessageChannel) { return IntegrationFlow.from(inboundFilesMessageChannel) .split(s -> s.delimiters("\n").removeFirstHeader()) // 拆分单行,可选跳过表头 .log(LoggingHandler.Level.DEBUG, "AcousticEngageDataSftpToKafkaIntegrationFlow", m -> "CSV Line: " + m.getPayload()) .transform(new ObjectToJsonTransformer()) // 将单行字符串转换为JSON .channel(kafkaPojoMessageChannel) .get(); } @Bean public IntegrationFlow publishToKafkaFlow(KafkaTemplate<String, String> kafkaTemplate, MessageChannel kafkaProducerErrorRecordChannel, QueueChannel kafkaPojoMessageChannel) { return IntegrationFlow.from(kafkaPojoMessageChannel) .log(LoggingHandler.Level.DEBUG, "AcousticEngageDataSftpToKafkaIntegrationFlow", e -> "JSON Payload: " + e.getPayload()) .handle(Kafka.outboundChannelAdapter(kafkaTemplate) .topic(KAFKA_TOPIC), e -> e.id("EngageKafkaProducer")) .routeByException(r -> r .channelMapping(KafkaProducerException.class, kafkaProducerErrorRecordChannel) .defaultOutputChannel("errorChannel")) .get(); } @Bean public IntegrationFlow logErrorsInErrorQueue() { return IntegrationFlow.from("errorChannel") .wireTap(f -> f.handle(m -> log.error("Error occurred in AcousticEngageDataSftpToKafkaIntegrationFlow: {} ", m))) .channel("kafkaProducerErrorRecordChannel") .get(); } @Bean public IntegrationFlow kafkaProducerErrorFlow( final MessageChannel kafkaProducerErrorRecordChannel) { return IntegrationFlow.from(kafkaProducerErrorRecordChannel) .handle("KafkaExceptionHandler", "kafkaProducerErrorChannel") .get(); } }
额外说明
- 若CSV文件采用非UTF-8编码,可指定
StreamTransformer的编码:new StreamTransformer(StandardCharsets.ISO_8859_1) - 若需要将单行CSV转换为实体类后再转JSON,可在拆分后添加
CsvToBeanTransformer,指定对应的POJO类 - 拆分配置中的
removeFirstHeader()用于跳过CSV表头,若无需跳过可移除该配置
内容的提问来源于stack exchange,提问作者Beez
相关产品推荐
相关产品推荐

