"""Model to handle kafka event source messages.""" from kafka_utils.consumer.source.msk import MSKMessage from kafka_utils.exceptions import MSKMessageException class EventSourceMessage: """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 MSKMessage class } } Provides a generator for events. Usage: my_message = LambdaMSKEventSourceMessage(lambda_event_input) for record_key, record in my_message: if record.value: // do something with record_key, record.key, and record.value """ 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, MSKMessage(single_event) if single_event else single_event