如何在FS2 Kafka中使用Vulcan为Scala样例类实现Codec
问题:使用FS2和Vulcan实现Avro Codec时触发空指针异常
样例类定义
case class People(name: String, address: Seq[Address]) case class Address(`type`: AddType.Value, street: String, suburb: String) object AddType extends Enumeration { type AddType = Value val H, P = Value } case class PeopleKey(eid: String)
原Codec实现
尝试编写的Codec定义如下,其中为了处理Seq[Address]手动添加了seqAddressAvroCodec,但运行时出错:
object codecPeople { implicit val keyEncoder: Codec[PeopleKey] = Codec.derive[PeopleKey] implicit val valueEncoder: Codec[People] = Codec.derive[People] implicit val seqAddressAvroCodec = Codec.seq[Address] implicit val addressAvroCodec: Codec[Address] = Codec.derive[Address] implicit val AddressTypeCodec: Aux[EnumSymbol, AddType] = Codec.enumeration[AddType]( name = "AddressTypeAvro", namespace = "au.com.avro", symbols = List("H", "P"), encode = { case AddType.H => "H" case AddType.P => "P" }, decode = { case "H" => Right(AddType.H) case "P" => Right(AddType.P) case e => Left(AvroError(s"$e is not a AddressTypeAvro")) }, default = Some(AddType.H) ) }
报错详情
运行时抛出空指针异常,错误栈指向Codec.derive[People]的初始化:
[info] Cause: java.lang.NullPointerException: [info] at vulcan.Codec$.collection(Codec.scala:759) [info] at vulcan.Codec$.list(Codec.scala:796) [info] at vulcan.Codec$.seq(Codec.scala:1128) ***[info] at au.com.stackoverflow.codecPeople$.<clinit>(codecPeople.scala:23)*** [info] at au.com.stackoverflow.producePeopleTest.produceCamCustData(producePeopleTest.scala:27) [info] at au.com.stackoverflow.producePeopleTest.$anonfun$new$2(producePeopleTest.scala:18) [info] at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.scala:18) [info] at org.scalatest.OutcomeOf.outcomeOf(OutcomeOf.scala:85) [info] at org.scalatest.OutcomeOf.outcomeOf$(OutcomeOf.scala:83) [info] at org.scalatest.OutcomeOf$.outcomeOf(OutcomeOf.scala:104)
Kafka消息生产代码
import cats.effect.unsafe.implicits.global import cats.effect.IO import fs2.kafka._ import fs2.kafka.vulcan.{AvroSettings, SchemaRegistryClientSettings, avroSerializer} import fs2.{Chunk, Pipe, Pure, Stream} import org.scalatest.freespec.AnyFreeSpec import scala.concurrent.duration.DurationInt class producePeopleTest extends AnyFreeSpec { "Integration test" - { "should be able to produce kafka records" in { val config = TestAppConfig("localhost:9092", "http://localhost:8081") produceCamCustData(config).compile.drain.unsafeRunSync() } } private def produceCamCustData(config: TestAppConfig): Stream[IO, ProducerResult[PeopleKey, People]] = { val tmsData: ProducerRecords[PeopleKey, People] = generateRecords() val avroSettings: AvroSettings[IO] = AvroSettings(SchemaRegistryClientSettings[IO](config.kafkaSchemaUrl)) implicit val keySerds = avroSerializer[PeopleKey].forKey(avroSettings) implicit val valueSerds = avroSerializer[People].forValue(avroSettings) val producerSettings = ProducerSettings[IO, PeopleKey, People].withBootstrapServers(config.kafkaBrokers) val producer : Pipe[IO, ProducerRecords[PeopleKey, People], ProducerResult[PeopleKey, People]] = KafkaProducer.pipe(producerSettings) val fs2Data: Stream[Pure, ProducerRecords[PeopleKey, People]] = fs2.Stream(tmsData) val res: Stream[IO, ProducerResult[PeopleKey, People]] = fs2Data.through(producer) producerHandleErrorRec(res) } private def producerHandleErrorRec[A](stream: Stream[IO, A]): Stream[IO, A] = stream.handleErrorWith(_ => Stream.sleep[IO](10.seconds) >> Stream.eval(IO.println("Produce Fail")) >> producerHandleErrorRec(stream)) private def generateRecords(): ProducerRecords[PeopleKey, People] = { val key = PeopleKey("1102388657") val address = Address( `type` = AddType.H, street = "1 Anna St", suburb = "Melbourne" ) val value = People( name = "Su", address = Seq(address) ) ProducerRecords( Chunk.from( List.fill(1)( ProducerRecord("test.people.topic.avro", key, value))) ) } }
解决方案
问题原因
Scala的单例对象是从上到下顺序初始化的,原代码中:
valueEncoder依赖Address的Codec,但此时addressAvroCodec和AddressTypeCodec还未初始化seqAddressAvroCodec依赖Address的Codec,同样在addressAvroCodec之前定义,导致初始化时找不到可用的隐式Codec[Address],触发空指针异常- 另外,Vulcan会自动为
Seq这类集合类型派生Codec,无需手动定义seqAddressAvroCodec
调整后的Codec代码
object codecPeople { // 先定义最底层的枚举类型Codec implicit val AddressTypeCodec: Aux[EnumSymbol, AddType] = Codec.enumeration[AddType]( name = "AddressTypeAvro", namespace = "au.com.avro", symbols = List("H", "P"), encode = { case AddType.H => "H" case AddType.P => "P" }, decode = { case "H" => Right(AddType.H) case "P" => Right(AddType.P) case e => Left(AvroError(s"$e is not a AddressTypeAvro")) }, default = Some(AddType.H) ) // 再定义Address的Codec,依赖上面的枚举Codec implicit val addressAvroCodec: Codec[Address] = Codec.derive[Address] // 最后定义依赖Address的People,以及独立的PeopleKey implicit val keyEncoder: Codec[PeopleKey] = Codec.derive[PeopleKey] implicit val valueEncoder: Codec[People] = Codec.derive[People] }
内容的提问来源于stack exchange,提问作者Chen Guo
相关产品推荐
相关产品推荐

