"""JSON Schema Deserializer.""" from confluent_kafka.schema_registry import SchemaRegistryClient from confluent_kafka.schema_registry.json_schema import JSONDeserializer as _JSONDeserializer from confluent_kafka.serialization import SerializationContext from kafka_utils.consumer.deserializer.base import BaseDeserializer from kafka_utils.consumer.deserializer.helpers import read_schema_id_from_message class JSONDeserializer(BaseDeserializer): """JSON Schema Deserializer.""" def __init__(self, schema_registry_url): """Init.""" super().__init__() self.schema_registry_client = SchemaRegistryClient({'url': schema_registry_url}) self.deserializers = {} def deserialize(self, value, topic, context): """Deserialize.""" deserializer = self.get_deserializer(value) serialization_context = SerializationContext(topic, context) return deserializer(value, serialization_context) def get_deserializer(self, value): """Get the proper deserializer by schema id.""" schema_id = self.get_schema_id(value) if self.deserializers.get(schema_id) is not None: return self.deserializers.get(schema_id) schema = self.schema_registry_client.get_schema(schema_id) self.deserializers.update({schema_id: _JSONDeserializer(schema.schema_str)}) return self.deserializers.get(schema_id) def get_schema_id(self, value): """Get the schema ID from the serialized message.""" return read_schema_id_from_message(value)