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

如何在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的单例对象是从上到下顺序初始化的,原代码中:

  1. valueEncoder依赖Address的Codec,但此时addressAvroCodec和AddressTypeCodec还未初始化
  2. seqAddressAvroCodec依赖Address的Codec,同样在addressAvroCodec之前定义,导致初始化时找不到可用的隐式Codec[Address],触发空指针异常
  3. 另外,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 23:05:59