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

如何为Kafka Avro主题生成Tombstone消息?

使用FS2 Kafka Vulcan生产Avro主题Tombstone消息时的错误解决

环境

  • Scala 2.13.10
  • FS2 Kafka v3.0.0-M8
  • Vulcan Avro序列化模块
  • 需求:从Kafka主题A消费消息,当消息值匹配特定条件时,向同一主题生产tombstone消息(用于压缩主题清理旧数据)

现有代码

val producerSettings =
  ProducerSettings(
    keySerializer = keySerializer,
    valueSerializer = Serializer.unit[IO]
  ).withBootstrapServers("localhost:9092")


def processRecord(committableRecord: CommittableConsumerRecord[IO, KeySchema, ValueSchema]
                     , producer: KafkaProducer.Metrics[IO, KeySchema, Unit]
                     ): IO[CommittableOffset[IO]] = {
      val key = committableRecord.record.key
      val value = committableRecord.record.value

      if(value.filterColumn.field1 == "<removable>") {
        val tombStone = ProducerRecord(committableRecord.record.topic, key, ())
        val producerRecord: ProducerRecords[CommittableOffset[IO], KeySchema, Unit] = ProducerRecords.one(tombStone, committableRecord.offset)
        producer.produce(producerRecord).flatten.map(_.flatMap(a => {
          IO(a.passthrough)
        }))
      }
      else
      IO(committableRecord.offset)
    }

错误信息

当生产tombstone时抛出以下异常:

java.lang.IllegalArgumentException: Invalid Avro record: bytes is null or empty
    at fs2.kafka.vulcan.AvroDeserializer$.$anonfun$using$4(AvroDeserializer.scala:32)
    at defer @ fs2.kafka.vulcan.AvroDeserializer$.$anonfun$using$3(AvroDeserializer.scala:29)
    at defer @ fs2.kafka.vulcan.AvroDeserializer$.$anonfun$using$3(AvroDeserializer.scala:29)
    at mapN @ fs2.kafka.KafkaProducerConnection$$anon$1.withSerializersFrom(KafkaProducerConnection.scala:141)
    at map @ fs2.kafka.ConsumerRecord$.fromJava(ConsumerRecord.scala:184)
    at map @ fs2.kafka.internal.KafkaConsumerActor.$anonfun$records$2(KafkaConsumerActor.scala:265)
    at traverse @ fs2.kafka.KafkaConsumer$$anon$1.$anonfun$partitionsMapStream$26(KafkaConsumer.scala:267)
    at defer @ fs2.kafka.vulcan.AvroDeserializer$.$anonfun$using$3(AvroDeserializer.scala:29)
    at defer @ fs2.kafka.vulcan.AvroDeserializer$.$anonfun$using$3(AvroDeserializer.scala:29)
    at mapN @ fs2.kafka.KafkaProducerConnection$$anon$1.withSerializersFrom(KafkaProducerConnection.scala:141)

目标Avro Schema

{
    "type": "record",
    "name": "SampleOrder",
    "namespace": "com.myschema.global",
    "fields": [
        {
            "name": "cust_id",
            "type": "int"
        },
        {
            "name": "month",
            "type": "int"
        },
        {
            "name": "expenses",
            "type": "double"
        },
        {
            "name": "filterColumn",
            "type": {
                "type": "record",
                "name": "filterColumn",
                "fields": [
                    {
                        "name": "id",
                        "type": "string"
                    },
                    {
                        "name": "field1",
                        "type": "string"
                    }
                ]
            }
        }
    ]
}

解决方案

问题根源

当前使用Serializer.unit[IO]序列化tombstone的value,该序列化器会把()转换为空字节数组,但Vulcan的Avro反序列化器期望接收合法的Avro编码数据,空字节不符合要求,因此抛出错误。

Kafka的tombstone消息定义是value为null,而非空字节。需要调整序列化器支持生成null值,同时反序列化器也要能处理null值。

步骤1:调整生产者序列化器

将value序列化器改为支持Option[ValueSchema]的Avro序列化器,传入None时会自动序列化为Kafka的null value(即tombstone):

// 假设你原本的Avro value序列化器是avroValueSerializer: Serializer[IO, ValueSchema]
val valueSerializer: Serializer[IO, Option[ValueSchema]] = Serializer.option(avroValueSerializer)

val producerSettings =
  ProducerSettings(
    keySerializer = keySerializer,
    valueSerializer = valueSerializer
  ).withBootstrapServers("localhost:9092")

步骤2:修改tombstone生产逻辑

生产tombstone时,将value设置为None而非():

def processRecord(committableRecord: CommittableConsumerRecord[IO, KeySchema, ValueSchema]
                     , producer: KafkaProducer.Metrics[IO, KeySchema, Option[ValueSchema]]
                     ): IO[CommittableOffset[IO]] = {
      val key = committableRecord.record.key
      val value = committableRecord.record.value

      if(value.filterColumn.field1 == "<removable>") {
        // 用None作为value生成tombstone
        val tombStone = ProducerRecord(committableRecord.record.topic, key, None)
        val producerRecord: ProducerRecords[CommittableOffset[IO], KeySchema, Option[ValueSchema]] = ProducerRecords.one(tombStone, committableRecord.offset)
        producer.produce(producerRecord).flatten.map(_.passthrough)
      } else {
        IO(committableRecord.offset)
      }
    }

步骤3:调整消费者反序列化器

为了让消费端能正确处理tombstone(null value),需要将消费者的value反序列化器也改为支持Option[ValueSchema]:

// 假设原本的Avro反序列化器是avroValueDeserializer: Deserializer[IO, ValueSchema]
val valueDeserializer: Deserializer[IO, Option[ValueSchema]] = Deserializer.option(avroValueDeserializer)

val consumerSettings =
  ConsumerSettings(
    keyDeserializer = keyDeserializer,
    valueDeserializer = valueDeserializer
  ).withBootstrapServers("localhost:9092")

额外说明

无需修改Avro Schema——Kafka的tombstone是独立于Avro Schema的消息类型,消费端用Option类型的反序列化器,就能正确识别null值(解析为None),正常消息则解析为Some(ValueSchema),不会影响原有逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 05:05:14