"""Example avro message consumer.""" from helpers import debug_msk_message from helpers import load_file_json from kafka_utils.consumer.deserializer.json_schema import JSONDeserializer from kafka_utils.consumer.deserializer.string import StringDeserializer from kafka_utils.consumer.source.mapping import EventSourceMessage SCHEMA_REGISTRY_URL = 'https://dev-schema-registry.dev.theorchard.io' EVENT_JSON = load_file_json('json_schema.json') def main(): """Process message.""" string_serializer = StringDeserializer() json_schema_serializer = JSONDeserializer(SCHEMA_REGISTRY_URL) event = EventSourceMessage(EVENT_JSON) for batch_key, msk_message in event: debug_msk_message( batch_key, msk_message, key_deserializer=string_serializer, value_deserializer=json_schema_serializer) if __name__ == '__main__': main()