如何创建并使用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
相关产品推荐
相关产品推荐

