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

Kafka消费者无法反序列化SFTP CSV Source Connector生成的JSON消息

问题:Spring Kafka消费者无法反序列化SFTP CSV Connector生成的JSON消息,字段全为null

SFTP CSV Source Connector配置

name= csvsourceconnector_POC
kafka.topic =sftp-testing-topic
batch.size=1000
tasks.max=1
connector.class=io.confluent.connect.sftp.SftpCsvSourceConnector
key.converter= org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
schema.generation.enabled=true
value.converter.schema.registry=localhost:8081/
value.converter.schemas.enable=true
errors.tolerance=NONE
errors.log.enable=true
errors.log.include.messages=true
cleanup.policy=MOVE
behavior.on.error=FAIL
sftp.host=abc.xyz.com
sftp.username=username
sftp.password=password
sftp.port=22
input.path=/path/to/data
error.path=/path/to/error
finished.path=/path/to/finished
input.file.pattern=csv-sftp-source.csv

Kafka主题中的消息格式

{
  "schema": "{\"type\": \"struct\", \"fields\" : [{\"type\" : \"string\", \"optional\" : true, \"field\":\"column01\"}, {\"type\" : \"string\", \"optional\" : true, \"field\":\"column02\"}, {\"type\" : \"string\", \"optional\" : true, \"field\":\"column03\"}], \"optional\" : false, \"name\" : \"defaultValueschemaname\"}", 
  "payload" : {
    "column01" : "C00", 
    "column02" : "priorityCode",
    "column03" : "US"
  }
}

Spring Kafka消费者代码

POCKafkaConsumerService.java

public class POCKafkaConsumerService {
    
     @KafkaListener(topics = "${spring.kafka.consumer.topic}", groupId= "${spring.kafka.consumer.group-id}", properties = {"spring.json.value.default.type=com.example.CsvRecord"})
    
      public void consumeMessage(@Payload CsvRecord message, @Header(KafkaHeaders.RECEIVED_TIMESTAMP) long timestamp) {
        
             System.out.println("Received message :  " + message);
        }
}

application.properties(原配置)

spring.json.value.default.type= com.example.CsvRecord  
spring.kafka.properties.schema.registry.url : localhost: 8081/  
spring.kafka.consumer.bootstrap-server = localhost: 9092  
group-id=abc-consumergroup-1  
topic = sftp-testing-topic  
key-deserializer = org.apache.kafka.common.serialization.StringDeserializer  
value-deserializer = org.springframework.kafka.support.serializer.JsonDeserializer  

CsvRecord.java

@Data
@AllArgsConstructor
@NoArgsConstructor
@Getter
@Setter
Public clcass CsvRecord {

    @JsonProperty("column01")
    private String column01;
    
    @JsonProperty("column02")
    private String column02;
    
    @JsonProperty("column03")
    private String column03;

}

解决方案

1. 修正application.properties的配置错误

原配置存在属性名前缀缺失、拼写错误,修正后Spring才能正确读取配置:

spring.json.value.default.type=com.example.CsvRecord  
spring.kafka.properties.schema.registry.url=localhost:8081/  
spring.kafka.consumer.bootstrap-servers=localhost:9092  
spring.kafka.consumer.group-id=abc-consumergroup-1  
spring.kafka.consumer.topic=sftp-testing-topic  
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer  
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer  

2. 针对双层JSON结构配置反序列化规则

Connector生成的消息是包含schema和payload的双层JSON,而CsvRecord只对应payload内的字段,需指定反序列化时只解析payload部分:

  • 方式一:在application.properties中添加
    spring.kafka.consumer.properties.spring.json.value.json.pointer=/payload
    
  • 方式二:在@KafkaListener的properties中添加
    @KafkaListener(
        topics = "${spring.kafka.consumer.topic}", 
        groupId= "${spring.kafka.consumer.group-id}", 
        properties = {
            "spring.json.value.default.type=com.example.CsvRecord",
            "spring.json.value.json.pointer=/payload"
        }
    )
    

3. 修复CsvRecord类的语法冗余与错误

  • 修正类名拼写错误:Public clcass改为public class
  • @Data注解已包含@Getter和@Setter,可移除重复注解
    修正后的代码:
@Data
@AllArgsConstructor
@NoArgsConstructor
public class CsvRecord {

    @JsonProperty("column01")
    private String column01;
    
    @JsonProperty("column02")
    private String column02;
    
    @JsonProperty("column03")
    private String column03;

}

4. 可选:简化Connector消息格式(无需schema)

如果不需要带schema的消息结构,可修改Connector配置关闭schema生成,直接发送纯payload JSON:

value.converter.schemas.enable=false
schema.generation.enabled=false

此方式下无需配置spring.json.value.json.pointer,但需确保Connector与消费者配置一致。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 01:30:59