使用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!
问题原因
- Debezium写入Kafka的Avro消息开头包含Schema Registry的ID前缀(4字节magic byte + 4字节schema ID),直接用
avro库解码会因为无法识别前缀报错。 - 方法2的错误通常是Schema Registry地址配置错误、网络连通性问题(比如Docker端口映射未生效),或者未正确处理Debezium的嵌套消息结构。
- 忽略了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
相关产品推荐
相关产品推荐

