#!/usr/bin/env python3
"""
Avro Consumer 示例。

依赖：pip install confluent-kafka[avro]

特点：
- Consumer 不需要预先知道 Schema，反序列化时根据消息里的 Schema ID 去 Schema Registry 拉取
- 第一次拉取后会本地缓存，性能近似无开销
"""

import os
from confluent_kafka import Consumer
from confluent_kafka.serialization import SerializationContext, MessageField, StringDeserializer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "learn.18.users"
GROUP = "learn-18-avro-consumer"


def main():
    sr = SchemaRegistryClient({"url": SR_URL})

    avro_deser = AvroDeserializer(
        schema_registry_client=sr,
        # 不传 schema_str → 反序列化时按消息内的 Schema ID 自动拉取
        from_dict=lambda d, ctx: d,
    )
    key_deser = StringDeserializer("utf_8")

    consumer = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    })
    consumer.subscribe([TOPIC])
    print(f"开始消费 {TOPIC} ...  Ctrl+C 退出")

    try:
        while True:
            msg = consumer.poll(1.0)
            if msg is None:
                continue
            if msg.error():
                print(f"[ERROR] {msg.error()}")
                continue

            key = key_deser(msg.key(), SerializationContext(TOPIC, MessageField.KEY))
            value = avro_deser(msg.value(), SerializationContext(TOPIC, MessageField.VALUE))

            # 解析消息头里的 Schema ID（前 5 字节是 magic + id）
            raw = msg.value()
            schema_id = int.from_bytes(raw[1:5], byteorder="big")

            print(
                f"[RECV] partition={msg.partition()} offset={msg.offset()} "
                f"schema_id={schema_id} key={key} value={value}"
            )
            consumer.commit(msg)
    except KeyboardInterrupt:
        pass
    finally:
        consumer.close()


if __name__ == "__main__":
    main()
