"""Model to handle kafka event source messages.""" from lambdacommon.models.kafka_message import KafkaMessage class MSKMessageException(Exception): """Invalid MSK Message.""" pass class LambdaMSKEventSourceMessage: """Lambda MSK Event Source Message Model Class. Utility class to read a list of kafka messages that have been packed into a lambda MSK event source message. Message structure is: { eventSource: {event source name} eventSourceArn: {event source arn} bootstrapServers: {cluster bootstrap servers} records: {list of records - see MSKEventSourceRecord class } } Provides a generator for events. Usage: my_message = LambdaMSKEventSourceMessage(lambda_event_input) for record_key, record in my_message: // do something with record_key and record """ def __init__(self, message): """Initialize.""" if not message.get('records'): raise MSKMessageException('Invalid MSK message structure.') self.records = message.get('records') def __iter__(self): """Iterate over records.""" for key, record in self.records.items(): for single_event in record: yield key, KafkaMessage(single_event) \ if single_event else single_event