Add mgbench tutorial (#836)
* Add Docker runner * Add Docker client * Add benchgraph.sh script * Add package script
This commit is contained in:
@@ -14,27 +14,28 @@
|
||||
import argparse
|
||||
import json
|
||||
import multiprocessing
|
||||
import pathlib
|
||||
import platform
|
||||
import random
|
||||
import sys
|
||||
|
||||
import helpers
|
||||
import log
|
||||
import runners
|
||||
import setup
|
||||
from benchmark_context import BenchmarkContext
|
||||
from workloads import *
|
||||
|
||||
WITH_FINE_GRAINED_AUTHORIZATION = "with_fine_grained_authorization"
|
||||
WITHOUT_FINE_GRAINED_AUTHORIZATION = "without_fine_grained_authorization"
|
||||
QUERY_COUNT_LOWER_BOUND = 30
|
||||
|
||||
|
||||
def parse_args():
|
||||
parser = argparse.ArgumentParser(description="Main parser.", add_help=False)
|
||||
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Memgraph benchmark executor.",
|
||||
formatter_class=argparse.ArgumentDefaultsHelpFormatter,
|
||||
)
|
||||
parser.add_argument(
|
||||
benchmark_parser = argparse.ArgumentParser(description="Benchmark arguments parser", add_help=False)
|
||||
|
||||
benchmark_parser.add_argument(
|
||||
"benchmarks",
|
||||
nargs="*",
|
||||
default=None,
|
||||
@@ -48,80 +49,65 @@ def parse_args():
|
||||
"the default group is '*' which selects all groups; the"
|
||||
"default query is '*' which selects all queries",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--vendor-binary",
|
||||
help="Vendor binary used for benchmarking, by default it is memgraph",
|
||||
default=helpers.get_binary_path("memgraph"),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--vendor-name",
|
||||
default="memgraph",
|
||||
choices=["memgraph", "neo4j"],
|
||||
help="Input vendor binary name (memgraph, neo4j)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--client-binary",
|
||||
default=helpers.get_binary_path("tests/mgbench/client"),
|
||||
help="Client binary used for benchmarking",
|
||||
)
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--num-workers-for-import",
|
||||
type=int,
|
||||
default=multiprocessing.cpu_count() // 2,
|
||||
help="number of workers used to import the dataset",
|
||||
)
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--num-workers-for-benchmark",
|
||||
type=int,
|
||||
default=1,
|
||||
help="number of workers used to execute the benchmark",
|
||||
)
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--single-threaded-runtime-sec",
|
||||
type=int,
|
||||
default=10,
|
||||
help="single threaded duration of each query",
|
||||
)
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--query-count-lower-bound",
|
||||
type=int,
|
||||
default=30,
|
||||
help="Lower bound for query count, minimum number of queries that will be executed. If approximated --single-threaded-runtime-sec query count is lower than this value, lower bound is used.",
|
||||
)
|
||||
benchmark_parser.add_argument(
|
||||
"--no-load-query-counts",
|
||||
action="store_true",
|
||||
default=False,
|
||||
help="disable loading of cached query counts",
|
||||
)
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--no-save-query-counts",
|
||||
action="store_true",
|
||||
default=False,
|
||||
help="disable storing of cached query counts",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--export-results",
|
||||
default=None,
|
||||
help="file path into which results should be exported",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--temporary-directory",
|
||||
default="/tmp",
|
||||
help="directory path where temporary data should be stored",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--no-authorization",
|
||||
action="store_false",
|
||||
default=True,
|
||||
help="Run each query with authorization",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--warm-up",
|
||||
default="cold",
|
||||
choices=["cold", "hot", "vulcanic"],
|
||||
help="Run different warmups before benchmarks sample starts",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--workload-realistic",
|
||||
nargs="*",
|
||||
type=int,
|
||||
@@ -134,7 +120,7 @@ def parse_args():
|
||||
70% read, 10% update and 0% analytical.""",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--workload-mixed",
|
||||
nargs="*",
|
||||
type=int,
|
||||
@@ -147,29 +133,66 @@ def parse_args():
|
||||
with the presence of 300 write queries from write type or 30%""",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--time-depended-execution",
|
||||
type=int,
|
||||
default=0,
|
||||
help="Execute defined number of queries (based on single-threaded-runtime-sec) for a defined duration in of wall-clock time",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--performance-tracking",
|
||||
action="store_true",
|
||||
default=False,
|
||||
help="Flag for runners performance tracking, this logs RES through time and vendor specific performance tracking.",
|
||||
)
|
||||
|
||||
parser.add_argument("--customer-workloads", default=None, help="Path to customers workloads")
|
||||
benchmark_parser.add_argument("--customer-workloads", default=None, help="Path to customers workloads")
|
||||
|
||||
parser.add_argument(
|
||||
benchmark_parser.add_argument(
|
||||
"--vendor-specific",
|
||||
nargs="*",
|
||||
default=[],
|
||||
help="Vendor specific arguments that can be applied to each vendor, format: [key=value, key=value ...]",
|
||||
)
|
||||
|
||||
subparsers = parser.add_subparsers(help="Subparsers", dest="run_option")
|
||||
|
||||
# Vendor native parser starts here
|
||||
parser_vendor_native = subparsers.add_parser(
|
||||
"vendor-native",
|
||||
help="Running database in binary native form",
|
||||
parents=[benchmark_parser],
|
||||
)
|
||||
parser_vendor_native.add_argument(
|
||||
"--vendor-name",
|
||||
default="memgraph",
|
||||
choices=["memgraph", "neo4j"],
|
||||
help="Input vendor binary name (memgraph, neo4j)",
|
||||
)
|
||||
parser_vendor_native.add_argument(
|
||||
"--vendor-binary",
|
||||
help="Vendor binary used for benchmarking, by default it is memgraph",
|
||||
default=helpers.get_binary_path("memgraph"),
|
||||
)
|
||||
|
||||
parser_vendor_native.add_argument(
|
||||
"--client-binary",
|
||||
default=helpers.get_binary_path("tests/mgbench/client"),
|
||||
help="Client binary used for benchmarking",
|
||||
)
|
||||
|
||||
# Vendor docker parsers starts here
|
||||
parser_vendor_docker = subparsers.add_parser(
|
||||
"vendor-docker", help="Running database in docker", parents=[benchmark_parser]
|
||||
)
|
||||
parser_vendor_docker.add_argument(
|
||||
"--vendor-name",
|
||||
default="memgraph",
|
||||
choices=["memgraph-docker", "neo4j-docker"],
|
||||
help="Input vendor name to run in docker (memgraph-docker, neo4j-docker)",
|
||||
)
|
||||
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
@@ -184,28 +207,28 @@ def get_queries(gen, count):
|
||||
|
||||
|
||||
def warmup(condition: str, client: runners.BaseRunner, queries: list = None):
|
||||
log.log("Database condition {} ".format(condition))
|
||||
log.init("Started warm-up procedure to match database condition: {} ".format(condition))
|
||||
if condition == "hot":
|
||||
log.log("Execute warm-up to match condition {} ".format(condition))
|
||||
log.log("Execute warm-up to match condition: {} ".format(condition))
|
||||
client.execute(
|
||||
queries=[
|
||||
("CREATE ();", {}),
|
||||
("CREATE ()-[:TempEdge]->();", {}),
|
||||
("MATCH (n) RETURN n LIMIT 1;", {}),
|
||||
("MATCH (n) RETURN count(n.prop) LIMIT 1;", {}),
|
||||
],
|
||||
num_workers=1,
|
||||
)
|
||||
elif condition == "vulcanic":
|
||||
log.log("Execute warm-up to match condition {} ".format(condition))
|
||||
log.log("Execute warm-up to match condition: {} ".format(condition))
|
||||
client.execute(queries=queries)
|
||||
else:
|
||||
log.log("No warm-up on condition {} ".format(condition))
|
||||
log.log("No warm-up on condition: {} ".format(condition))
|
||||
log.log("Finished warm-up procedure to match database condition: {} ".format(condition))
|
||||
|
||||
|
||||
def mixed_workload(
|
||||
vendor: runners.BaseRunner, client: runners.BaseClient, dataset, group, queries, benchmark_context: BenchmarkContext
|
||||
):
|
||||
|
||||
num_of_queries = benchmark_context.mode_config[0]
|
||||
percentage_distribution = benchmark_context.mode_config[1:]
|
||||
if sum(percentage_distribution) != 100:
|
||||
@@ -233,7 +256,7 @@ def mixed_workload(
|
||||
"analytical": [],
|
||||
}
|
||||
|
||||
for (_, funcname) in queries[group]:
|
||||
for _, funcname in queries[group]:
|
||||
for key in queries_by_type.keys():
|
||||
if key in funcname:
|
||||
queries_by_type[key].append(funcname)
|
||||
@@ -252,8 +275,7 @@ def mixed_workload(
|
||||
full_workload = []
|
||||
|
||||
log.info(
|
||||
"Running query in mixed workload:",
|
||||
"{}/{}/{}".format(
|
||||
"Running query in mixed workload: {}/{}/{}".format(
|
||||
group,
|
||||
query,
|
||||
funcname,
|
||||
@@ -278,7 +300,7 @@ def mixed_workload(
|
||||
additional_query = getattr(dataset, funcname)
|
||||
full_workload.append(additional_query())
|
||||
|
||||
vendor.start_benchmark(
|
||||
vendor.start_db(
|
||||
dataset.NAME + dataset.get_variant() + "_" + "mixed" + "_" + query + "_" + config_distribution
|
||||
)
|
||||
warmup(benchmark_context.warm_up, client=client)
|
||||
@@ -286,7 +308,7 @@ def mixed_workload(
|
||||
queries=full_workload,
|
||||
num_workers=benchmark_context.num_workers_for_benchmark,
|
||||
)[0]
|
||||
usage_workload = vendor.stop(
|
||||
usage_workload = vendor.stop_db(
|
||||
dataset.NAME + dataset.get_variant() + "_" + "mixed" + "_" + query + "_" + config_distribution
|
||||
)
|
||||
|
||||
@@ -313,13 +335,13 @@ def mixed_workload(
|
||||
additional_query = getattr(dataset, funcname)
|
||||
full_workload.append(additional_query())
|
||||
|
||||
vendor.start_benchmark(dataset.NAME + dataset.get_variant() + "_" + "realistic" + "_" + config_distribution)
|
||||
vendor.start_db(dataset.NAME + dataset.get_variant() + "_" + "realistic" + "_" + config_distribution)
|
||||
warmup(benchmark_context.warm_up, client=client)
|
||||
ret = client.execute(
|
||||
queries=full_workload,
|
||||
num_workers=benchmark_context.num_workers_for_benchmark,
|
||||
)[0]
|
||||
usage_workload = vendor.stop(
|
||||
usage_workload = vendor.stop_db(
|
||||
dataset.NAME + dataset.get_variant() + "_" + "realistic" + "_" + config_distribution
|
||||
)
|
||||
mixed_workload = {
|
||||
@@ -349,7 +371,6 @@ def get_query_cache_count(
|
||||
config_key: list,
|
||||
benchmark_context: BenchmarkContext,
|
||||
):
|
||||
|
||||
cached_count = config.get_value(*config_key)
|
||||
if cached_count is None:
|
||||
log.info(
|
||||
@@ -357,8 +378,9 @@ def get_query_cache_count(
|
||||
benchmark_context.single_threaded_runtime_sec
|
||||
)
|
||||
)
|
||||
log.log("Running query to prime the query cache...")
|
||||
# First run to prime the query caches.
|
||||
vendor.start_benchmark("cache")
|
||||
vendor.start_db("cache")
|
||||
client.execute(queries=queries, num_workers=1)
|
||||
# Get a sense of the runtime.
|
||||
count = 1
|
||||
@@ -379,11 +401,10 @@ def get_query_cache_count(
|
||||
break
|
||||
else:
|
||||
count = count * 10
|
||||
vendor.stop("cache")
|
||||
vendor.stop_db("cache")
|
||||
|
||||
QUERY_COUNT_LOWER_BOUND = 30
|
||||
if count < QUERY_COUNT_LOWER_BOUND:
|
||||
count = QUERY_COUNT_LOWER_BOUND
|
||||
if count < benchmark_context.query_count_lower_bound:
|
||||
count = benchmark_context.query_count_lower_bound
|
||||
|
||||
config.set_value(
|
||||
*config_key,
|
||||
@@ -394,7 +415,7 @@ def get_query_cache_count(
|
||||
)
|
||||
else:
|
||||
log.log(
|
||||
"Using cached query count of {} queries for {} seconds of single-threaded runtime.".format(
|
||||
"Using cached query count of {} queries for {} seconds of single-threaded runtime to extrapolate .".format(
|
||||
cached_count["count"], cached_count["duration"]
|
||||
),
|
||||
)
|
||||
@@ -403,14 +424,10 @@ def get_query_cache_count(
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
args = parse_args()
|
||||
vendor_specific_args = helpers.parse_kwargs(args.vendor_specific)
|
||||
|
||||
assert args.benchmarks != None, helpers.list_available_workloads()
|
||||
assert args.vendor_name == "memgraph" or args.vendor_name == "neo4j", "Unsupported vendors"
|
||||
assert args.vendor_binary != None, "Pass database binary for runner"
|
||||
assert args.client_binary != None, "Pass client binary for benchmark client "
|
||||
assert args.num_workers_for_import > 0
|
||||
assert args.num_workers_for_benchmark > 0
|
||||
assert args.export_results != None, "Pass where will results be saved"
|
||||
@@ -421,17 +438,21 @@ if __name__ == "__main__":
|
||||
args.workload_realistic == None or args.workload_mixed == None
|
||||
), "Cannot run both realistic and mixed workload, only one mode run at the time"
|
||||
|
||||
temp_dir = pathlib.Path.cwd() / ".temp"
|
||||
temp_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
benchmark_context = BenchmarkContext(
|
||||
benchmark_target_workload=args.benchmarks,
|
||||
vendor_binary=args.vendor_binary,
|
||||
vendor_name=args.vendor_name,
|
||||
client_binary=args.client_binary,
|
||||
vendor_binary=args.vendor_binary if args.run_option == "vendor-native" else None,
|
||||
vendor_name=args.vendor_name.replace("-", ""),
|
||||
client_binary=args.client_binary if args.run_option == "vendor-native" else None,
|
||||
num_workers_for_import=args.num_workers_for_import,
|
||||
num_workers_for_benchmark=args.num_workers_for_benchmark,
|
||||
single_threaded_runtime_sec=args.single_threaded_runtime_sec,
|
||||
query_count_lower_bound=args.query_count_lower_bound,
|
||||
no_load_query_counts=args.no_load_query_counts,
|
||||
export_results=args.export_results,
|
||||
temporary_directory=args.temporary_directory,
|
||||
temporary_directory=temp_dir.absolute(),
|
||||
workload_mixed=args.workload_mixed,
|
||||
workload_realistic=args.workload_realistic,
|
||||
time_dependent_execution=args.time_depended_execution,
|
||||
@@ -444,10 +465,16 @@ if __name__ == "__main__":
|
||||
|
||||
log.init("Executing benchmark with following arguments: ")
|
||||
for key, value in benchmark_context.__dict__.items():
|
||||
log.log(str(key) + " : " + str(value))
|
||||
log.log("{:<30} : {:<30}".format(str(key), str(value)))
|
||||
|
||||
log.init("Check requirements for running benchmark")
|
||||
if setup.check_requirements(benchmark_context=benchmark_context):
|
||||
log.success("Requirements satisfied... ")
|
||||
else:
|
||||
log.warning("Requirements not satisfied...")
|
||||
sys.exit(1)
|
||||
|
||||
log.log("Creating cache folder for: dataset, configurations, indexes, results etc. ")
|
||||
# Create cache, config and results objects.
|
||||
cache = helpers.Cache()
|
||||
log.init("Folder in use: " + cache.get_default_cache_directory())
|
||||
if not benchmark_context.no_load_query_counts:
|
||||
@@ -457,15 +484,11 @@ if __name__ == "__main__":
|
||||
config = helpers.RecursiveDict()
|
||||
results = helpers.RecursiveDict()
|
||||
|
||||
log.init("Creating vendor runner for DB: " + benchmark_context.vendor_name)
|
||||
vendor_runner = runners.BaseRunner.create(
|
||||
benchmark_context=benchmark_context,
|
||||
)
|
||||
log.log("Class in use: " + str(vendor_runner))
|
||||
|
||||
run_config = {
|
||||
"vendor": benchmark_context.vendor_name,
|
||||
"condition": benchmark_context.warm_up,
|
||||
"num_workers_for_benchmark": benchmark_context.num_workers_for_benchmark,
|
||||
"single_threaded_runtime_sec": benchmark_context.single_threaded_runtime_sec,
|
||||
"benchmark_mode": benchmark_context.mode,
|
||||
"benchmark_mode_config": benchmark_context.mode_config,
|
||||
"platform": platform.platform(),
|
||||
@@ -475,53 +498,78 @@ if __name__ == "__main__":
|
||||
|
||||
available_workloads = helpers.get_available_workloads(benchmark_context.customer_workloads)
|
||||
|
||||
log.init("Currently available workloads: ")
|
||||
log.log(helpers.list_available_workloads(benchmark_context.customer_workloads))
|
||||
|
||||
# Filter out the workloads based on the pattern
|
||||
target_workloads = helpers.filter_workloads(
|
||||
available_workloads=available_workloads, benchmark_context=benchmark_context
|
||||
)
|
||||
|
||||
if len(target_workloads) == 0:
|
||||
log.error("No workloads matched the pattern: " + str(benchmark_context.benchmark_target_workload))
|
||||
log.error("Please check the pattern and workload NAME property, query group and query name.")
|
||||
log.info("Currently available workloads: ")
|
||||
log.log(helpers.list_available_workloads(benchmark_context.customer_workloads))
|
||||
sys.exit(1)
|
||||
|
||||
# Run all target workloads.
|
||||
for workload, queries in target_workloads:
|
||||
log.info("Started running following workload: " + str(workload.NAME))
|
||||
|
||||
benchmark_context.set_active_workload(workload.NAME)
|
||||
benchmark_context.set_active_variant(workload.get_variant())
|
||||
|
||||
log.init("Creating vendor runner for DB: " + benchmark_context.vendor_name)
|
||||
vendor_runner = runners.BaseRunner.create(
|
||||
benchmark_context=benchmark_context,
|
||||
)
|
||||
log.log("Class in use: " + str(vendor_runner.__class__.__name__))
|
||||
|
||||
log.info("Cleaning the database from any previous data")
|
||||
vendor_runner.clean_db()
|
||||
|
||||
client = vendor_runner.fetch_client()
|
||||
log.log("Get appropriate client for vendor " + str(client))
|
||||
log.log("Get appropriate client for vendor " + str(client.__class__.__name__))
|
||||
|
||||
ret = None
|
||||
usage = None
|
||||
|
||||
log.init("Preparing workload: " + workload.NAME + "/" + workload.get_variant())
|
||||
workload.prepare(cache.cache_directory("datasets", workload.NAME, workload.get_variant()))
|
||||
generated_queries = workload.dataset_generator()
|
||||
if generated_queries:
|
||||
vendor_runner.start_preparation("import")
|
||||
print("\n")
|
||||
log.info("Using workload as dataset generator...")
|
||||
if workload.get_index():
|
||||
log.info("Using index from specified file: {}".format(workload.get_index()))
|
||||
client.execute(file_path=workload.get_index(), num_workers=benchmark_context.num_workers_for_import)
|
||||
else:
|
||||
log.warning("Make sure proper indexes/constraints are created in generated queries!")
|
||||
|
||||
vendor_runner.start_db_init("import")
|
||||
|
||||
log.warning("Using following indexes...")
|
||||
log.info(workload.indexes_generator())
|
||||
log.info("Executing database index setup...")
|
||||
ret = client.execute(queries=workload.indexes_generator(), num_workers=1)
|
||||
log.log("Finished setting up indexes...")
|
||||
for row in ret:
|
||||
log.success(
|
||||
"Executed {} queries in {} seconds using {} workers with a total throughput of {} Q/S.".format(
|
||||
row["count"], row["duration"], row["num_workers"], row["throughput"]
|
||||
)
|
||||
)
|
||||
|
||||
log.info("Importing dataset...")
|
||||
ret = client.execute(queries=generated_queries, num_workers=benchmark_context.num_workers_for_import)
|
||||
usage = vendor_runner.stop("import")
|
||||
log.log("Finished importing dataset...")
|
||||
usage = vendor_runner.stop_db_init("import")
|
||||
else:
|
||||
log.init("Preparing workload: " + workload.NAME + "/" + workload.get_variant())
|
||||
workload.prepare(cache.cache_directory("datasets", workload.NAME, workload.get_variant()))
|
||||
log.info("Using workload dataset information for import...")
|
||||
imported = workload.custom_import()
|
||||
if not imported:
|
||||
log.log("Basic import execution")
|
||||
vendor_runner.start_preparation("import")
|
||||
vendor_runner.start_db_init("import")
|
||||
log.log("Executing database index setup...")
|
||||
client.execute(file_path=workload.get_index(), num_workers=benchmark_context.num_workers_for_import)
|
||||
client.execute(file_path=workload.get_index(), num_workers=1)
|
||||
log.log("Importing dataset...")
|
||||
ret = client.execute(
|
||||
file_path=workload.get_file(), num_workers=benchmark_context.num_workers_for_import
|
||||
)
|
||||
usage = vendor_runner.stop("import")
|
||||
usage = vendor_runner.stop_db_init("import")
|
||||
else:
|
||||
log.info("Custom import executed...")
|
||||
|
||||
@@ -531,7 +579,7 @@ if __name__ == "__main__":
|
||||
# Display import statistics.
|
||||
for row in ret:
|
||||
log.success(
|
||||
"Executed {} queries in {} seconds using {} workers with a total throughput of {} + Q/S.".format(
|
||||
"Executed {} queries in {} seconds using {} workers with a total throughput of {} Q/S.".format(
|
||||
row["count"], row["duration"], row["num_workers"], row["throughput"]
|
||||
)
|
||||
)
|
||||
@@ -539,7 +587,7 @@ if __name__ == "__main__":
|
||||
log.success(
|
||||
"The database used {} seconds of CPU time and peaked at {} MiB of RAM".format(
|
||||
usage["cpu"], usage["memory"] / 1024 / 1024
|
||||
),
|
||||
)
|
||||
)
|
||||
|
||||
results.set_value(*import_key, value={"client": ret, "database": usage})
|
||||
@@ -548,6 +596,7 @@ if __name__ == "__main__":
|
||||
|
||||
# Run all benchmarks in all available groups.
|
||||
for group in sorted(queries.keys()):
|
||||
print("\n")
|
||||
log.init("Running benchmark in " + benchmark_context.mode)
|
||||
if benchmark_context.mode == "Mixed":
|
||||
mixed_workload(vendor_runner, client, workload, group, queries, benchmark_context)
|
||||
@@ -568,12 +617,13 @@ if __name__ == "__main__":
|
||||
group,
|
||||
query,
|
||||
]
|
||||
log.init("Determining query count for benchmark based on --single-threaded-runtime argument")
|
||||
count = get_query_cache_count(
|
||||
vendor_runner, client, get_queries(func, 1), config_key, benchmark_context
|
||||
)
|
||||
|
||||
# Benchmark run.
|
||||
log.info("Sample query:{}".format(get_queries(func, 1)[0][0]))
|
||||
sample_query = get_queries(func, 1)[0][0]
|
||||
log.info("Sample query:{}".format(sample_query))
|
||||
log.log(
|
||||
"Executing benchmark with {} queries that should yield a single-threaded runtime of {} seconds.".format(
|
||||
count, benchmark_context.single_threaded_runtime_sec
|
||||
@@ -584,11 +634,11 @@ if __name__ == "__main__":
|
||||
benchmark_context.num_workers_for_benchmark
|
||||
)
|
||||
)
|
||||
vendor_runner.start_benchmark(
|
||||
vendor_runner.start_db(
|
||||
workload.NAME + workload.get_variant() + "_" + "_" + benchmark_context.mode + "_" + query
|
||||
)
|
||||
|
||||
warmup(condition=benchmark_context.warm_up, client=client, queries=get_queries(func, count))
|
||||
log.init("Executing benchmark queries...")
|
||||
if benchmark_context.time_dependent_execution != 0:
|
||||
ret = client.execute(
|
||||
queries=get_queries(func, count),
|
||||
@@ -600,8 +650,8 @@ if __name__ == "__main__":
|
||||
queries=get_queries(func, count),
|
||||
num_workers=benchmark_context.num_workers_for_benchmark,
|
||||
)[0]
|
||||
|
||||
usage = vendor_runner.stop(
|
||||
log.info("Benchmark execution finished...")
|
||||
usage = vendor_runner.stop_db(
|
||||
workload.NAME + workload.get_variant() + "_" + benchmark_context.mode + "_" + query
|
||||
)
|
||||
ret["database"] = usage
|
||||
@@ -610,15 +660,28 @@ if __name__ == "__main__":
|
||||
log.log("Executed {} queries in {} seconds.".format(ret["count"], ret["duration"]))
|
||||
log.log("Queries have been retried {} times".format(ret["retries"]))
|
||||
log.log("Database used {:.3f} seconds of CPU time.".format(usage["cpu"]))
|
||||
log.log("Database peaked at {:.3f} MiB of memory.".format(usage["memory"] / 1024.0 / 1024.0))
|
||||
log.log("{:<31} {:>20} {:>20} {:>20}".format("Metadata:", "min", "avg", "max"))
|
||||
metadata = ret["metadata"]
|
||||
for key in sorted(metadata.keys()):
|
||||
log.log(
|
||||
"{name:>30}: {minimum:>20.06f} {average:>20.06f} "
|
||||
"{maximum:>20.06f}".format(name=key, **metadata[key])
|
||||
)
|
||||
log.info("Database peaked at {:.3f} MiB of memory.".format(usage["memory"] / 1024.0 / 1024.0))
|
||||
if "docker" not in benchmark_context.vendor_name:
|
||||
log.log("{:<31} {:>20} {:>20} {:>20}".format("Metadata:", "min", "avg", "max"))
|
||||
metadata = ret["metadata"]
|
||||
for key in sorted(metadata.keys()):
|
||||
log.log(
|
||||
"{name:>30}: {minimum:>20.06f} {average:>20.06f} "
|
||||
"{maximum:>20.06f}".format(name=key, **metadata[key])
|
||||
)
|
||||
print("\n")
|
||||
log.info("Result:")
|
||||
log.info(funcname)
|
||||
log.info(sample_query)
|
||||
log.success("Latency statistics:")
|
||||
for key, value in ret["latency_stats"].items():
|
||||
if key == "iterations":
|
||||
log.success("{:<10} {:>10}".format(key, value))
|
||||
else:
|
||||
log.success("{:<10} {:>10.06f} seconds".format(key, value))
|
||||
|
||||
log.success("Throughput: {:02f} QPS".format(ret["throughput"]))
|
||||
print("\n\n")
|
||||
|
||||
# Save results.
|
||||
results_key = [
|
||||
@@ -632,8 +695,9 @@ if __name__ == "__main__":
|
||||
|
||||
# If there is need for authorization testing.
|
||||
if benchmark_context.no_authorization:
|
||||
log.info("Running queries with authorization...")
|
||||
vendor_runner.start_benchmark("authorization")
|
||||
log.init("Running queries with authorization...")
|
||||
log.info("Setting USER and PRIVILEGES...")
|
||||
vendor_runner.start_db("authorization")
|
||||
client.execute(
|
||||
queries=[
|
||||
("CREATE USER user IDENTIFIED BY 'test';", {}),
|
||||
@@ -644,13 +708,11 @@ if __name__ == "__main__":
|
||||
)
|
||||
|
||||
client.set_credentials(username="user", password="test")
|
||||
vendor_runner.stop("authorization")
|
||||
vendor_runner.stop_db("authorization")
|
||||
|
||||
for query, funcname in queries[group]:
|
||||
|
||||
log.info(
|
||||
"Running query:",
|
||||
"{}/{}/{}/{}".format(group, query, funcname, WITH_FINE_GRAINED_AUTHORIZATION),
|
||||
log.init(
|
||||
"Running query:" + "{}/{}/{}/{}".format(group, query, funcname, WITH_FINE_GRAINED_AUTHORIZATION)
|
||||
)
|
||||
func = getattr(workload, funcname)
|
||||
|
||||
@@ -664,13 +726,14 @@ if __name__ == "__main__":
|
||||
vendor_runner, client, get_queries(func, 1), config_key, benchmark_context
|
||||
)
|
||||
|
||||
vendor_runner.start_benchmark("authorization")
|
||||
vendor_runner.start_db("authorization")
|
||||
warmup(condition=benchmark_context.warm_up, client=client, queries=get_queries(func, count))
|
||||
|
||||
ret = client.execute(
|
||||
queries=get_queries(func, count),
|
||||
num_workers=benchmark_context.num_workers_for_benchmark,
|
||||
)[0]
|
||||
usage = vendor_runner.stop("authorization")
|
||||
usage = vendor_runner.stop_db("authorization")
|
||||
ret["database"] = usage
|
||||
# Output summary.
|
||||
log.log("Executed {} queries in {} seconds.".format(ret["count"], ret["duration"]))
|
||||
@@ -695,8 +758,8 @@ if __name__ == "__main__":
|
||||
]
|
||||
results.set_value(*results_key, value=ret)
|
||||
|
||||
# Clean up database from any roles and users job
|
||||
vendor_runner.start_benchmark("authorizations")
|
||||
log.info("Deleting USER and PRIVILEGES...")
|
||||
vendor_runner.start_db("authorizations")
|
||||
ret = client.execute(
|
||||
queries=[
|
||||
("REVOKE LABELS * FROM user;", {}),
|
||||
@@ -704,7 +767,7 @@ if __name__ == "__main__":
|
||||
("DROP USER user;", {}),
|
||||
]
|
||||
)
|
||||
vendor_runner.stop("authorization")
|
||||
vendor_runner.stop_db("authorization")
|
||||
|
||||
# Save configuration.
|
||||
if not benchmark_context.no_save_query_counts:
|
||||
@@ -714,3 +777,30 @@ if __name__ == "__main__":
|
||||
if benchmark_context.export_results:
|
||||
with open(benchmark_context.export_results, "w") as f:
|
||||
json.dump(results.get_data(), f)
|
||||
|
||||
# Results summary.
|
||||
log.init("~" * 45)
|
||||
log.info("Benchmark finished.")
|
||||
log.init("~" * 45)
|
||||
log.log("\n")
|
||||
log.summary("Benchmark summary")
|
||||
log.log("-" * 90)
|
||||
log.summary("{:<20} {:>30} {:>30}".format("Query name", "Throughput", "Peak Memory usage"))
|
||||
with open(benchmark_context.export_results, "r") as f:
|
||||
results = json.load(f)
|
||||
for dataset, variants in results.items():
|
||||
if dataset == "__run_configuration__":
|
||||
continue
|
||||
for variant, groups in variants.items():
|
||||
for group, queries in groups.items():
|
||||
if group == "__import__":
|
||||
continue
|
||||
for query, auth in queries.items():
|
||||
for key, value in auth.items():
|
||||
log.log("-" * 90)
|
||||
log.summary(
|
||||
"{:<20} {:>26.2f} QPS {:>27.2f} MB".format(
|
||||
query, value["throughput"], value["database"]["memory"] / 1024.0 / 1024.0
|
||||
)
|
||||
)
|
||||
log.log("-" * 90)
|
||||
|
||||
Reference in New Issue
Block a user