import base64 import os from confluent_kafka.schema_registry import SchemaRegistryClient from confluent_kafka.schema_registry.avro import AvroDeserializer SCHEMA_REGISTRY_CONFIG = { 'url': os.environ.get('SCHEMA_REGISTRY_URL') } schema_registry_client = SchemaRegistryClient(SCHEMA_REGISTRY_CONFIG) deserializer = AvroDeserializer( schema_registry_client=schema_registry_client) def handler(event, context): print(event) for records in event['records'].values(): for record in records: key = base64.b64decode(record['key']) value = base64.b64decode(record['value']) msg = deserializer(value, None) print(f'key: {key} value: {msg}')