Scala:如何简化Kafka序列化中含副作用的模式匹配逻辑?
重构Avro4s Kafka序列化代码的方案
核心思路是把重复的序列化逻辑抽象出来,仅将「Avro输出流类型」作为可变参数传入,彻底消除两个序列化类里的模式匹配和流创建重复代码。
1. 抽象通用序列化基类
创建一个基础序列化类,封装所有通用逻辑(字节流处理、资源关闭、Kafka Serializer接口实现),仅留一个抽象方法让子类指定具体的Avro输出流类型:
import org.apache.kafka.common.serialization.Serializer import com.sksamuel.avro4s.{AvroOutputStream, SchemaFor} import java.io.ByteArrayOutputStream abstract class BaseKafkaAvroSerializer[T: SchemaFor] extends Serializer[T] { // 由子类实现:指定创建二进制/JSON流的逻辑 protected def createOutputStream(out: ByteArrayOutputStream): AvroOutputStream[T] override def serialize(topic: String, data: T): Array[Byte] = { val baos = new ByteArrayOutputStream() val avroOut = createOutputStream(baos) try { avroOut.write(data) avroOut.flush() baos.toByteArray } finally { avroOut.close() } } // Kafka Serializer接口的默认实现 override def configure(configs: java.util.Map[String, _], isKey: Boolean): Unit = {} override def close(): Unit = {} }
2. 实现具体序列化类
基于基类快速实现二进制和JSON序列化,仅需一行代码指定流类型:
// 二进制序列化类 class KafkaBinarySerializer[T: SchemaFor] extends BaseKafkaAvroSerializer[T] { override protected def createOutputStream(out: ByteArrayOutputStream): AvroOutputStream[T] = { AvroOutputStream.binary[T](out) } } // JSON序列化类 class KafkaJsonSerializer[T: SchemaFor] extends BaseKafkaAvroSerializer[T] { override protected def createOutputStream(out: ByteArrayOutputStream): AvroOutputStream[T] = { AvroOutputStream.json[T](out) } }
3. 处理多样例类的类型支持
如果你的序列化是针对UserCreated、UserDeleted这类实现了同一密封特质的样例类(比如sealed trait KafkaEvent),avro4s会自动为每个样例类生成SchemaFor隐式实例,无需手动写模式匹配。直接使用时指定泛型即可:
// 使用示例:为UserCreated创建二进制序列化器 val userCreatedBinarySerializer = new KafkaBinarySerializer[UserCreated]() // 为UserDeleted创建JSON序列化器 val userDeletedJsonSerializer = new KafkaJsonSerializer[UserDeleted]()
替代方案:用高阶函数封装
如果不想用类继承,也可以用高阶函数直接生成序列化器,代码更简洁:
import org.apache.kafka.common.serialization.Serializer import com.sksamuel.avro4s.{AvroOutputStream, SchemaFor} import java.io.ByteArrayOutputStream object KafkaAvroSerializers { private def buildSerializer[T: SchemaFor]( streamFactory: ByteArrayOutputStream => AvroOutputStream[T] ): Serializer[T] = new Serializer[T] { override def serialize(topic: String, data: T): Array[Byte] = { val baos = new ByteArrayOutputStream() val avroOut = streamFactory(baos) try { avroOut.write(data) avroOut.flush() baos.toByteArray } finally { avroOut.close() } } override def configure(configs: java.util.Map[String, _], isKey: Boolean): Unit = {} override def close(): Unit = {} } // 对外暴露的序列化器创建方法 def binary[T: SchemaFor]: Serializer[T] = buildSerializer(AvroOutputStream.binary[T]) def json[T: SchemaFor]: Serializer[T] = buildSerializer(AvroOutputStream.json[T]) } // 使用示例 val userCreatedBinary = KafkaAvroSerializers.binary[UserCreated] val userDeletedJson = KafkaAvroSerializers.json[UserDeleted]
重构优势
- 彻底消除重复代码:核心序列化逻辑仅实现一次,流类型的差异通过参数/子类实现隔离
- 扩展性强:如果后续需要支持Avro的其他输出格式(如格式化JSON),只需新增一个流工厂方法或子类
- 类型安全:依赖avro4s的
SchemaFor隐式机制自动处理样例类的Schema生成,无需手动维护模式匹配
内容的提问来源于stack exchange,提问作者ticofab
相关产品推荐
相关产品推荐

