Kafka消费非序列化JSON消息遇Unexpected end-of-input异常求助
问题根因
你使用的org.apache.kafka.connect.json.JsonDeserializer默认是解析带Schema的JSON格式(即Kafka Connect生成的包含schema和payload的结构),但你的Topic中存储的是纯原生JSON字符串,导致解析时因找不到Schema部分抛出Unexpected end-of-input异常。
解决方案
方案1:用StringDeserializer手动反序列化(最稳妥适配纯JSON场景)
先将消息以字符串形式读取,再用Jackson手动转成目标DTO类,完全适配原生JSON场景。
修改application.yml配置
topics: input: datasource: test-topic --- kafka: bootstrap: servers: localhost:9092 consumers: consumer: key: deserializer: org.apache.kafka.common.serialization.StringDeserializer value: deserializer: org.apache.kafka.common.serialization.StringDeserializer
修改监听代码
import com.fasterxml.jackson.databind.ObjectMapper import com.fasterxml.jackson.databind.PropertyNamingStrategies import com.fasterxml.jackson.databind.DeserializationFeature import org.apache.kafka.clients.consumer.Consumer import org.springframework.kafka.annotation.KafkaKey import org.springframework.kafka.annotation.Topic import org.springframework.messaging.handler.annotation.MessageBody import kotlinx.coroutines.runBlocking // 初始化Jackson ObjectMapper,适配蛇形命名和忽略未知字段 private val objectMapper = ObjectMapper().apply { propertyNamingStrategy = PropertyNamingStrategies.SNAKE_CASE configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false) } @Topic(value = ["\${topics.input.datasource}"]) fun receiveNotifications( @Suppress("UNUSED_PARAMETER") @KafkaKey keys: List<String?>, @MessageBody notifications: List<String>, // 改为接收字符串列表 topics: List<String>, partitions: List<Int>, offsets: List<Long>, kafkaConsumer: Consumer<String, String> // 泛型同步改为String ) = runBlocking { notifications.forEach { jsonStr -> val guestDTO = objectMapper.readValue(jsonStr, GuestDTO::class.java) // 执行业务逻辑 } logger.info("Delete Request Processed -> Commiting Offset") kafkaConsumer.commitSync() } // 保持原DTO类不变 @JsonIgnoreProperties(ignoreUnknown = true) @JsonNaming(PropertyNamingStrategies.SnakeCaseStrategy::class) data class GuestDTO( var guests: List<Guest>? = null, var caseId: String = "" ) @JsonIgnoreProperties(ignoreUnknown = true) @JsonNaming(PropertyNamingStrategies.SnakeCaseStrategy::class) data class Guest( var guestRefId: String = "", var ids: Map<String, List<String>>? = null )
方案2:改用Spring Kafka的JsonDeserializer(自动反序列化纯JSON)
Spring Kafka自带的org.springframework.kafka.support.serializer.JsonDeserializer专门为纯JSON消息设计,无需依赖Schema即可自动反序列化。
修改application.yml配置
topics: input: datasource: test-topic --- kafka: bootstrap: servers: localhost:9092 consumers: consumer: key: deserializer: org.apache.kafka.common.serialization.StringDeserializer value: deserializer: org.springframework.kafka.support.serializer.JsonDeserializer properties: spring.json.value.default.type: com.yourpackage.GuestDTO # 替换为你的GuestDTO全类名 spring.json.trusted.packages: "*" # 允许反序列化指定包下的类,*表示全部
监听代码无需修改
保持原有的List<GuestDTO>接收逻辑即可,Spring会自动完成反序列化。
方案3:适配Kafka Connect的JsonDeserializer(不推荐)
如果必须使用org.apache.kafka.connect.json.JsonDeserializer,可配置它忽略Schema直接解析内容,但该方案稳定性不如前两者(因该反序列化器原生为带Schema场景设计):
kafka: consumers: consumer: value: deserializer: org.apache.kafka.connect.json.JsonDeserializer properties: json.value.type: com.yourpackage.GuestDTO # 替换为你的GuestDTO全类名 specific.avro.reader: false
内容的提问来源于stack exchange,提问作者Mohamed Niyaz
相关产品推荐
相关产品推荐

