Loading src/telemetry/backend/service/TelemetryBackendService.py +5 −6 Original line number Diff line number Diff line Loading @@ -88,7 +88,6 @@ class TelemetryBackendService(GenericGrpcService): """ Method receives collector request and initiates collecter backend. """ # print("Initiating backend for collector: ", collector_id) LOGGER.info("Initiating backend for collector: {:s}".format(str(collector_id))) start_time = time.time() self.emulatorCollector = NetworkMetricsEmulator( Loading @@ -103,13 +102,13 @@ class TelemetryBackendService(GenericGrpcService): if not self.metric_queue.empty(): metric_value = self.metric_queue.get() LOGGER.debug("Metric: {:} - Value : {:}".format(collector['kpi_id'], metric_value)) self.GenerateCollectorResponse(collector_id, collector['kpi_id'] , metric_value) self.GenerateKpiValue(collector_id, collector['kpi_id'] , metric_value) time.sleep(1) self.TerminateCollectorBackend(collector_id) def GenerateCollectorResponse(self, collector_id: str, kpi_id: str, measured_kpi_value: Any): def GenerateKpiValue(self, collector_id: str, kpi_id: str, measured_kpi_value: Any): """ Method to write kpi value on TELEMETRY_RESPONSE Kafka topic Method to write kpi value on VALUE Kafka topic """ producer = self.kafka_producer kpi_value : Dict = { Loading @@ -118,7 +117,7 @@ class TelemetryBackendService(GenericGrpcService): "kpi_value" : measured_kpi_value } producer.produce( KafkaTopic.VALUE.value, # TODO: to the topic ... KafkaTopic.VALUE.value, key = collector_id, value = json.dumps(kpi_value), callback = self.delivery_callback Loading Loading @@ -146,7 +145,7 @@ class TelemetryBackendService(GenericGrpcService): "kpi_value" : measured_kpi_value, } producer.produce( KafkaTopic.TELEMETRY_RESPONSE.value, # TODO: to the topic ... KafkaTopic.TELEMETRY_RESPONSE.value, key = collector_id, value = json.dumps(kpi_value), callback = self.delivery_callback Loading Loading
src/telemetry/backend/service/TelemetryBackendService.py +5 −6 Original line number Diff line number Diff line Loading @@ -88,7 +88,6 @@ class TelemetryBackendService(GenericGrpcService): """ Method receives collector request and initiates collecter backend. """ # print("Initiating backend for collector: ", collector_id) LOGGER.info("Initiating backend for collector: {:s}".format(str(collector_id))) start_time = time.time() self.emulatorCollector = NetworkMetricsEmulator( Loading @@ -103,13 +102,13 @@ class TelemetryBackendService(GenericGrpcService): if not self.metric_queue.empty(): metric_value = self.metric_queue.get() LOGGER.debug("Metric: {:} - Value : {:}".format(collector['kpi_id'], metric_value)) self.GenerateCollectorResponse(collector_id, collector['kpi_id'] , metric_value) self.GenerateKpiValue(collector_id, collector['kpi_id'] , metric_value) time.sleep(1) self.TerminateCollectorBackend(collector_id) def GenerateCollectorResponse(self, collector_id: str, kpi_id: str, measured_kpi_value: Any): def GenerateKpiValue(self, collector_id: str, kpi_id: str, measured_kpi_value: Any): """ Method to write kpi value on TELEMETRY_RESPONSE Kafka topic Method to write kpi value on VALUE Kafka topic """ producer = self.kafka_producer kpi_value : Dict = { Loading @@ -118,7 +117,7 @@ class TelemetryBackendService(GenericGrpcService): "kpi_value" : measured_kpi_value } producer.produce( KafkaTopic.VALUE.value, # TODO: to the topic ... KafkaTopic.VALUE.value, key = collector_id, value = json.dumps(kpi_value), callback = self.delivery_callback Loading Loading @@ -146,7 +145,7 @@ class TelemetryBackendService(GenericGrpcService): "kpi_value" : measured_kpi_value, } producer.produce( KafkaTopic.TELEMETRY_RESPONSE.value, # TODO: to the topic ... KafkaTopic.TELEMETRY_RESPONSE.value, key = collector_id, value = json.dumps(kpi_value), callback = self.delivery_callback Loading