"""Abstracted debezium message.""" from kafka_utils.consumer.deserializer.simple_json import JSONDeserializer from kafka_utils.consumer.message.base import BaseEventMessage from kafka_utils.exceptions import DebeziumMessageException from kafka_utils.exceptions import IneligibleEventException class DebeziumMessage(BaseEventMessage): """Debezium mysql replication message. Message structure: { schema: schema information payload: { op: "c|u|d", before: affected table row as a dictionary, null when op is c after: affected table row as a dictionary, null when op is d source: message metadata } } """ def __init__(self, message, topic, allowed_event_ops=None): """Parse message. Args: message (dict): the string for JSON representing the message topic (str): The topic for the message. allowed_event_ops (list): a list of allowable operations (eg: c|u|d) """ super().__init__(message, topic) self.value_deserializer = JSONDeserializer() self.message = self.value_deserializer.deserialize(self.message) if 'payload' not in self.message or not self.message.get('payload'): raise DebeziumMessageException('Missing data: payload') payload = self.message.get('payload') if 'op' not in payload or not payload.get('op'): raise DebeziumMessageException('Missing data: op') self.operation = payload.get('op') if allowed_event_ops and self.operation not in allowed_event_ops: raise IneligibleEventException(f'Event table operation ignored: {self.operation}') if self.operation != 'c': if 'before' not in payload or payload.get('before') is None: raise DebeziumMessageException('Missing data: before') self.before = payload.get('before') if self.operation != 'd': if 'after' not in payload or payload.get('after') is None: raise DebeziumMessageException('Missing data: after') self.after = payload.get('after') def get_field_values(self, field_name): """Get the before/after values for a field by name.""" values = {'before': None, 'after': None} if self.operation != 'c': values.update({'before': self.before.get(field_name)}) if self.operation != 'd': values.update({'after': self.after.get(field_name)}) return values