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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 16:35:04