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

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。


解决方案

问题根源

  1. Sftp.inboundStreamingAdapter输出的消息payload是SftpInputStreamAsync流对象,而非文件内容字符串,直接拆分或转JSON会导致类型不匹配。
  2. 独立配置的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 21:32:01