如何为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
相关产品推荐
相关产品推荐

