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

如何创建并使用JSON反序列化器与Jackson自定义序列化器(Kafka流场景)

一、如何创建并使用JSON反序列化器

我以Java/Scala生态里常用的Jackson(结合Scala模块)为例,给你一步步拆解:

  • 1. 先定义实体类
    就用你代码里的Person类,用Option来处理空值是非常合适的,能避免空指针问题:

    case class Person(user: Option[String])
    
  • 2. 实现自定义反序列化器(按需)
    如果默认的Jackson反序列化逻辑能满足需求,其实不用自定义,但如果有特殊规则(比如要处理非标准的空值格式、字段映射),可以自己写:

    import com.fasterxml.jackson.databind.{DeserializationContext, JsonDeserializer}
    import com.fasterxml.jackson.core.JsonParser
    import scala.util.Try
    
    class PersonDeserializer extends JsonDeserializer[Person] {
      override def deserialize(p: JsonParser, ctxt: DeserializationContext): Person = {
        val node = p.getCodec.readTree(p)
        // 处理user字段不存在、值为null或者格式错误的情况
        val userOpt = Option(node.get("user")).flatMap(n => Try(n.asText()).toOption)
        Person(userOpt)
      }
    }
    
  • 3. 配置并实际使用
    把自定义反序列化器注册到Jackson的ObjectMapper里,之后就可以用来解析JSON了:

    import com.fasterxml.jackson.databind.ObjectMapper
    import com.fasterxml.jackson.module.scala.DefaultScalaModule
    
    val mapper = new ObjectMapper()
    mapper.registerModule(DefaultScalaModule) // 必须注册Scala模块,不然处理不了Option等Scala类型
    mapper.addDeserializationProblemHandler(new PersonDeserializer())
    
    // 测试反序列化,包括空值场景
    val nullUserJson = """{"user": null}"""
    val noUserJson = """{}"""
    val validUserJson = """{"user": "srinivas"}"""
    
    println(mapper.readValue(nullUserJson, classOf[Person])) // 输出 Person(None)
    println(mapper.readValue(noUserJson, classOf[Person]))   // 输出 Person(None)
    println(mapper.readValue(validUserJson, classOf[Person]))// 输出 Person(Some(srinivas))
    
二、Kafka流的Jackson自定义序列化器(解决空值失败问题)

看了你贴的代码,问题出在你的CustomSerializer只处理了user字段是字符串的情况,完全没覆盖user字段为null、字段不存在或者输入为空的场景,这就导致遇到空值时匹配不到任何case,直接抛出异常。下面是修正后的完整实现:

1. 修正后的json4s CustomSerializer

import org.json4s._
import org.json4s.jackson.JsonMethods._

case class Person(user: Option[String])

object PersonSerializer extends CustomSerializer[Person](formats => (
  // 反序列化逻辑:覆盖所有可能的输入情况
  {
    // 情况1:只有user字段且是字符串
    case JObject(JField("user", JString(user)) :: Nil) => Person(Some(user))
    // 情况2:只有user字段但值为null
    case JObject(JField("user", JNull) :: Nil) => Person(None)
    // 情况3:没有任何字段(空对象)
    case JObject(Nil) => Person(None)
    // 情况4:有多个字段,只提取user字段
    case JObject(fields) => 
      val userOpt = fields.find(_._1 == "user").flatMap {
        case (_, JString(u)) => Some(u)
        case _ => None // 不管是null还是其他类型,都返回None
      }
      Person(userOpt)
    // 情况5:输入是null(Kafka可能传递null消息)
    case JNull => Person(None)
  },
  // 序列化逻辑:把Person转成JSON,空值转成null
  {
    case Person(Some(user)) => JObject(JField("user", JString(user)) :: Nil)
    case Person(None) => JObject(JField("user", JNull) :: Nil)
  }
))

2. 在Kafka Streams中配置使用这个序列化器

要让Kafka流能识别这个序列化器,需要把它包装成Kafka的Serde(序列化/反序列化器对):

步骤1:创建自定义PersonSerde

import org.apache.kafka.common.serialization.{Deserializer, Serde, Serializer}
import org.json4s._
import org.json4s.jackson.Serialization

class PersonSerde extends Serde[Person] {
  // 配置json4s的Formats,加上我们的自定义序列化器
  private implicit val formats: Formats = DefaultFormats + PersonSerializer

  override def serializer(): Serializer[Person] = new Serializer[Person] {
    override def serialize(topic: String, data: Person): Array[Byte] = {
      // 处理data为null的情况
      if (data == null) Array.emptyByteArray
      else Serialization.write(data).getBytes("UTF-8")
    }
  }

  override def deserializer(): Deserializer[Person] = new Deserializer[Person] {
    override def deserialize(topic: String, data: Array[Byte]): Person = {
      // 处理空字节数组或者null的情况
      if (data == null || data.isEmpty) Person(None)
      else Serialization.read[Person](new String(data, "UTF-8"))
    }
  }
}

步骤2:在Kafka Streams拓扑中使用

import org.apache.kafka.streams.KafkaStreams
import org.apache.kafka.streams.StreamsBuilder
import org.apache.kafka.streams.kstream.Consumed

// 构建流拓扑
val builder = new StreamsBuilder()
// 使用自定义的PersonSerde消费主题
val personStream = builder.stream("your-input-topic", Consumed.`with`(new PersonSerde()))

// 这里可以加你的业务处理逻辑,比如打印、过滤、聚合等
personStream.foreach((key, person) => println(s"Received person from Kafka: $person"))

// 启动Kafka流
val streamsConfig = yourKafkaStreamsConfig // 这里替换成你的实际配置
val streams = new KafkaStreams(builder.build(), streamsConfig)
streams.start()

// 优雅关闭的钩子(可选,但建议加)
sys.addShutdownHook(streams.close())

几个关键提醒

  • 一定要覆盖所有空值/异常场景:包括消息体为空、字段不存在、字段值为null,不然Kafka流任务很容易崩溃
  • json4s的DefaultFormats必须加上你的自定义序列化器,不然不会生效
  • 在Serde里要处理data == null的情况,因为Kafka允许传递null消息

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:04:58