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

