Loading src/analytics/backend/service/AnalyticsBackendService.py +16 −8 Original line number Diff line number Diff line Loading @@ -33,11 +33,18 @@ class AnalyticsBackendService(GenericGrpcService): 'group.id' : 'analytics-frontend', 'auto.offset.reset' : 'latest'}) def RunSparkStreamer(self, kpi_list, oper_list, thresholds_dict): print ("Received parameters: {:} - {:} - {:}".format(kpi_list, oper_list, thresholds_dict)) LOGGER.debug ("Received parameters: {:} - {:} - {:}".format(kpi_list, oper_list, thresholds_dict)) def RunSparkStreamer(self, analyzer): kpi_list = analyzer['input_kpis'] oper_list = [s.replace('_value', '') for s in list(analyzer["thresholds"].keys())] thresholds = analyzer['thresholds'] window_size = analyzer['window_size'] window_slider = analyzer['window_slider'] print ("Received parameters: {:} - {:} - {:} - {:} - {:}".format( kpi_list, oper_list, thresholds, window_size, window_slider)) LOGGER.debug ("Received parameters: {:} - {:} - {:} - {:} - {:}".format( kpi_list, oper_list, thresholds, window_size, window_slider)) threading.Thread(target=SparkStreamer, args=(kpi_list, oper_list, None, None, thresholds_dict, None) args=(kpi_list, oper_list, thresholds, window_size, window_slider, None) ).start() return True Loading @@ -63,6 +70,7 @@ class AnalyticsBackendService(GenericGrpcService): break analyzer = json.loads(receive_msg.value().decode('utf-8')) analyzer_id = receive_msg.key().decode('utf-8') LOGGER.debug('Recevied Collector: {:} - {:}'.format(analyzer_id, analyzer)) print('Recevied Collector: {:} - {:} - {:}'.format(analyzer_id, analyzer, analyzer['input_kpis'])) self.RunSparkStreamer(analyzer['input_kpis']) # TODO: Add active analyzer to list LOGGER.debug('Recevied Analyzer: {:} - {:}'.format(analyzer_id, analyzer)) print('Recevied Analyzer: {:} - {:}'.format(analyzer_id, analyzer)) # TODO: Add active analyzer to list self.RunSparkStreamer(analyzer) src/analytics/backend/service/SparkStreaming.py +2 −1 Original line number Diff line number Diff line Loading @@ -73,7 +73,8 @@ def ApplyThresholds(aggregated_df, thresholds): ) return aggregated_df def SparkStreamer(kpi_list, oper_list, thresholds, window_size=None, win_slide_duration=None, time_stamp_col=None): def SparkStreamer(kpi_list, oper_list, thresholds, window_size=None, win_slide_duration=None, time_stamp_col=None): """ Method to perform Spark operation Kafka stream. NOTE: Kafka topic to be processesd should have atleast one row before initiating the spark session. Loading src/analytics/backend/tests/messages.py +1 −1 Original line number Diff line number Diff line Loading @@ -26,7 +26,7 @@ def get_threshold_dict(): 'max_value' : (45, 50), 'first_value' : (00, 10), 'last_value' : (40, 50), 'stddev_value' : (00, 10), 'stdev_value' : (00, 10), } # Filter threshold_dict based on the operation_list return { Loading src/analytics/frontend/service/AnalyticsFrontendServiceServicerImpl.py +8 −4 Original line number Diff line number Diff line Loading @@ -67,7 +67,11 @@ class AnalyticsFrontendServiceServicerImpl(AnalyticsFrontendServiceServicer): "algo_name" : analyzer_obj.algorithm_name, "input_kpis" : [k.kpi_id.uuid for k in analyzer_obj.input_kpi_ids], "output_kpis" : [k.kpi_id.uuid for k in analyzer_obj.output_kpi_ids], "oper_mode" : analyzer_obj.operation_mode "oper_mode" : analyzer_obj.operation_mode, "thresholds" : json.loads(analyzer_obj.parameters["thresholds"]), "window_size" : analyzer_obj.parameters["window_size"], "window_slider" : analyzer_obj.parameters["window_slider"], # "store_aggregate" : analyzer_obj.parameters["store_aggregate"] } self.kafka_producer.produce( KafkaTopic.ANALYTICS_REQUEST.value, Loading src/analytics/frontend/tests/messages.py +4 −4 Original line number Diff line number Diff line Loading @@ -33,10 +33,10 @@ def create_analyzer(): _kpi_id = KpiId() # input IDs to analyze _kpi_id.kpi_id.uuid = str(uuid.uuid4()) _kpi_id.kpi_id.uuid = "1e22f180-ba28-4641-b190-2287bf446666" _kpi_id.kpi_id.uuid = "6e22f180-ba28-4641-b190-2287bf448888" _create_analyzer.input_kpi_ids.append(_kpi_id) _kpi_id.kpi_id.uuid = str(uuid.uuid4()) _kpi_id.kpi_id.uuid = "6e22f180-ba28-4641-b190-2287bf448888" _kpi_id.kpi_id.uuid = "1e22f180-ba28-4641-b190-2287bf446666" _create_analyzer.input_kpi_ids.append(_kpi_id) _kpi_id.kpi_id.uuid = str(uuid.uuid4()) _create_analyzer.input_kpi_ids.append(_kpi_id) Loading @@ -47,8 +47,8 @@ def create_analyzer(): _create_analyzer.output_kpi_ids.append(_kpi_id) # parameter _threshold_dict = { 'avg_value' :(20, 30), 'min_value' :(00, 10), 'max_value' :(45, 50), 'first_value' :(00, 10), 'last_value' :(40, 50), 'stddev_value':(00, 10)} # 'avg_value' :(20, 30), 'min_value' :(00, 10), 'max_value' :(45, 50), 'first_value' :(00, 10), 'last_value' :(40, 50), 'stdev_value':(00, 10)} _create_analyzer.parameters['thresholds'] = json.dumps(_threshold_dict) _create_analyzer.parameters['window_size'] = "60 seconds" # Such as "10 seconds", "2 minutes", "3 hours", "4 days" or "5 weeks" _create_analyzer.parameters['window_slider'] = "30 seconds" # should be less than window size Loading Loading
src/analytics/backend/service/AnalyticsBackendService.py +16 −8 Original line number Diff line number Diff line Loading @@ -33,11 +33,18 @@ class AnalyticsBackendService(GenericGrpcService): 'group.id' : 'analytics-frontend', 'auto.offset.reset' : 'latest'}) def RunSparkStreamer(self, kpi_list, oper_list, thresholds_dict): print ("Received parameters: {:} - {:} - {:}".format(kpi_list, oper_list, thresholds_dict)) LOGGER.debug ("Received parameters: {:} - {:} - {:}".format(kpi_list, oper_list, thresholds_dict)) def RunSparkStreamer(self, analyzer): kpi_list = analyzer['input_kpis'] oper_list = [s.replace('_value', '') for s in list(analyzer["thresholds"].keys())] thresholds = analyzer['thresholds'] window_size = analyzer['window_size'] window_slider = analyzer['window_slider'] print ("Received parameters: {:} - {:} - {:} - {:} - {:}".format( kpi_list, oper_list, thresholds, window_size, window_slider)) LOGGER.debug ("Received parameters: {:} - {:} - {:} - {:} - {:}".format( kpi_list, oper_list, thresholds, window_size, window_slider)) threading.Thread(target=SparkStreamer, args=(kpi_list, oper_list, None, None, thresholds_dict, None) args=(kpi_list, oper_list, thresholds, window_size, window_slider, None) ).start() return True Loading @@ -63,6 +70,7 @@ class AnalyticsBackendService(GenericGrpcService): break analyzer = json.loads(receive_msg.value().decode('utf-8')) analyzer_id = receive_msg.key().decode('utf-8') LOGGER.debug('Recevied Collector: {:} - {:}'.format(analyzer_id, analyzer)) print('Recevied Collector: {:} - {:} - {:}'.format(analyzer_id, analyzer, analyzer['input_kpis'])) self.RunSparkStreamer(analyzer['input_kpis']) # TODO: Add active analyzer to list LOGGER.debug('Recevied Analyzer: {:} - {:}'.format(analyzer_id, analyzer)) print('Recevied Analyzer: {:} - {:}'.format(analyzer_id, analyzer)) # TODO: Add active analyzer to list self.RunSparkStreamer(analyzer)
src/analytics/backend/service/SparkStreaming.py +2 −1 Original line number Diff line number Diff line Loading @@ -73,7 +73,8 @@ def ApplyThresholds(aggregated_df, thresholds): ) return aggregated_df def SparkStreamer(kpi_list, oper_list, thresholds, window_size=None, win_slide_duration=None, time_stamp_col=None): def SparkStreamer(kpi_list, oper_list, thresholds, window_size=None, win_slide_duration=None, time_stamp_col=None): """ Method to perform Spark operation Kafka stream. NOTE: Kafka topic to be processesd should have atleast one row before initiating the spark session. Loading
src/analytics/backend/tests/messages.py +1 −1 Original line number Diff line number Diff line Loading @@ -26,7 +26,7 @@ def get_threshold_dict(): 'max_value' : (45, 50), 'first_value' : (00, 10), 'last_value' : (40, 50), 'stddev_value' : (00, 10), 'stdev_value' : (00, 10), } # Filter threshold_dict based on the operation_list return { Loading
src/analytics/frontend/service/AnalyticsFrontendServiceServicerImpl.py +8 −4 Original line number Diff line number Diff line Loading @@ -67,7 +67,11 @@ class AnalyticsFrontendServiceServicerImpl(AnalyticsFrontendServiceServicer): "algo_name" : analyzer_obj.algorithm_name, "input_kpis" : [k.kpi_id.uuid for k in analyzer_obj.input_kpi_ids], "output_kpis" : [k.kpi_id.uuid for k in analyzer_obj.output_kpi_ids], "oper_mode" : analyzer_obj.operation_mode "oper_mode" : analyzer_obj.operation_mode, "thresholds" : json.loads(analyzer_obj.parameters["thresholds"]), "window_size" : analyzer_obj.parameters["window_size"], "window_slider" : analyzer_obj.parameters["window_slider"], # "store_aggregate" : analyzer_obj.parameters["store_aggregate"] } self.kafka_producer.produce( KafkaTopic.ANALYTICS_REQUEST.value, Loading
src/analytics/frontend/tests/messages.py +4 −4 Original line number Diff line number Diff line Loading @@ -33,10 +33,10 @@ def create_analyzer(): _kpi_id = KpiId() # input IDs to analyze _kpi_id.kpi_id.uuid = str(uuid.uuid4()) _kpi_id.kpi_id.uuid = "1e22f180-ba28-4641-b190-2287bf446666" _kpi_id.kpi_id.uuid = "6e22f180-ba28-4641-b190-2287bf448888" _create_analyzer.input_kpi_ids.append(_kpi_id) _kpi_id.kpi_id.uuid = str(uuid.uuid4()) _kpi_id.kpi_id.uuid = "6e22f180-ba28-4641-b190-2287bf448888" _kpi_id.kpi_id.uuid = "1e22f180-ba28-4641-b190-2287bf446666" _create_analyzer.input_kpi_ids.append(_kpi_id) _kpi_id.kpi_id.uuid = str(uuid.uuid4()) _create_analyzer.input_kpi_ids.append(_kpi_id) Loading @@ -47,8 +47,8 @@ def create_analyzer(): _create_analyzer.output_kpi_ids.append(_kpi_id) # parameter _threshold_dict = { 'avg_value' :(20, 30), 'min_value' :(00, 10), 'max_value' :(45, 50), 'first_value' :(00, 10), 'last_value' :(40, 50), 'stddev_value':(00, 10)} # 'avg_value' :(20, 30), 'min_value' :(00, 10), 'max_value' :(45, 50), 'first_value' :(00, 10), 'last_value' :(40, 50), 'stdev_value':(00, 10)} _create_analyzer.parameters['thresholds'] = json.dumps(_threshold_dict) _create_analyzer.parameters['window_size'] = "60 seconds" # Such as "10 seconds", "2 minutes", "3 hours", "4 days" or "5 weeks" _create_analyzer.parameters['window_slider'] = "30 seconds" # should be less than window size Loading