- vnfd = self.common_db_client.get_vnfd(vnfr['vnfd-id'])
- for vdur in vnfr['vdur']:
- # This avoids errors when vdur records have not been completely filled
- if 'name' not in vdur:
- continue
- vdu = next(
- filter(lambda vdu: vdu['id'] == vdur['vdu-id-ref'], vnfd['vdu'])
- )
- vnf_member_index = vnfr['member-vnf-index-ref']
- vdu_name = vdur['name']
- if 'monitoring-param' in vdu:
- for param in vdu['monitoring-param']:
- metric_name = param['nfvi-metric']
- payload = await self._generate_read_metric_payload(metric_name, nsr_id, vdu_name,
- vnf_member_index)
- producer.send(topic='metric_request', key='read_metric_data_request',
- value=json.dumps(payload))
- producer.flush(5)
- for message in consumer:
- if message.key == 'read_metric_data_response':
- content = json.loads(message.value)
- if content['correlation_id'] == payload['correlation_id']:
- log.debug("Found read_metric_data_response with same correlation_id")
- if len(content['metrics_data']['metrics_series']):
- metric_reading = content['metrics_data']['metrics_series'][-1]
- if metric_name not in metrics.keys():
- metrics[metric_name] = GaugeMetricFamily(
- metric_name,
- 'OSM metric',
- labels=['ns_id', 'vnf_member_index', 'vdu_name']
- )
- metrics[metric_name].add_metric([nsr_id, vnf_member_index, vdu_name],
- metric_reading)
- break
- if 'vdu-configuration' in vdu and 'metrics' in vdu['vdu-configuration']:
- vnf_name_vca = await self._generate_vca_vdu_name(vdu_name)
- vnf_metrics = await self.n2vc.GetMetrics(vca_model_name, vnf_name_vca)
- log.debug('VNF Metrics: %s', vnf_metrics)
- for vnf_metric_list in vnf_metrics.values():
- for vnf_metric in vnf_metric_list:
- log.debug("VNF Metric: %s", vnf_metric)
- if vnf_metric['key'] not in metrics.keys():
- metrics[vnf_metric['key']] = GaugeMetricFamily(
- vnf_metric['key'],
- 'OSM metric',
- labels=['ns_id', 'vnf_member_index', 'vdu_name']
- )
- metrics[vnf_metric['key']].add_metric([nsr_id, vnf_member_index, vdu_name],
- float(vnf_metric['value']))
- consumer.close()
- producer.close(5)
- log.debug("metric.values = %s", metrics.values())
- return metrics.values()
-
- @staticmethod
- async def _generate_vca_vdu_name(vdu_name) -> str:
- """
- Replaces all digits in vdu name for corresponding ascii characters. This is the format required by N2VC.
- :param vdu_name: Vdu name according to the vdur
- :return: Name with digits replaced with characters
- """
- vnf_name_vca = ''.join(
- ascii_lowercase[int(char)] if char.isdigit() else char for char in vdu_name)
- vnf_name_vca = re.sub(r'-[a-z]+$', '', vnf_name_vca)
- return vnf_name_vca
-
- @staticmethod
- async def _generate_read_metric_payload(metric_name, nsr_id, vdu_name, vnf_member_index) -> dict:
- """
- Builds JSON payload for asking for a metric measurement in MON. It follows the model defined in core.models.
- :param metric_name: OSM metric name (e.g.: cpu_utilization)
- :param nsr_id: NSR ID
- :param vdu_name: Vdu name according to the vdur
- :param vnf_member_index: Index of the VNF in the NS according to the vnfr
- :return: JSON payload as dict
- """
- cor_id = random.randint(1, 10e7)
- payload = {
- 'correlation_id': cor_id,
- 'metric_name': metric_name,
- 'ns_id': nsr_id,
- 'vnf_member_index': vnf_member_index,
- 'vdu_name': vdu_name,
- 'collection_period': 1,
- 'collection_unit': 'DAY',
- }
- return payload
+ vnf_member_index = vnfr['member-vnf-index-ref']
+ vim_account_id = self.common_db.get_vim_account_id(nsr_id, vnf_member_index)
+ p = multiprocessing.Process(target=self._collect_vim_metrics,
+ args=(vnfr, vim_account_id))
+ processes.append(p)
+ p.start()
+ p = multiprocessing.Process(target=self._collect_vca_metrics,
+ args=(vnfr,))
+ processes.append(p)
+ p.start()
+ vims = self.common_db.get_vim_accounts()
+ for vim in vims:
+ p = multiprocessing.Process(target=self._collect_vim_infra_metrics,
+ args=(vim['_id'],))
+ processes.append(p)
+ p.start()
+ sdncs = self.common_db.get_sdncs()
+ for sdnc in sdncs:
+ p = multiprocessing.Process(target=self._collect_sdnc_infra_metrics,
+ args=(sdnc['_id'],))
+ processes.append(p)
+ p.start()
+ for process in processes:
+ process.join(timeout=10)
+ metrics = []
+ while not self.queue.empty():
+ metrics.append(self.queue.get())
+ for plugin in self.plugins:
+ plugin.handle(metrics)
+
+ def _init_backends(self):
+ for backend in METRIC_BACKENDS:
+ self.plugins.append(backend())