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

使用Python从Kafka解码Debezium生成的Avro数据失败求助

PostgreSQL Debezium Kafka Avro 解码失败问题解决

问题背景

用Debezium监听PostgreSQL数据变更,消息成功写入Kafka Topic,kafkacat能正常解析消息,但用Python解码Avro格式的payload时始终失败,试了两种方法都没得到正确结果。

相关配置

PostgreSQL 表结构

CREATE TABLE public.users (
    id SERIAL PRIMARY KEY,
    name VARCHAR(50) NOT NULL,
    email VARCHAR(100) UNIQUE NOT NULL,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

Docker Compose 配置(Debezium + Kafka + Schema Registry)

version: '3.8'
services:
  postgres:
    image: postgres:14-alpine
    environment:
      POSTGRES_USER: postgres
      POSTGRES_PASSWORD: postgres
      POSTGRES_DB: demo
    ports:
      - "5432:5432"
    command: ["postgres", "-c", "wal_level=logical"]

  zookeeper:
    image: confluentinc/cp-zookeeper:7.4.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  kafka:
    image: confluentinc/cp-kafka:7.4.0
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
      - "29092:29092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1

  schema-registry:
    image: confluentinc/cp-schema-registry:7.4.0
    depends_on:
      - kafka
    ports:
      - "8081:8081"
    environment:
      SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://kafka:29092
      SCHEMA_REGISTRY_HOST_NAME: schema-registry
      SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081

  debezium-connect:
    image: debezium/connect:2.4
    depends_on:
      - kafka
      - postgres
    ports:
      - "8083:8083"
    environment:
      BOOTSTRAP_SERVERS: kafka:29092
      GROUP_ID: 1
      CONFIG_STORAGE_TOPIC: connect_configs
      OFFSET_STORAGE_TOPIC: connect_offsets
      STATUS_STORAGE_TOPIC: connect_statuses
      KEY_CONVERTER: io.confluent.connect.avro.AvroConverter
      KEY_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081
      VALUE_CONVERTER: io.confluent.connect.avro.AvroConverter
      VALUE_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081

尝试的解码方法及错误输出

方法1:直接用avro库解码(未处理Schema Registry前缀)

import avro.schema
from avro.io import DatumReader
import io
import kafka

consumer = kafka.KafkaConsumer(
    'postgres.demo.public.users',
    bootstrap_servers=['localhost:9092'],
    auto_offset_reset='earliest'
)

schema = avro.schema.parse(open("user_schema.avsc", "r").read())

for msg in consumer:
    reader = DatumReader(schema)
    decoded = reader.read(io.BytesIO(msg.value))
    print(decoded)

输出:

avro.io.AvroTypeException: Invalid data: b'\x00\x00\x00\x00\x01' does not match union type

方法2:用confluent-kafka的AvroConsumer但配置/处理有误

from confluent_kafka.avro import AvroConsumer
from confluent_kafka.avro.serializer import SerializerError

c = AvroConsumer({
    'bootstrap.servers': 'localhost:9092',
    'group.id': 'test-group',
    'auto.offset.reset': 'earliest',
    'schema.registry.url': 'http://localhost:8081'
})

c.subscribe(['postgres.demo.public.users'])

while True:
    try:
        msg = c.poll(1.0)
        if msg is None:
            continue
        if msg.error():
            print("Consumer error: {}".format(msg.error()))
            continue
        print(msg.value())
    except SerializerError as e:
        print("Serializer error: {}".format(e))
        break
c.close()

输出:

Serializer error: Unknown magic byte!

问题原因

  1. Debezium写入Kafka的Avro消息开头包含Schema Registry的ID前缀(4字节magic byte + 4字节schema ID),直接用avro库解码会因为无法识别前缀报错。
  2. 方法2的错误通常是Schema Registry地址配置错误、网络连通性问题(比如Docker端口映射未生效),或者未正确处理Debezium的嵌套消息结构。
  3. 忽略了Debezium消息的层级结构:实际业务数据在after字段中,顶层还有before、source、op等元数据字段。

正确解决方案

步骤1:确认Schema Registry可访问

本地执行以下命令,能返回包含postgres.demo.public.users-value的Subject列表说明连通正常:

curl http://localhost:8081/subjects

步骤2:正确配置AvroConsumer并处理Debezium消息结构

from confluent_kafka.avro import AvroConsumer
from confluent_kafka.avro.serializer import SerializerError

def main():
    consumer_config = {
        'bootstrap.servers': 'localhost:9092',
        'group.id': 'debezium-avro-consumer',
        'auto.offset.reset': 'earliest',
        'schema.registry.url': 'http://localhost:8081'
    }

    consumer = AvroConsumer(consumer_config)
    consumer.subscribe(['postgres.demo.public.users'])

    try:
        while True:
            msg = consumer.poll(1.0)
            if msg is None:
                continue
            if msg.error():
                print(f"Error: {msg.error()}")
                continue
            
            # 提取Debezium消息中的业务数据(after字段)
            payload = msg.value()
            print("原始消息结构:", payload)
            print("变更后的数据:", payload.get('after'))

    except SerializerError as e:
        print(f"序列化错误:{e}")
    except KeyboardInterrupt:
        pass
    finally:
        consumer.close()

if __name__ == "__main__":
    main()

步骤3:验证解码结果

当PostgreSQL插入一条数据:

INSERT INTO users (name, email) VALUES ('Alice', 'alice@example.com');

Python脚本会输出:

原始消息结构: {'before': None, 'after': {'id': 1, 'name': 'Alice', 'email': 'alice@example.com', 'created_at': 1699999999000}, 'source': {...}, 'op': 'c', 'ts_ms': 1699999999000, ...}
变更后的数据: {'id': 1, 'name': 'Alice', 'email': 'alice@example.com', 'created_at': 1699999999000}

为什么kafkacat能正常解析?

kafkacat通过-s value=avro -r http://localhost:8081参数自动处理了Avro消息的前缀,从Schema Registry获取对应schema并完成解码,无需手动处理前缀和结构。例如执行:

kafkacat -b localhost:9092 -t postgres.demo.public.users -s value=avro -r http://localhost:8081 -C -o beginning

就能直接看到结构化的消息内容。

内容的提问来源于stack exchange,提问作者Soumil Nitin Shah

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 09:50:29