X-Git-Url: https://osm.etsi.org/gitweb/?a=blobdiff_plain;f=osm_mon%2Fcore%2Fmessage_bus%2Fproducer.py;h=f6feba16352ffd22f2df43301a39ad20cb50d106;hb=refs%2Fchanges%2F70%2F6670%2F2;hp=f04ecf820da89584582814bdbe3ad4e2c21f6a52;hpb=b5b7819197730f5000d90a60ed13b32ba4e18fad;p=osm%2FMON.git diff --git a/osm_mon/core/message_bus/producer.py b/osm_mon/core/message_bus/producer.py index f04ecf8..f6feba1 100755 --- a/osm_mon/core/message_bus/producer.py +++ b/osm_mon/core/message_bus/producer.py @@ -61,22 +61,16 @@ class KafkaProducer(object): self.producer = kaf( key_serializer=str.encode, value_serializer=str.encode, - bootstrap_servers=broker, api_version=(0, 10)) + bootstrap_servers=broker, api_version=(0, 10, 1)) def publish(self, key, value, topic=None): """Send the required message on the Kafka message bus.""" try: future = self.producer.send(topic=topic, key=key, value=value) - record_metadata = future.get(timeout=10) + future.get(timeout=10) except Exception: logging.exception("Error publishing to {} topic." .format(topic)) raise - try: - logging.debug("TOPIC:", record_metadata.topic) - logging.debug("PARTITION:", record_metadata.partition) - logging.debug("OFFSET:", record_metadata.offset) - except KafkaError: - pass def publish_alarm_request(self, key, message): """Publish an alarm request."""