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

如何正确反序列化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.datetime
  • io.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值转LocalDate
  • io.debezium.time.ZonedTimestamp:字符串转ZonedDateTime
  • io.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 06:15:39