diff --git a/.gitignore b/.gitignore index c510e7ea..4fb279b1 100644 --- a/.gitignore +++ b/.gitignore @@ -16,6 +16,7 @@ __pycache__/* *.out kube-burner* kube_burner* +recommender_*.json # Project files .ropeproject diff --git a/config/recommender_config.yaml b/config/recommender_config.yaml index 5c235faa..4bf60ec0 100644 --- a/config/recommender_config.yaml +++ b/config/recommender_config.yaml @@ -7,6 +7,8 @@ auth_token: scrape_duration: 10m chaos_library: "kraken" log_level: INFO +json_output_file: False +json_output_folder_path: # for output purpose only do not change if not needed chaos_tests: @@ -26,4 +28,8 @@ chaos_tests: - pod_network_chaos MEM: - node_memory_hog - - pvc_disk_fill \ No newline at end of file + - pvc_disk_fill + +threshold: .7 +cpu_threshold: .5 +mem_threshold: .5 diff --git a/kraken/chaos_recommender/analysis.py b/kraken/chaos_recommender/analysis.py index 5a6267da..270f7e9b 100644 --- a/kraken/chaos_recommender/analysis.py +++ b/kraken/chaos_recommender/analysis.py @@ -6,7 +6,8 @@ import time KRAKEN_TESTS_PATH = "./kraken_chaos_tests.txt" -#Placeholder, this should be done with topology + +# Placeholder, this should be done with topology def return_critical_services(): return ["web", "cart"] @@ -15,6 +16,7 @@ def load_telemetry_data(file_path): data = pd.read_csv(file_path, delimiter=r"\s+") return data + def calculate_zscores(data): zscores = pd.DataFrame() zscores["Service"] = data["service"] @@ -23,6 +25,7 @@ def calculate_zscores(data): zscores["Network"] = (data["NETWORK"] - data["NETWORK"].mean()) / data["NETWORK"].std() return zscores + def identify_outliers(data, threshold): outliers_cpu = data[data["CPU"] > threshold]["Service"].tolist() outliers_memory = data[data["Memory"] > threshold]["Service"].tolist() @@ -45,42 +48,62 @@ def get_services_above_heatmap_threshold(dataframe, cpu_threshold, mem_threshold def analysis(file_path, chaos_tests_config, threshold, heatmap_cpu_threshold, heatmap_mem_threshold): # Load the telemetry data from file + logging.info("Fetching the Telemetry data") data = load_telemetry_data(file_path) # Calculate Z-scores for CPU, Memory, and Network columns zscores = calculate_zscores(data) # Identify outliers + logging.info("Identifying outliers") outliers_cpu, outliers_memory, outliers_network = identify_outliers(zscores, threshold) cpu_services, mem_services = get_services_above_heatmap_threshold(data, heatmap_cpu_threshold, heatmap_mem_threshold) - # Display the identified outliers - logging.info("======================== Profiling ==================================") - logging.info(f"CPU outliers: {outliers_cpu}") - logging.info(f"Memory outliers: {outliers_memory}") - logging.info(f"Network outliers: {outliers_network}") - logging.info("===================== HeatMap Analysis ==============================") + analysis_data = analysis_json(outliers_cpu, outliers_memory, + outliers_network, cpu_services, + mem_services, chaos_tests_config) + + if not cpu_services: + logging.info("There are no services that are using significant CPU compared to their assigned limits (infinite in case no limits are set).") + if not mem_services: + logging.info("There are no services that are using significant MEMORY compared to their assigned limits (infinite in case no limits are set).") + time.sleep(2) + + logging.info("Please check data in utilisation.txt for further analysis") + + return analysis_data + + +def analysis_json(outliers_cpu, outliers_memory, outliers_network, + cpu_services, mem_services, chaos_tests_config): + + profiling = { + "cpu_outliers": outliers_cpu, + "memory_outliers": outliers_memory, + "network_outliers": outliers_network + } + + heatmap = { + "services_with_cpu_heatmap_above_threshold": cpu_services, + "services_with_mem_heatmap_above_threshold": mem_services + } + + recommendations = {} if cpu_services: - logging.info("Services with CPU_HEATMAP above threshold:", cpu_services) - else: - logging.info("There are no services that are using siginificant CPU compared to their assigned limits (infinite in case no limits are set).") + cpu_recommend = {"services": cpu_services, + "tests": chaos_tests_config['CPU']} + recommendations["cpu_services_recommendations"] = cpu_recommend + if mem_services: - logging.info("Services with MEM_HEATMAP above threshold:", mem_services) - else: - logging.info("There are no services that are using siginificant MEMORY compared to their assigned limits (infinite in case no limits are set).") - time.sleep(2) - logging.info("======================= Recommendations =============================") - if cpu_services: - logging.info(f"Recommended tests for {str(cpu_services)} :\n {chaos_tests_config['CPU']}") - logging.info("\n") - if mem_services: - logging.info(f"Recommended tests for {str(mem_services)} :\n {chaos_tests_config['MEM']}") - logging.info("\n") + mem_recommend = {"services": mem_services, + "tests": chaos_tests_config['MEM']} + recommendations["mem_services_recommendations"] = mem_recommend if outliers_network: - logging.info(f"Recommended tests for str(outliers_network) :\n {chaos_tests_config['NETWORK']}") - logging.info("\n") + outliers_network_recommend = {"outliers_networks": outliers_network, + "tests": chaos_tests_config['NETWORK']} + recommendations["outliers_network_recommendations"] = ( + outliers_network_recommend) - logging.info("\n") - logging.info("Please check data in utilisation.txt for further analysis") + return [profiling, heatmap, recommendations] diff --git a/kraken/chaos_recommender/prometheus.py b/kraken/chaos_recommender/prometheus.py index ba9d913e..e45cb843 100644 --- a/kraken/chaos_recommender/prometheus.py +++ b/kraken/chaos_recommender/prometheus.py @@ -1,6 +1,5 @@ import logging -import pandas from prometheus_api_client import PrometheusConnect import pandas as pd import urllib3 @@ -8,6 +7,7 @@ import urllib3 saved_metrics_path = "./utilisation.txt" + def convert_data_to_dataframe(data, label): df = pd.DataFrame() df['service'] = [item['metric']['pod'] for item in data] @@ -25,6 +25,7 @@ def convert_data(data, service): result[pod_name] = value return result.get(service, '100000000000') # for those pods whose limits are not defined they can take as much resources, there assigning a very high value + def save_utilization_to_file(cpu_data, cpu_limits_result, mem_data, mem_limits_result, network_data, filename): df_cpu = convert_data_to_dataframe(cpu_data, "CPU") merged_df = pd.DataFrame(columns=['service','CPU','CPU_LIMITS','MEM','MEM_LIMITS','NETWORK']) @@ -39,8 +40,6 @@ def save_utilization_to_file(cpu_data, cpu_limits_result, mem_data, mem_limits_r "NETWORK" : convert_data(network_data, s)}, index=[0]) merged_df = pd.concat([merged_df, new_row_df], ignore_index=True) - - # Convert columns to string merged_df['CPU'] = merged_df['CPU'].astype(str) merged_df['MEM'] = merged_df['MEM'].astype(str) @@ -57,40 +56,39 @@ def save_utilization_to_file(cpu_data, cpu_limits_result, mem_data, mem_limits_r merged_df.to_csv(filename, sep='\t', index=False) + def fetch_utilization_from_prometheus(prometheus_endpoint, auth_token, namespace, scrape_duration): urllib3.disable_warnings() prometheus = PrometheusConnect(url=prometheus_endpoint, headers={'Authorization':'Bearer {}'.format(auth_token)}, disable_ssl=True) # Fetch CPU utilization + logging.info("Fetching utilization") cpu_query = 'sum (rate (container_cpu_usage_seconds_total{image!="", namespace="%s"}[%s])) by (pod) *1000' % (namespace,scrape_duration) - logging.info(cpu_query) cpu_result = prometheus.custom_query(cpu_query) - cpu_data = cpu_result - cpu_limits_query = '(sum by (pod) (kube_pod_container_resource_limits{resource="cpu", namespace="%s"}))*1000' %(namespace) - logging.info(cpu_limits_query) cpu_limits_result = prometheus.custom_query(cpu_limits_query) - mem_query = 'sum by (pod) (avg_over_time(container_memory_usage_bytes{image!="", namespace="%s"}[%s]))' % (namespace, scrape_duration) - logging.info(mem_query) mem_result = prometheus.custom_query(mem_query) - mem_data = mem_result mem_limits_query = 'sum by (pod) (kube_pod_container_resource_limits{resource="memory", namespace="%s"}) ' %(namespace) - logging.info(mem_limits_query) mem_limits_result = prometheus.custom_query(mem_limits_query) - network_query = 'sum by (pod) ((avg_over_time(container_network_transmit_bytes_total{namespace="%s"}[%s])) + \ (avg_over_time(container_network_receive_bytes_total{namespace="%s"}[%s])))' % (namespace, scrape_duration, namespace, scrape_duration) network_result = prometheus.custom_query(network_query) - logging.info(network_query) - network_data = network_result - - save_utilization_to_file(cpu_data, cpu_limits_result, mem_data, mem_limits_result, network_data, saved_metrics_path) - return saved_metrics_path + save_utilization_to_file(cpu_result, cpu_limits_result, mem_result, mem_limits_result, network_result, saved_metrics_path) + queries = json_queries(cpu_query, cpu_limits_query, mem_query, mem_limits_query) + return saved_metrics_path, queries +def json_queries(cpu_query, cpu_limits_query, mem_query, mem_limits_query): + queries = { + "cpu_query": cpu_query, + "cpu_limit_query": cpu_limits_query, + "memory_query": mem_query, + "memory_limit_query": mem_limits_query + } + return queries diff --git a/utils/chaos_recommender/README.md b/utils/chaos_recommender/README.md index 58b38cba..cc59fc70 100644 --- a/utils/chaos_recommender/README.md +++ b/utils/chaos_recommender/README.md @@ -39,6 +39,8 @@ You can customize the default values by editing the `krkn/config/recommender_con - `auth_token`: Auth token to connect to prometheus endpoint (must). - `scrape_duration`: For how long data should be fetched, e.g., '1m' (must). - `chaos_library`: "kraken" (currently it only supports kraken). + - `json_output_file`: True or False (by default False). + - `json_output_folder_path`: Specify folder path where output should be saved. If empty the default path is used. - `chaos_tests`: (for output purpose only do not change if not needed) - `GENERAL`: list of general purpose tests available in Krkn - `MEM`: list of memory related tests available in Krkn @@ -79,6 +81,8 @@ You can also provide the input values through command-line arguments launching t Chaos library -L LOG_LEVEL, --log-level LOG_LEVEL log level (DEBUG, INFO, WARNING, ERROR, CRITICAL + -J [FOLDER_PATH], --json-output-file [FOLDER_PATH] + Create output file, the path to the folder can be specified, if not specified the default folder is used. -M MEM [MEM ...], --MEM MEM [MEM ...] Memory related chaos tests (space separated list) -C CPU [CPU ...], --CPU CPU [CPU ...] diff --git a/utils/chaos_recommender/chaos_recommender.py b/utils/chaos_recommender/chaos_recommender.py index d7565e51..40eeb275 100644 --- a/utils/chaos_recommender/chaos_recommender.py +++ b/utils/chaos_recommender/chaos_recommender.py @@ -1,7 +1,9 @@ import argparse +import json import logging import os.path import sys +import time import yaml # kraken module import for running the recommender # both from the root directory and the recommender @@ -9,12 +11,13 @@ import yaml sys.path.insert(0, './') sys.path.insert(0, '../../') +from krkn_lib.utils import get_yaml_item_value + import kraken.chaos_recommender.analysis as analysis import kraken.chaos_recommender.prometheus as prometheus from kubernetes import config as kube_config - def parse_arguments(parser): # command line options @@ -26,6 +29,10 @@ def parse_arguments(parser): parser.add_argument("-t", "--token", action="store", default="", help="Kubernetes authentication token") parser.add_argument("-s", "--scrape-duration", action="store", default="10m", help="Prometheus scrape duration") parser.add_argument("-L", "--log-level", action="store", default="INFO", help="log level (DEBUG, INFO, WARNING, ERROR, CRITICAL") + + parser.add_argument("-J", "--json-output-file", default=False, nargs="?", action="store", + help="Create output file, the path to the folder can be specified, if not specified the default folder is used") + parser.add_argument("-M", "--MEM", nargs='+', action="store", default=[], help="Memory related chaos tests (space separated list)") parser.add_argument("-C", "--CPU", nargs='+', action="store", default=[], @@ -35,11 +42,12 @@ def parse_arguments(parser): parser.add_argument("-G", "--GENERIC", nargs='+', action="store", default=[], help="Memory related chaos tests (space separated list)") parser.add_argument("--threshold", action="store", default="", help="Threshold") - parser.add_argument("--cpu_threshold", action="store", default="", help="CPU threshold") - parser.add_argument("--mem_threshold", action="store", default="", help="Memory threshold") + parser.add_argument("--cpu-threshold", action="store", default="", help="CPU threshold") + parser.add_argument("--mem-threshold", action="store", default="", help="Memory threshold") return parser.parse_args() + def read_configuration(config_file_path): if not os.path.exists(config_file_path): logging.error(f"Config file not found: {config_file_path}") @@ -49,18 +57,25 @@ def read_configuration(config_file_path): config = yaml.safe_load(config_file) log_level = config.get("log level", "INFO") - namespace = config.get("namespace", "") - kubeconfig = config.get("kubeconfig", kube_config.KUBE_CONFIG_DEFAULT_LOCATION) + namespace = config.get("namespace") + kubeconfig = get_yaml_item_value(config, "kubeconfig", kube_config.KUBE_CONFIG_DEFAULT_LOCATION) - prometheus_endpoint = config.get("prometheus_endpoint", "") - auth_token = config.get("auth_token", "") - scrape_duration = config.get("scrape_duration", "10m") - chaos_tests = config.get("chaos_tests" , {}) - threshold = config.get("threshold", ".7") - heatmap_cpu_threshold = config.get("cpu_threshold", ".5") - heatmap_mem_threshold = config.get("mem_threshold", ".5") + prometheus_endpoint = config.get("prometheus_endpoint") + auth_token = config.get("auth_token") + scrape_duration = get_yaml_item_value(config, "scrape_duration", "10m") + threshold = get_yaml_item_value(config, "threshold", ".7") + heatmap_cpu_threshold = get_yaml_item_value(config, "cpu_threshold", ".5") + heatmap_mem_threshold = get_yaml_item_value(config, "mem_threshold", ".3") + output_file = config.get("json_output_file", False) + if output_file is True: + output_path = config.get("json_output_folder_path") + else: + output_path = False + chaos_tests = config.get("chaos_tests", {}) return (namespace, kubeconfig, prometheus_endpoint, auth_token, scrape_duration, - chaos_tests, log_level, threshold, heatmap_cpu_threshold, heatmap_mem_threshold) + chaos_tests, log_level, threshold, heatmap_cpu_threshold, + heatmap_mem_threshold, output_path) + def prompt_input(prompt, default_value): user_input = input(f"{prompt} [{default_value}]: ") @@ -68,6 +83,44 @@ def prompt_input(prompt, default_value): return user_input return default_value + +def make_json_output(inputs, queries, analysis_data, output_path): + time_str = time.strftime("%Y-%m-%d_%H-%M-%S", time.localtime()) + + data = { + "inputs": inputs, + "queries": queries, + "profiling": analysis_data[0], + "heatmap_analysis": analysis_data[1], + "recommendations": analysis_data[2] + } + + logging.info(f"Summary\n{json.dumps(data, indent=4)}") + + if output_path is not False: + file = f"recommender_{inputs['namespace']}_{time_str}.json" + path = f"{os.path.expanduser(output_path)}/{file}" + + with open(path, "w") as json_output: + logging.info(f"Saving output file in {output_path} folder...") + json_output.write(json.dumps(data, indent=4)) + logging.info(f"Recommendation output saved in {file}.") + + +def json_inputs(namespace, kubeconfig, prometheus_endpoint, scrape_duration, chaos_tests, threshold, heatmap_cpu_threshold, heatmap_mem_threshold): + inputs = { + "namespace": namespace, + "kubeconfig": kubeconfig, + "prometheus_endpoint": prometheus_endpoint, + "scrape_duration": scrape_duration, + "chaos_tests": chaos_tests, + "threshold": threshold, + "heatmap_cpu_threshold": heatmap_cpu_threshold, + "heatmap_mem_threshold": heatmap_mem_threshold + } + return inputs + + def main(): parser = argparse.ArgumentParser(description="Krkn Chaos Recommender Command-Line tool") args = parse_arguments(parser) @@ -88,7 +141,8 @@ def main(): log_level, threshold, heatmap_cpu_threshold, - heatmap_mem_threshold + heatmap_mem_threshold, + output_path ) = read_configuration(args.config_file) if args.options: @@ -98,30 +152,35 @@ def main(): scrape_duration = args.scrape_duration log_level = args.log_level prometheus_endpoint = args.prometheus_endpoint + output_path = args.json_output_file chaos_tests = {"MEM": args.MEM, "GENERIC": args.GENERIC, "CPU": args.CPU, "NETWORK": args.NETWORK} threshold = args.threshold - heatmap_mem_threshold = args.heatmap_mem_threshold - heatmap_cpu_threshold = args.heatmap_cpu_threshold + heatmap_mem_threshold = args.mem_threshold + heatmap_cpu_threshold = args.cpu_threshold - if log_level not in ["DEBUG","INFO", "WARNING", "ERROR","CRITICAL"]: + if log_level not in ["DEBUG", "INFO", "WARNING", "ERROR", "CRITICAL"]: logging.error(f"{log_level} not a valid log level") sys.exit(1) logging.basicConfig(level=log_level) - logging.info("============================INPUTS===================================") - logging.info(f"Namespace: {namespace}") - logging.info(f"Kubeconfig: {kubeconfig}") - logging.info(f"Prometheus endpoint: {prometheus_endpoint}") - logging.info(f"Scrape duration: {scrape_duration}") - for test in chaos_tests.keys(): - logging.info(f"Chaos tests {test}: {chaos_tests[test]}") - logging.info("=====================================================================") + if output_path is not False: + if output_path is None: + output_path = "./recommender_output" + logging.info(f"Path for output file not specified. " + f"Using default folder {output_path}") + if not os.path.exists(os.path.expanduser(output_path)): + logging.error(f"Folder {output_path} for output not found.") + sys.exit(1) + logging.info("Loading inputs...") + inputs = json_inputs(namespace, kubeconfig, prometheus_endpoint, scrape_duration, chaos_tests, threshold, heatmap_cpu_threshold, heatmap_mem_threshold) logging.info("Starting Analysis ...") - logging.info("Fetching the Telemetry data") - file_path = prometheus.fetch_utilization_from_prometheus(prometheus_endpoint, auth_token, namespace, scrape_duration) - analysis(file_path, chaos_tests, threshold, heatmap_cpu_threshold, heatmap_mem_threshold) + file_path, queries = prometheus.fetch_utilization_from_prometheus(prometheus_endpoint, auth_token, namespace, scrape_duration) + analysis_data = analysis(file_path, chaos_tests, threshold, heatmap_cpu_threshold, heatmap_mem_threshold) + + make_json_output(inputs, queries, analysis_data, output_path) + if __name__ == "__main__": main()