Loading src/analytics/backend/service/SparkStreaming.py +4 −4 Original line number Diff line number Diff line Loading @@ -44,7 +44,7 @@ def DefiningRequestSchema(): StructField("kpi_value" , DoubleType() , True) ]) def get_aggregations(oper_list): def GetAggregations(oper_list): # Define the possible aggregation functions agg_functions = { 'avg' : round(avg ("kpi_value"), 3) .alias("avg_value"), Loading @@ -56,7 +56,7 @@ def get_aggregations(oper_list): } return [agg_functions[op] for op in oper_list if op in agg_functions] # Filter and return only the selected aggregations def apply_thresholds(aggregated_df, thresholds): def ApplyThresholds(aggregated_df, thresholds): # Apply thresholds (TH-Fail and TH-RAISE) based on the thresholds dictionary on the aggregated DataFrame. # Loop through each column name and its associated thresholds Loading Loading @@ -114,9 +114,9 @@ def SparkStreamer(kpi_list, oper_list, window_size=None, win_slide_duration=None ), col("kpi_id") ) \ .agg(*get_aggregations(oper_list)) .agg(*GetAggregations(oper_list)) # Apply thresholds to the aggregated data thresholded_stream_data = apply_thresholds(windowed_stream_data, thresholds) thresholded_stream_data = ApplyThresholds(windowed_stream_data, thresholds) # --- This will write output on console: FOR TESTING PURPOSES # Start the Spark streaming query Loading Loading
src/analytics/backend/service/SparkStreaming.py +4 −4 Original line number Diff line number Diff line Loading @@ -44,7 +44,7 @@ def DefiningRequestSchema(): StructField("kpi_value" , DoubleType() , True) ]) def get_aggregations(oper_list): def GetAggregations(oper_list): # Define the possible aggregation functions agg_functions = { 'avg' : round(avg ("kpi_value"), 3) .alias("avg_value"), Loading @@ -56,7 +56,7 @@ def get_aggregations(oper_list): } return [agg_functions[op] for op in oper_list if op in agg_functions] # Filter and return only the selected aggregations def apply_thresholds(aggregated_df, thresholds): def ApplyThresholds(aggregated_df, thresholds): # Apply thresholds (TH-Fail and TH-RAISE) based on the thresholds dictionary on the aggregated DataFrame. # Loop through each column name and its associated thresholds Loading Loading @@ -114,9 +114,9 @@ def SparkStreamer(kpi_list, oper_list, window_size=None, win_slide_duration=None ), col("kpi_id") ) \ .agg(*get_aggregations(oper_list)) .agg(*GetAggregations(oper_list)) # Apply thresholds to the aggregated data thresholded_stream_data = apply_thresholds(windowed_stream_data, thresholds) thresholded_stream_data = ApplyThresholds(windowed_stream_data, thresholds) # --- This will write output on console: FOR TESTING PURPOSES # Start the Spark streaming query Loading