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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 18:36:40