"""Event Producer.""" from confluent_kafka import Producer from confluent_kafka.serialization import MessageField from confluent_kafka.serialization import SerializationContext class EventProducer: """Event Producer.""" def __init__( self, bootstrap_servers, key_serializer=None, value_serializer=None, security_protocol='SSL', producer_options=None ): """Initialize event producer.""" self.key_serializer = key_serializer self.value_serializer = value_serializer producer_config = { 'bootstrap.servers': bootstrap_servers, 'security.protocol': security_protocol } if producer_options is not None: producer_config.update(**producer_options) self.producer = Producer(producer_config) def produce( self, topic, event_key, event_value, callback=None, key_serializer=None, value_serializer=None, auto_flush=True, **kwargs ): """Produce an event.""" if not key_serializer: key_serializer = self.key_serializer key = self._serialize(event_key, key_serializer, topic, MessageField.KEY) if not value_serializer: value_serializer = self.value_serializer value = self._serialize(event_value, value_serializer, topic, MessageField.VALUE) self.producer.produce( topic=topic, value=value, key=key, on_delivery=callback, **kwargs ) if auto_flush: self.producer.flush() def _serialize(self, value, serializer, topic, field): """Serialize a field.""" if serializer: ctx = SerializationContext(topic, field) return serializer(value, ctx) return value def __enter__(self): """Enter producer context.""" return self def __exit__(self, exc_type, exc_val, exc_tb): """Exit producer context.""" self.producer.flush()