如何使用confluent-kafka Python客户端为主题配置多Schema?
处理Kafka单主题多Avro Schema的配置问题(confluent-kafka==2.2.0)
问题解答
1. 如何传入多个Avro Schema到AvroSerializer?
AvroSerializer初始化时的schema参数是序列化默认使用的Schema,但如果采用RecordNameStrategy,序列化过程会根据待发送记录对应Avro Schema的name字段,自动匹配Schema Registry中已注册的对应Subject(格式为<record-name>)。
你不需要在初始化时传入多个Schema,只需要:
- 确保所有需要使用的Schema已经注册到Schema Registry(可通过代码或Registry UI提前完成)
- 发送记录时,保证记录结构与对应Schema匹配,Serializer会自动识别并使用正确的Schema序列化
2. record_subject_name_strategy的参数如何传递?
record_subject_name_strategy是一个回调函数,由AvroSerializer内部调用,不需要你手动传入ctx和record_name参数。你只需将函数本身赋值给subject.name.strategy配置项,不要加括号调用。
该函数内部会自动接收两个参数:
ctx:包含当前topic、partition等上下文信息的对象record_name:从待序列化记录对应的Avro Schema中提取的name字段值
完整代码示例
1. 工具函数与Schema注册(提前注册所需Schema)
import asyncio from confluent_kafka.admin import AdminClient, NewTopic from confluent_kafka import SerializingProducer, DeserializingConsumer from confluent_kafka.serialization import StringSerializer, StringDeserializer from confluent_kafka.schema_registry import record_subject_name_strategy, SchemaRegistryClient from confluent_kafka.schema_registry.avro import AvroSerializer, AvroDeserializer from confluent_kafka.schema_registry.schema import Schema from config import Config def create_admin(config: dict): return AdminClient(config) async def create_new_topic(admin: AdminClient, topic_name: str): if topic_name in admin.list_topics().topics: print("Topic already exists!") return futures = admin.create_topics([NewTopic(topic_name, num_partitions=1, replication_factor=1)]) await asyncio.wrap_future(futures[topic_name]) print("Topic created successfully!") def register_schema(schema_registry_client: SchemaRegistryClient, schema_str: str, schema_type: str = "AVRO"): """注册Schema到Schema Registry""" schema = Schema(schema_str, schema_type) subject_name = schema.name # RecordNameStrategy对应的Subject名是Schema的name字段 try: # 检查Schema是否已存在 schema_id = schema_registry_client.lookup_schema(subject_name, schema).schema_id print(f"Schema {subject_name} already registered with ID: {schema_id}") except Exception: # 不存在则注册 schema_id = schema_registry_client.register_schema(subject_name, schema).schema_id print(f"Schema {subject_name} registered with ID: {schema_id}") return schema_id
2. Producer配置(支持多Schema序列化)
async def setup_producer(schema_registry_client: SchemaRegistryClient): # 读取并注册两个Schema with open("event1.avsc", "r") as f: schema_1_str = f.read() register_schema(schema_registry_client, schema_1_str) with open("event2.avsc", "r") as f: schema_2_str = f.read() register_schema(schema_registry_client, schema_2_str) # 初始化Producer,注意subject.name.strategy的赋值方式 producer_config = Config.KAFKA | { "key.serializer": StringSerializer(), "value.serializer": AvroSerializer( schema_registry_client, schema_1_str, # 默认Schema,当发送的记录无法自动匹配时使用 conf={ "subject.name.strategy": record_subject_name_strategy # 不要加括号! } ) } producer = SerializingProducer(producer_config) return producer async def produce_messages(producer: SerializingProducer): # 发送对应event1的记录 event1_record = { "id": 1, "name": "event1_test", "timestamp": 1690000000 } producer.produce( topic="test_multischema", key="key1", value=event1_record ) # 发送对应event2的记录 event2_record = { "order_id": "ORD-123", "amount": 99.99, "status": "PAID" } producer.produce( topic="test_multischema", key="key2", value=event2_record ) producer.flush() print("Messages produced successfully!")
3. Consumer配置(支持多Schema反序列化)
async def setup_consumer(schema_registry_client: SchemaRegistryClient): consumer_config = Config.KAFKA | { "group.id": "multischema-group", "auto.offset.reset": "earliest", "key.deserializer": StringDeserializer(), "value.deserializer": AvroDeserializer( schema_registry_client, # 反序列化时不需要指定固定Schema,会根据消息中的Schema ID自动从Registry拉取 conf={ "subject.name.strategy": record_subject_name_strategy # 同样不要加括号 } ) } consumer = DeserializingConsumer(consumer_config) consumer.subscribe(["test_multischema"]) return consumer async def consume_messages(consumer: DeserializingConsumer): print("Starting consumer...") try: while True: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): print(f"Consumer error: {msg.error()}") continue print(f"Received message - Key: {msg.key()}, Value: {msg.value()}") except KeyboardInterrupt: pass finally: consumer.close()
4. 主函数
async def main(): admin = create_admin(Config.KAFKA) await create_new_topic(admin, "test_multischema") schema_registry_client = SchemaRegistryClient(Config.SCHEMA_REGISTRY) # 初始化Producer并发送消息 producer = await setup_producer(schema_registry_client) await produce_messages(producer) # 初始化Consumer并消费消息 consumer = await setup_consumer(schema_registry_client) await consume_messages(consumer) if __name__ == "__main__": asyncio.run(main())
关键注意事项
- 确保Avro Schema文件(event1.avsc、event2.avsc)中定义了唯一的
name字段,RecordNameStrategy会以此作为Subject名在Schema Registry中查找对应Schema - 反序列化时,AvroDeserializer会自动根据消息中的Schema ID从Registry拉取对应的Schema,无需手动指定多个Schema
subject.name.strategy配置项必须赋值为函数本身,不能加括号调用(否则会立即执行并传递错误参数)
内容的提问来源于stack exchange,提问作者Ivan Priz
相关产品推荐
相关产品推荐

