如何正确反序列化Debezium生成的Avro数据?
解决方案
核心原因
Debezium为MySQL时间类型定义了自定义逻辑类型(如io.debezium.time.Date、io.debezium.time.ZonedTimestamp等),这些类型并非标准Avro类型,Confluent的Avro反序列化器只会按原始Schema类型(int32/string/int64)解析,不会自动转换为datetime对象,需要手动实现解码逻辑。
Debezium官方没有直接提供Python/Scala的解码插件,但可以根据其类型定义规则手动编写转换函数。
Python解码示例
解码逻辑说明
io.debezium.time.Date:int32值是从1970-01-01开始的天数,需转换为datetime.date对象io.debezium.time.ZonedTimestamp:string值是ISO 8601格式的带时区字符串,可直接解析为datetime.datetimeio.debezium.time.MicroTimestamp:int64值是从1970-01-01开始的微秒数,需转换为UTC时区的datetime.datetime
代码实现
import datetime from confluent_kafka.avro import AvroConsumer def decode_debezium_time_fields(record: dict) -> dict: """递归解码Debezium自定义时间类型字段""" decoded_record = {} for field, value in record.items(): if isinstance(value, dict): # 处理嵌套结构 decoded_record[field] = decode_debezium_time_fields(value) continue # 根据字段名+值类型匹配Debezium时间类型 if field == "created_date" and isinstance(value, int): decoded_record[field] = datetime.date(1970, 1, 1) + datetime.timedelta(days=value) elif field == "updated_time" and isinstance(value, str): # 兼容Python 3.11以下版本对Z时区的解析 normalized_str = value[:-1] + "+00:00" if value.endswith("Z") else value decoded_record[field] = datetime.datetime.fromisoformat(normalized_str) elif field == "created_datetime" and isinstance(value, int): decoded_record[field] = datetime.datetime.fromtimestamp( value / 1_000_000, datetime.timezone.utc ) else: decoded_record[field] = value return decoded_record # 消费并解码数据示例 consumer = AvroConsumer({ 'bootstrap.servers': 'kafka:9092', 'group.id': 'avro-mysql-payments-consumer', 'schema.registry.url': 'http://schema-registry:8081' }) consumer.subscribe(['avro.mysql.cdc.payments']) try: while True: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): print(f"消费错误: {msg.error()}") continue raw_data = msg.value() decoded_data = decode_debezium_time_fields(raw_data) print(f"解码后数据: {decoded_data}") finally: consumer.close()
Scala解码示例
解码逻辑说明
利用Java/Scala的java.time API转换:
io.debezium.time.Date:int值转LocalDateio.debezium.time.ZonedTimestamp:字符串转ZonedDateTimeio.debezium.time.MicroTimestamp:long值转Instant再转为UTC时区的ZonedDateTime
代码实现
import io.confluent.kafka.serializers.KafkaAvroDeserializer import org.apache.kafka.clients.consumer.{ConsumerConfig, KafkaConsumer} import java.time.{LocalDate, ZonedDateTime, Instant, ZoneOffset} import java.time.format.DateTimeFormatter import java.util.Properties import scala.collection.JavaConverters._ object DebeziumAvroTimeDecoder { def decodeTimeFields(record: java.util.Map[String, Any]): java.util.Map[String, Any] = { val decoded = new java.util.HashMap[String, Any]() record.asScala.foreach { case (field, value) => value match { case nested: java.util.Map[String, @unchecked Any] => decoded.put(field, decodeTimeFields(nested)) case days: Int if field == "created_date" => decoded.put(field, LocalDate.ofEpochDay(days)) case tsStr: String if field == "updated_time" => decoded.put(field, ZonedDateTime.parse(tsStr, DateTimeFormatter.ISO_ZONED_DATE_TIME)) case micros: Long if field == "created_datetime" => decoded.put(field, Instant.ofEpochMilli(micros / 1000).atZone(ZoneOffset.UTC)) case _ => decoded.put(field, value) } } decoded } def main(args: Array[String]): Unit = { val props = new Properties() props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092") props.put(ConsumerConfig.GROUP_ID_CONFIG, "avro-mysql-scala-consumer") props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, classOf[KafkaAvroDeserializer]) props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, classOf[KafkaAvroDeserializer]) props.put("schema.registry.url", "http://schema-registry:8081") props.put("specific.avro.reader", "false") val consumer = new KafkaConsumer[String, java.util.Map[String, Any]](props) consumer.subscribe(List("avro.mysql.cdc.payments").asJava) try { while (true) { val records = consumer.poll(java.time.Duration.ofMillis(1000)) records.asScala.foreach { record => val rawData = record.value() val decodedData = decodeTimeFields(rawData) println(s"解码后数据: $decodedData") } } } finally { consumer.close() } } }
进阶优化
如果需要通用解码(不依赖字段名),可在反序列化时获取字段的Schema元数据,检查字段的name属性是否为Debezium时间类型,再执行转换。例如在Python中,通过msg.value_schema()获取Schema,遍历字段的name属性判断类型后处理。
内容的提问来源于stack exchange,提问作者Aaron Li
相关产品推荐
相关产品推荐

