From d81fb90e24a1927ae6b65662f456afbc77c33034 Mon Sep 17 00:00:00 2001 From: Dave Marion Date: Mon, 28 Sep 2026 21:51:48 +0000 Subject: [PATCH] Add Manager Stress Testing framework This commit adds two components that are designed to be used together. The first is the ManagerStress component that creates a configurable number of Accumulo clients to drive a weighted workload of table operations, including create, delete, split, merge, availability changes, compactions, and bulk imports. The second component is the StressCompactor, which does not perform a compaction, but returns success, failure, or cancellation based on a weighted criteria. It starts a configurable number of Compactor instances in a single JVM, each running in their own thread and with their own Thrift RPC server. The bin/manager-stress/configure-mgr-stress-test.sh script is the main entrypoint for the user. It prompts the user for information and outputs commands to be run by the user to start the StressCompactor's and ManagerStress test. Before using this feature the user will need to create a directory in HDFS that will be used by the workers to create RFiles for the bulk import simulation and create a Resource Group for the compactors. Assisted-By: OpenAI GPT-6 Luna --- README.md | 118 +++++ .../configure-mgr-stress-test.sh | 215 ++++++++ bin/manager-stress/managerstress.sh | 40 ++ bin/manager-stress/stresscompactor.sh | 44 ++ pom.xml | 8 + src/build/checkstyle/import-control.xml | 19 + .../testing/manager/stress/ManagerStress.java | 457 +++++++++++++++++ .../manager/stress/ManagerStressMetrics.java | 181 +++++++ .../manager/stress/ManagerStressOptions.java | 134 +++++ .../manager/stress/ManagerStressRows.java | 201 ++++++++ .../manager/stress/ManagerStressWorker.java | 420 +++++++++++++++ .../testing/manager/stress/Operation.java | 77 +++ .../testing/manager/stress/TablePool.java | 294 +++++++++++ .../stress/compactor/CompactionOutcome.java | 46 ++ .../stress/compactor/SimulatedCompactor.java | 485 ++++++++++++++++++ .../stress/compactor/StressCompactor.java | 105 ++++ .../compactor/StressCompactorOptions.java | 112 ++++ .../stress/ManagerStressCleanupTest.java | 62 +++ .../stress/ManagerStressMetricsTest.java | 72 +++ .../stress/ManagerStressOptionsTest.java | 73 +++ .../manager/stress/ManagerStressRowsTest.java | 119 +++++ .../testing/manager/stress/OperationTest.java | 56 ++ .../testing/manager/stress/TablePoolTest.java | 33 ++ .../compactor/CompactionOutcomeTest.java | 37 ++ 24 files changed, 3408 insertions(+) create mode 100755 bin/manager-stress/configure-mgr-stress-test.sh create mode 100755 bin/manager-stress/managerstress.sh create mode 100755 bin/manager-stress/stresscompactor.sh create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStress.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetrics.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptions.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressRows.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressWorker.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/Operation.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/TablePool.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcome.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/compactor/SimulatedCompactor.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactor.java create mode 100644 src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactorOptions.java create mode 100644 src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressCleanupTest.java create mode 100644 src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetricsTest.java create mode 100644 src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptionsTest.java create mode 100644 src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressRowsTest.java create mode 100644 src/test/java/org/apache/accumulo/testing/manager/stress/OperationTest.java create mode 100644 src/test/java/org/apache/accumulo/testing/manager/stress/TablePoolTest.java create mode 100644 src/test/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcomeTest.java diff --git a/README.md b/README.md index 9f0eaf24..d04bfca0 100644 --- a/README.md +++ b/README.md @@ -57,6 +57,9 @@ Tests are run using the following scripts in `bin/`: * `agitator` - Runs agitator * `gcs` - Runs garbage collection simulation * `monitor` - Runs availability monitor probe + * `manager-stress/managerstress.sh` - Stresses Manager and Compaction Coordinator operations + * `manager-stress/stresscompactor.sh` - Starts in-JVM Compactor simulators in a separate JVM + * `manager-stress/configure-mgr-stress-test.sh` - Prompts for settings and prints both commands Run the scripts without arguments to view usage. @@ -176,6 +179,121 @@ verification test. See [gcs.md](docs/gcs.md). +## Manager and Compaction Coordinator stress test + +The `managerstress` command starts local client JVMs that issue weighted create-table, delete-table, +split, merge, tablet-availability, compaction, and bulk-import tasks against an existing Accumulo +instance. It creates the requested number of uniquely named tables before starting workers and +replenishes the table pool after deletes. A selected create task creates a new table generation +before deleting the old generation so the pool returns to its configured size. Weights apply to +selected workload tasks; the follow-up table creations after deletes are unweighted pool maintenance. + +All created table generations start with `TabletAvailability.UNHOSTED` and 10–999 evenly spaced +initial splits, leaving at most 1,000 tablets. Split tasks add a random number of unique split rows +without exceeding that tablet limit. Merge tasks sample a target from 1–100, clamp it to a valid +2–100 tablet range, and use the table's actual split boundaries for the merge. Split rows and the +bulk-import seed key use the same fixed-width hexadecimal row-key scheme. + +Availability tasks sample a random number of unique tablets from one through the table's full +tablet count. Each selected tablet is assigned a random availability state different from its +current state. Nonadjacent selections remain separate; adjacent selections with the same target +state are grouped into one row-range request. + +Bulk import requires a shared HDFS directory accessible to all local workers. Each worker creates +one one-entry RFile there at startup and stages a copy of it for every bulk import. Compaction tasks +submit non-blocking major compactions. The coordinator removes the run-specific HDFS staging +directory during shutdown, including any files left by workers that did not exit cleanly. + +`managerstress` temporarily configures the `mgrstress` compaction service in the system +configuration. It uses `org.apache.accumulo.core.spi.compaction.RatioBasedCompactionPlanner`, with +`compaction.service.mgrstress.planner.opts.groups` set to the group specified by +the required `--compactor-resource-group` option. The stress namespace is created with +`table.compaction.dispatcher.opts.service=mgrstress`. Previous ZooKeeper system-property overrides +are restored when the test finishes; pre-existing values do not prevent startup. + +The required `--namespace` must name a namespace that does not already exist. `managerstress` creates +it before creating tables, places all stress tables in it, and removes its tables and namespace when +the run finishes. The configured Accumulo user needs permission to create and drop namespaces and +tables, alter tables and namespaces, modify system configuration, compact, and bulk import. It also +needs access to the configured HDFS directory. + +For example, the command below runs four client processes for 30 minutes, maintaining eight tables. +Operation weights are relative; setting a weight to zero disables that operation. + +```bash +./bin/manager-stress/managerstress.sh \ + --namespace mgrstress_example \ + --compactor-resource-group default \ + --duration 30m \ + --clients 4 \ + --tables 8 \ + --hdfs-dir hdfs://namenode:8020/tmp/accumulo-managerstress \ + --create-weight 1 \ + --delete-weight 1 \ + --split-weight 2 \ + --merge-weight 2 \ + --availability-weight 1 \ + --compact-weight 4 \ + --bulk-import-weight 1 +``` + +Use `--table-prefix` to choose a prefix for the generated tables; if omitted, a unique prefix is +created for the run. Choose a namespace name that is unused on the Accumulo instance; the test fails +if it already exists. The test removes all tables in its namespace and deletes the namespace after +workers stop. Per-worker logs include operation attempts, completions, skips, races, errors, and +average latency. + +When workers stop, `managerstress` logs totals aggregated from each worker's atomic metrics snapshot. +For each operation, `submitted` counts workload operations that reached their Manager API call, +`completed` counts successful synchronous operations, and `asyncAccepted` counts compaction and +tablet-availability requests accepted by Accumulo. Availability and compaction requests do not wait +for tablets to reach their target state. `failed` counts operation errors; `outcomeUnknown` is the +subset with a transport failure or interruption after submission. Snapshot files are refreshed +periodically in the run's temporary control directory and removed after the summary is logged. + +### Compactor simulator + +`stresscompactor` is designed to be used with the `managerstress` test framework. Run it alongside +`managerstress` to provide Compactor simulators for external compaction requests generated by the +stress tables. Each invocation starts the configured number of independently registered Compactor +simulators as threads in that JVM. Start additional JVMs by invoking the command again (or on other +hosts); set `JAVA_OPTS` separately for each process to tune its heap and garbage collector. + +Use the same group for `managerstress --compactor-resource-group` and +`stresscompactor --resource-group`. Run `./bin/manager-stress/configure-mgr-stress-test.sh` to +interactively enter settings for both processes and print the two commands. Press Enter to accept +the displayed defaults. The client config prompt defaults to `../../conf/accumulo-client.properties` +relative to the helper script; credentials should be stored in that file. The helper prints commands +to stdout without starting the processes. + +The Manager must be able to reach the advertised host and callback ports. Configure the resource +group to match the Compactor resource group used by the external compaction service for the stress +tables. The configured Accumulo user must be permitted to perform system actions. When using SASL, +the advertised host should be the host's canonical name. + +```bash +JAVA_OPTS="-Xms2g -Xmx2g -XX:+UseG1GC" ./bin/manager-stress/stresscompactor.sh \ + --compactors-per-jvm 16 \ + --duration 30m \ + --resource-group default \ + --host compactor-client01.example.net \ + --port-range 9600-9699 \ + --success-weight 1 \ + --failure-weight 1 \ + --cancellation-weight 1 +``` + +Each simulator applies `--resource-group` to its own `compactor.group` configuration property, +registers its Compactor address, and serves the Manager's status, wake, and cancel RPCs. On receiving +a job it immediately reports a result according to the configured relative weights. A successful +result reports zero output entries, which commits as a zero-output compaction and removes the +input-file references from the stress table. Simulators unregister and close their connections when +the JVM's duration expires or it is stopped. + +Each simulator logs its job outcomes on shutdown, followed by a per-JVM aggregate. These counts are +external compaction jobs, which may be multiple for one table-level compaction request; they are +reported separately from `managerstress`'s accepted asynchronous request count. + ## Agitator The agitator will periodically kill the Accumulo manager, tablet server, and Hadoop data node diff --git a/bin/manager-stress/configure-mgr-stress-test.sh b/bin/manager-stress/configure-mgr-stress-test.sh new file mode 100755 index 00000000..176df54b --- /dev/null +++ b/bin/manager-stress/configure-mgr-stress-test.sh @@ -0,0 +1,215 @@ +#! /usr/bin/env bash +# +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +# + +set -euo pipefail + +script_dir=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +default_client_config="${script_dir}/../../conf/accumulo-client.properties" +default_namespace="mgrstress_$(date +%Y%m%d_%H%M%S)_${RANDOM}" +default_host=$(hostname -f 2>/dev/null || hostname) + +prompt_default() { + local label=$1 + local default=$2 + local answer + printf '%s [%s]: ' "$label" "$default" >&2 + IFS= read -r answer + REPLY=${answer:-$default} +} + +prompt_optional() { + local label=$1 + local default_description=$2 + printf '%s [%s]: ' "$label" "$default_description" >&2 + IFS= read -r REPLY +} + +shell_quote() { + local value=$1 + value=${value//\'/\'\\\'\'} + printf "'%s'" "$value" +} + +printf '%s\n' \ + 'Enter the ManagerStress settings. Press Enter to accept the displayed defaults.' \ + 'Client passwords should be stored in the selected client properties file.' >&2 + +prompt_default 'ManagerStress duration' '5m' +manager_duration=$REPLY +prompt_default 'ManagerStress client JVMs' '2' +manager_clients=$REPLY +prompt_default 'ManagerStress table count' '4' +manager_tables=$REPLY +prompt_default 'Namespace (generated names must not already exist)' "$default_namespace" +manager_namespace=$REPLY +prompt_default 'Compactor resource group (shared with StressCompactor)' 'default' +compactor_resource_group=$REPLY +prompt_optional 'Table prefix (blank generates one)' 'generated automatically' +manager_table_prefix=$REPLY +prompt_default 'Shared HDFS directory for bulk imports' 'hdfs:///tmp/accumulo-managerstress' +manager_hdfs_dir=$REPLY +prompt_default 'Create-table weight' '1' +manager_create_weight=$REPLY +prompt_default 'Delete-table weight' '1' +manager_delete_weight=$REPLY +prompt_default 'Split weight' '1' +manager_split_weight=$REPLY +prompt_default 'Merge weight' '1' +manager_merge_weight=$REPLY +prompt_default 'Tablet-availability weight' '1' +manager_availability_weight=$REPLY +prompt_default 'Compaction weight' '1' +manager_compact_weight=$REPLY +prompt_default 'Bulk-import weight' '1' +manager_bulk_import_weight=$REPLY +prompt_optional 'ManagerStress random seed (blank uses a random seed)' 'random per run' +manager_seed=$REPLY +prompt_optional 'ManagerStress JAVA_OPTS (blank uses JVM defaults)' 'empty' +manager_java_opts=$REPLY + +printf '\n%s\n' 'Enter the StressCompactor settings.' >&2 +prompt_default 'StressCompactor JVM duration' '5m' +compactor_duration=$REPLY +prompt_default 'Compactors per JVM' '1' +compactors_per_jvm=$REPLY +prompt_default 'Compactor callback host' "$default_host" +compactor_host=$REPLY +prompt_default 'Compactor callback port range' '0' +compactor_port_range=$REPLY +prompt_default 'Compaction success weight' '1' +compactor_success_weight=$REPLY +prompt_default 'Compaction failure weight' '1' +compactor_failure_weight=$REPLY +prompt_default 'Compaction cancellation weight' '1' +compactor_cancellation_weight=$REPLY +prompt_optional 'StressCompactor random seed (blank uses a random seed)' 'random per run' +compactor_seed=$REPLY +prompt_optional 'StressCompactor JAVA_OPTS (blank uses JVM defaults)' 'empty' +compactor_java_opts=$REPLY + +printf '\n%s\n' 'Enter shared Accumulo client settings.' >&2 +prompt_default 'Accumulo client config file' "$default_client_config" +client_config=$REPLY +prompt_optional 'Accumulo user override (blank uses client config)' 'from client config' +client_user=$REPLY +prompt_optional 'Authorizations (blank uses none)' 'empty' +client_auths=$REPLY +prompt_default 'Enable distributed tracing? (yes/no)' 'no' +trace_choice=${REPLY,,} +case "$trace_choice" in + yes | y) enable_trace=true ;; + no | n) enable_trace=false ;; + *) + printf 'Please answer yes or no.\n' >&2 + exit 2 + ;; +esac + +client_args=(--config-file "$client_config") +if [[ -n "$client_user" ]]; then + client_args+=(--user "$client_user") +fi +if [[ -n "$client_auths" ]]; then + client_args+=(--auths "$client_auths") +fi +if [[ "$enable_trace" == true ]]; then + client_args+=(--trace) +fi +while true; do + prompt_optional 'Additional client property override (key=value; blank to finish)' 'none' + if [[ -z "$REPLY" ]]; then + break + fi + client_args+=(-o "$REPLY") +done + +manager_args=( + --namespace "$manager_namespace" + --compactor-resource-group "$compactor_resource_group" + --duration "$manager_duration" + --clients "$manager_clients" + --tables "$manager_tables" + --hdfs-dir "$manager_hdfs_dir" + --create-weight "$manager_create_weight" + --delete-weight "$manager_delete_weight" + --split-weight "$manager_split_weight" + --merge-weight "$manager_merge_weight" + --availability-weight "$manager_availability_weight" + --compact-weight "$manager_compact_weight" + --bulk-import-weight "$manager_bulk_import_weight" +) +if [[ -n "$manager_table_prefix" ]]; then + manager_args+=(--table-prefix "$manager_table_prefix") +fi +if [[ -n "$manager_seed" ]]; then + manager_args+=(--seed "$manager_seed") +fi +manager_args+=("${client_args[@]}") + +compactor_args=( + --resource-group "$compactor_resource_group" + --duration "$compactor_duration" + --compactors-per-jvm "$compactors_per_jvm" + --host "$compactor_host" + --port-range "$compactor_port_range" + --success-weight "$compactor_success_weight" + --failure-weight "$compactor_failure_weight" + --cancellation-weight "$compactor_cancellation_weight" +) +if [[ -n "$compactor_seed" ]]; then + compactor_args+=(--seed "$compactor_seed") +fi +compactor_args+=("${client_args[@]}") + +print_command() { + local label=$1 + local java_opts=$2 + local script=$3 + local argument + local index=0 + shift 3 + local -a arguments=("$@") + + printf '# %s\n' "$label" + if [[ -n "$java_opts" ]]; then + printf 'JAVA_OPTS=%s ' "$(shell_quote "$java_opts")" + fi + printf '%q' "$script" + while ((index < ${#arguments[@]})); do + argument=${arguments[index]} + printf ' \\\n %q' "$argument" + if [[ "$argument" == --trace ]]; then + index=$((index + 1)) + continue + fi + if ((index + 1 >= ${#arguments[@]})); then + printf 'Option %s is missing its value.\n' "$argument" >&2 + return 2 + fi + printf ' %q' "${arguments[index + 1]}" + index=$((index + 2)) + done + printf '\n\n' +} + +print_command 'ManagerStress' "$manager_java_opts" "${script_dir}/managerstress.sh" \ + "${manager_args[@]}" +print_command 'StressCompactor' "$compactor_java_opts" "${script_dir}/stresscompactor.sh" \ + "${compactor_args[@]}" diff --git a/bin/manager-stress/managerstress.sh b/bin/manager-stress/managerstress.sh new file mode 100755 index 00000000..4218c5c0 --- /dev/null +++ b/bin/manager-stress/managerstress.sh @@ -0,0 +1,40 @@ +#! /usr/bin/env bash +# +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +# + +bin_dir=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +source "${bin_dir}/../build" + +export CLASSPATH="$TEST_JAR_PATH:$HADOOP_API_JAR:$HADOOP_RUNTIME_JAR:$CLASSPATH" + +client_config_set=false +for arg in "$@"; do + case "$arg" in + -c | --config-file | --config-file=*) + client_config_set=true + break + ;; + esac +done +if [[ "$client_config_set" == false ]]; then + set -- "$@" -c "$ACCUMULO_CLIENT_PROPS" +fi + +java $JAVA_OPTS -Dlog4j.configurationFile="file:$TEST_LOG4J" \ + org.apache.accumulo.testing.manager.stress.ManagerStress "$@" diff --git a/bin/manager-stress/stresscompactor.sh b/bin/manager-stress/stresscompactor.sh new file mode 100755 index 00000000..705a1f62 --- /dev/null +++ b/bin/manager-stress/stresscompactor.sh @@ -0,0 +1,44 @@ +#! /usr/bin/env bash +# +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +# + +bin_dir=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd) +source "${bin_dir}/../build" + +export CLASSPATH="$ACCUMULO_HOME/conf:$ACCUMULO_HOME/lib/*:$TEST_JAR_PATH:$HADOOP_API_JAR:$HADOOP_RUNTIME_JAR:$CLASSPATH" + +if [[ "$#" -eq 0 ]]; then + set -- --help +fi + +client_config_set=false +for arg in "$@"; do + case "$arg" in + -c | --config-file | --config-file=*) + client_config_set=true + break + ;; + esac +done +if [[ "$client_config_set" == false ]]; then + set -- "$@" -c "$ACCUMULO_CLIENT_PROPS" +fi + +java $JAVA_OPTS -Dlog4j.configurationFile="file:$TEST_LOG4J" \ + org.apache.accumulo.testing.manager.stress.compactor.StressCompactor "$@" diff --git a/pom.xml b/pom.xml index c27351cd..ece487b5 100644 --- a/pom.xml +++ b/pom.xml @@ -105,6 +105,11 @@ org.slf4j slf4j-api + + org.apache.accumulo + accumulo-server-base + provided + org.apache.accumulo accumulo-minicluster @@ -156,6 +161,8 @@ SCRIPT_STYLE SCRIPT_STYLE SCRIPT_STYLE + SCRIPT_STYLE + SCRIPT_STYLE SCRIPT_STYLE SCRIPT_STYLE SCRIPT_STYLE @@ -165,6 +172,7 @@ SCRIPT_STYLE SCRIPT_STYLE SLASHSTAR_STYLE + SCRIPT_STYLE SCRIPT_STYLE SCRIPT_STYLE diff --git a/src/build/checkstyle/import-control.xml b/src/build/checkstyle/import-control.xml index 43008f5c..b602a2a3 100644 --- a/src/build/checkstyle/import-control.xml +++ b/src/build/checkstyle/import-control.xml @@ -30,6 +30,7 @@ + @@ -41,6 +42,24 @@ + + + + + + + + + + + + + + + + + + diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStress.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStress.java new file mode 100644 index 00000000..57be1341 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStress.java @@ -0,0 +1,457 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static java.nio.file.StandardCopyOption.ATOMIC_MOVE; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.NoSuchFileException; +import java.nio.file.Path; +import java.time.Duration; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.UUID; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.apache.accumulo.core.client.Accumulo; +import org.apache.accumulo.core.client.AccumuloClient; +import org.apache.accumulo.core.conf.Property; +import org.apache.accumulo.core.data.ResourceGroupId; +import org.apache.accumulo.core.spi.compaction.RatioBasedCompactionPlanner; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** Starts local client processes which issue Manager and Compaction Coordinator work. */ +public class ManagerStress { + + private static final Logger log = LoggerFactory.getLogger(ManagerStress.class); + private static final Duration WORKER_STARTUP_TIMEOUT = Duration.ofMinutes(2); + private static final Duration WORKER_SHUTDOWN_TIMEOUT = Duration.ofSeconds(30); + private static final String COMPACTION_SERVICE_NAME = "mgrstress"; + + private record PreviousSystemProperty(boolean wasSet, String value) { + } + + public static void main(String[] args) throws Exception { + ManagerStressOptions options = new ManagerStressOptions(); + options.parseArgs(ManagerStress.class.getName(), args); + if (options.workerId >= 0) { + options.validate(false); + new ManagerStressWorker(options).run(); + } else { + options.validate(true); + new ManagerStress(options, args).run(); + } + } + + private final ManagerStressOptions options; + private final String[] originalArgs; + private final List workers = new CopyOnWriteArrayList<>(); + private final AtomicBoolean stopping = new AtomicBoolean(); + private final AtomicBoolean metricsReported = new AtomicBoolean(); + private volatile boolean namespaceCreated; + private volatile boolean compactionServicePropertiesConfigured; + private volatile Map previousCompactionServiceProperties = + Map.of(); + private volatile Map desiredCompactionServiceProperties = Map.of(); + private Path stateDir; + private String runId; + + private ManagerStress(ManagerStressOptions options, String[] originalArgs) { + this.options = options; + this.originalArgs = originalArgs.clone(); + } + + private void run() throws Exception { + String tablePrefix = options.tablePrefix; + if (tablePrefix == null || tablePrefix.isBlank()) { + tablePrefix = "mgrstress_" + UUID.randomUUID().toString().replace("-", "").substring(0, 12); + options.tablePrefix = tablePrefix; + } + runId = UUID.randomUUID().toString().replace("-", ""); + stateDir = Files.createTempDirectory("accumulo-manager-stress-"); + + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + stopWorkers(); + try { + reportRunMetrics(); + } catch (Exception e) { + log.error("Could not aggregate manager-stress metrics", e); + } + try { + cleanupNamespace(); + } catch (Exception e) { + log.error("Could not clean up manager-stress namespace {}", options.namespace, e); + } finally { + try { + cleanupBulkImportDirectory(); + } finally { + try { + restoreCompactionServiceProperties(); + } catch (Exception e) { + log.error("Could not restore compaction service system properties", e); + } finally { + deleteControlDirectory(); + } + } + } + }, "manager-stress-shutdown")); + + try { + try (AccumuloClient client = Accumulo.newClient().from(options.getClientProps()).build()) { + createNamespace(client); + configureCompactionService(client); + new TablePool(stateDir, options.namespace, tablePrefix, options.tables, options.seed) + .initialize(client); + } + + log.info("Created {} stress tables with prefix {}", options.tables, tablePrefix); + startWorkers(); + waitForWorkersReady(); + + long deadline = Math.addExact(System.currentTimeMillis(), options.duration); + Path startFile = stateDir.resolve("start"); + Path tempStartFile = stateDir.resolve("start.tmp"); + Files.writeString(tempStartFile, Long.toString(deadline), UTF_8); + Files.move(tempStartFile, startFile, ATOMIC_MOVE); + log.info("Started {} client processes for {} ms", options.clients, options.duration); + awaitWorkers(deadline); + } finally { + stopWorkers(); + try { + reportRunMetrics(); + } finally { + try { + cleanupNamespace(); + } finally { + try { + cleanupBulkImportDirectory(); + } finally { + try { + restoreCompactionServiceProperties(); + } finally { + deleteControlDirectory(); + } + } + } + } + } + } + + private synchronized void reportRunMetrics() { + if (!metricsReported.compareAndSet(false, true)) { + return; + } + List snapshots = new ArrayList<>(); + for (int workerId = 0; workerId < options.clients; workerId++) { + Path snapshotPath = ManagerStressMetrics.snapshotPath(stateDir, workerId); + if (Files.exists(snapshotPath)) { + try { + snapshots.add(ManagerStressMetrics.readSnapshot(snapshotPath)); + } catch (Exception e) { + log.warn("Could not read metrics snapshot for worker {}", workerId, e); + } + } else { + log.warn("No final metrics snapshot was available for worker {}", workerId); + } + } + + ManagerStressMetrics.RunSnapshot totals = ManagerStressMetrics.aggregate(snapshots); + log.info("ManagerStress totals from {}/{} worker snapshots:", totals.workers(), + options.clients); + for (Operation operation : Operation.values()) { + ManagerStressMetrics.OperationSnapshot counts = totals.operations().get(operation); + double averageMillis = + counts.attempts() == 0 ? 0 : (double) counts.totalNanos() / counts.attempts() / 1_000_000; + log.info( + "ManagerStress {}: attempts={}, submitted={}, completed={}, asyncAccepted={}, " + + "skipped={}, races={}, failed={}, outcomeUnknown={}, avgMs={}", + operation, counts.attempts(), counts.submitted(), counts.completed(), + counts.asyncAccepted(), counts.skipped(), counts.races(), counts.failed(), + counts.outcomeUnknown(), String.format("%.3f", averageMillis)); + } + log.info("ManagerStress replacement-table maintenance creates={}", totals.maintenanceCreates()); + } + + private void createNamespace(AccumuloClient client) throws Exception { + if (client.namespaceOperations().exists(options.namespace)) { + throw new IllegalStateException("Refusing to use existing namespace " + options.namespace); + } + client.namespaceOperations().create(options.namespace); + namespaceCreated = true; + client.namespaceOperations().setProperty(options.namespace, + Property.TABLE_COMPACTION_DISPATCHER_OPTS.getKey() + "service", COMPACTION_SERVICE_NAME); + log.info("Created namespace {} for manager-stress run", options.namespace); + } + + static Map compactionServiceProperties(String compactorResourceGroup) { + String servicePrefix = Property.COMPACTION_SERVICE_PREFIX.getKey() + COMPACTION_SERVICE_NAME; + String group = ResourceGroupId.of(compactorResourceGroup).canonical(); + return Map.of(servicePrefix + ".planner", RatioBasedCompactionPlanner.class.getName(), + servicePrefix + ".planner.opts.groups", "[{\"group\":\"" + group + "\"}]"); + } + + private synchronized void configureCompactionService(AccumuloClient client) throws Exception { + Map desired = compactionServiceProperties(options.compactorResourceGroup); + Map previous = new HashMap<>(); + desiredCompactionServiceProperties = desired; + client.instanceOperations().modifyProperties(systemProperties -> { + previous.clear(); + desired.forEach((key, value) -> { + previous.put(key, new PreviousSystemProperty(systemProperties.containsKey(key), + systemProperties.get(key))); + systemProperties.put(key, value); + }); + previousCompactionServiceProperties = Map.copyOf(previous); + compactionServicePropertiesConfigured = true; + }); + log.info("Configured compaction service {} for resource group {}", COMPACTION_SERVICE_NAME, + options.compactorResourceGroup); + } + + private synchronized void restoreCompactionServiceProperties() throws Exception { + if (!compactionServicePropertiesConfigured) { + return; + } + try (AccumuloClient client = Accumulo.newClient().from(options.getClientProps()).build()) { + client.instanceOperations().modifyProperties(systemProperties -> { + previousCompactionServiceProperties.forEach((key, previous) -> { + if (!Objects.equals(systemProperties.get(key), + desiredCompactionServiceProperties.get(key))) { + log.warn( + "System property {} changed during the ManagerStress run; preserving its value", + key); + return; + } + if (previous.wasSet()) { + systemProperties.put(key, previous.value()); + } else { + systemProperties.remove(key); + } + }); + }); + compactionServicePropertiesConfigured = false; + log.info("Restored compaction service system properties after ManagerStress run"); + } + } + + private synchronized void cleanupNamespace() throws Exception { + if (!namespaceCreated) { + return; + } + try (AccumuloClient client = Accumulo.newClient().from(options.getClientProps()).build()) { + var tableOps = client.tableOperations(); + String namespacePrefix = options.namespace + "."; + List tables = tableOps.list().stream() + .filter(tableName -> tableName.startsWith(namespacePrefix)).toList(); + Exception cleanupError = null; + for (String table : tables) { + try { + tableOps.delete(table); + } catch (Exception e) { + if (cleanupError == null) { + cleanupError = new IllegalStateException( + "Could not remove all tables from manager-stress namespace " + options.namespace); + } + cleanupError.addSuppressed(e); + } + } + if (cleanupError != null) { + throw cleanupError; + } + client.namespaceOperations().delete(options.namespace); + namespaceCreated = false; + log.info("Deleted namespace {} after manager-stress run", options.namespace); + } + } + + private void cleanupBulkImportDirectory() { + if (options.hdfsDir == null || options.hdfsDir.isBlank() || runId == null) { + return; + } + try { + deleteBulkImportRunDirectory(options.hdfsDir, runId); + log.info("Deleted HDFS bulk-import directory for run {}", runId); + } catch (IOException e) { + log.warn("Could not remove HDFS bulk-import directory for run {}", runId, e); + } + } + + static void deleteBulkImportRunDirectory(String configuredDirectory, String runId) + throws IOException { + var runDirectory = + new org.apache.hadoop.fs.Path(new org.apache.hadoop.fs.Path(configuredDirectory), runId); + var configuration = new org.apache.hadoop.conf.Configuration(); + try (var fileSystem = + org.apache.hadoop.fs.FileSystem.newInstance(runDirectory.toUri(), configuration)) { + boolean deleted = fileSystem.delete(runDirectory, true); + if (!deleted && fileSystem.exists(runDirectory)) { + throw new IOException("Could not delete HDFS bulk-import directory " + runDirectory); + } + } + } + + private void startWorkers() throws IOException { + for (int i = 0; i < options.clients; i++) { + List command = new ArrayList<>(); + command.add(Path.of(System.getProperty("java.home"), "bin", "java").toString()); + String logConfig = System.getProperty("log4j.configurationFile"); + if (logConfig != null) { + command.add("-Dlog4j.configurationFile=" + logConfig); + } + command.add("-cp"); + command.add(System.getProperty("java.class.path")); + command.add(ManagerStress.class.getName()); + command.addAll(workerArguments()); + command.add("--table-prefix"); + command.add(options.tablePrefix); + command.add("--internal-worker-id"); + command.add(Integer.toString(i)); + command.add("--internal-state-dir"); + command.add(stateDir.toString()); + command.add("--internal-run-id"); + command.add(runId); + + Process process = new ProcessBuilder(command).inheritIO().start(); + workers.add(process); + log.info("Started worker {} with pid {}", i, process.pid()); + } + } + + private void waitForWorkersReady() throws Exception { + long deadline = System.nanoTime() + WORKER_STARTUP_TIMEOUT.toNanos(); + while (System.nanoTime() < deadline) { + int ready = 0; + for (int i = 0; i < workers.size(); i++) { + if (Files.exists(stateDir.resolve("ready-" + i))) { + ready++; + } else if (!workers.get(i).isAlive()) { + throw new IllegalStateException("Worker " + i + " exited before startup completed"); + } + } + if (ready == workers.size()) { + return; + } + Thread.sleep(100); + } + throw new IllegalStateException("Timed out waiting for client workers to initialize"); + } + + private void awaitWorkers(long deadline) throws InterruptedException { + long untilDeadline = deadline - System.currentTimeMillis(); + if (untilDeadline > 0) { + Thread.sleep(untilDeadline); + } + for (int i = 0; i < workers.size(); i++) { + Process process = workers.get(i); + if (!process.waitFor(WORKER_SHUTDOWN_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS)) { + log.warn("Worker {} did not stop before the shutdown timeout", i); + process.destroy(); + try { + if (!process.waitFor(5, TimeUnit.SECONDS)) { + process.destroyForcibly(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + process.destroyForcibly(); + throw e; + } + } else if (process.exitValue() != 0) { + log.warn("Worker {} exited with status {}", i, process.exitValue()); + } + } + } + + private List workerArguments() { + List result = new ArrayList<>(); + for (int i = 0; i < originalArgs.length; i++) { + String arg = originalArgs[i]; + if (arg.equals("--table-prefix") || arg.equals("--internal-worker-id") + || arg.equals("--internal-state-dir") || arg.equals("--internal-run-id") + || arg.startsWith("--table-prefix=") || arg.startsWith("--internal-worker-id=") + || arg.startsWith("--internal-state-dir=") || arg.startsWith("--internal-run-id=")) { + if (arg.contains("=")) { + continue; + } + i++; + } else { + result.add(arg); + } + } + return result; + } + + private synchronized void stopWorkers() { + if (!stopping.compareAndSet(false, true)) { + return; + } + for (Process process : workers) { + if (process.isAlive()) { + process.destroy(); + } + } + for (Process process : workers) { + if (process.isAlive()) { + try { + if (!process.waitFor(WORKER_SHUTDOWN_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS)) { + process.destroyForcibly(); + process.waitFor(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + process.destroyForcibly(); + } + } + } + } + + private synchronized void deleteControlDirectory() { + Path directory = stateDir; + stateDir = null; + if (directory == null) { + return; + } + try { + deleteRecursivelyIfExists(directory); + } catch (IOException e) { + log.warn("Could not remove temporary manager-stress directory {}", directory, e); + } + } + + static void deleteRecursivelyIfExists(Path directory) throws IOException { + try (var paths = Files.walk(directory)) { + for (Path path : paths.sorted((a, b) -> b.compareTo(a)).toList()) { + try { + Files.deleteIfExists(path); + } catch (NoSuchFileException ignored) { + // The entry may have been removed by another cleanup invocation. + } + } + } catch (NoSuchFileException ignored) { + // The root directory may already have been removed. + } + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetrics.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetrics.java new file mode 100644 index 00000000..a4aee5a0 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetrics.java @@ -0,0 +1,181 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static java.nio.file.StandardCopyOption.ATOMIC_MOVE; +import static java.nio.file.StandardCopyOption.REPLACE_EXISTING; + +import java.io.IOException; +import java.io.OutputStream; +import java.nio.file.AtomicMoveNotSupportedException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Collection; +import java.util.EnumMap; +import java.util.Map; +import java.util.Properties; + +/** Per-worker workload metrics and their atomic on-disk snapshots. */ +final class ManagerStressMetrics { + + static final class Counts { + long attempts; + long submitted; + long completed; + long asyncAccepted; + long skipped; + long races; + long failed; + long outcomeUnknown; + long totalNanos; + + OperationSnapshot snapshot() { + return new OperationSnapshot(attempts, submitted, completed, asyncAccepted, skipped, races, + failed, outcomeUnknown, totalNanos); + } + } + + record OperationSnapshot(long attempts, long submitted, long completed, long asyncAccepted, + long skipped, long races, long failed, long outcomeUnknown, long totalNanos) { + } + + record Snapshot(int workerId, long maintenanceCreates, + Map operations) { + } + + record RunSnapshot(int workers, long maintenanceCreates, + Map operations) { + } + + private final EnumMap counts = new EnumMap<>(Operation.class); + private long maintenanceCreates; + + ManagerStressMetrics() { + for (Operation operation : Operation.values()) { + counts.put(operation, new Counts()); + } + } + + Counts counts(Operation operation) { + return counts.get(operation); + } + + void addMaintenanceCreates(long creates) { + maintenanceCreates += creates; + } + + long maintenanceCreates() { + return maintenanceCreates; + } + + void writeSnapshot(Path stateDir, int workerId) throws IOException { + Properties properties = new Properties(); + properties.setProperty("workerId", Integer.toString(workerId)); + properties.setProperty("maintenanceCreates", Long.toString(maintenanceCreates)); + for (Operation operation : Operation.values()) { + Counts operationCounts = counts.get(operation); + String prefix = operation.name() + "."; + properties.setProperty(prefix + "attempts", Long.toString(operationCounts.attempts)); + properties.setProperty(prefix + "submitted", Long.toString(operationCounts.submitted)); + properties.setProperty(prefix + "completed", Long.toString(operationCounts.completed)); + properties.setProperty(prefix + "asyncAccepted", + Long.toString(operationCounts.asyncAccepted)); + properties.setProperty(prefix + "skipped", Long.toString(operationCounts.skipped)); + properties.setProperty(prefix + "races", Long.toString(operationCounts.races)); + properties.setProperty(prefix + "failed", Long.toString(operationCounts.failed)); + properties.setProperty(prefix + "outcomeUnknown", + Long.toString(operationCounts.outcomeUnknown)); + properties.setProperty(prefix + "totalNanos", Long.toString(operationCounts.totalNanos)); + } + + Path target = snapshotPath(stateDir, workerId); + Path temporary = Files.createTempFile(stateDir, "worker-metrics-", ".tmp"); + try { + try (OutputStream output = Files.newOutputStream(temporary)) { + properties.store(output, "ManagerStress worker metrics"); + } + try { + Files.move(temporary, target, ATOMIC_MOVE, REPLACE_EXISTING); + } catch (AtomicMoveNotSupportedException e) { + Files.move(temporary, target, REPLACE_EXISTING); + } + } finally { + Files.deleteIfExists(temporary); + } + } + + static Path snapshotPath(Path stateDir, int workerId) { + return stateDir.resolve("worker-metrics-" + workerId + ".properties"); + } + + static Snapshot readSnapshot(Path snapshotPath) throws IOException { + Properties properties = new Properties(); + try (var input = Files.newInputStream(snapshotPath)) { + properties.load(input); + } + int workerId = Integer.parseInt(properties.getProperty("workerId")); + long maintenanceCreates = Long.parseLong(properties.getProperty("maintenanceCreates")); + EnumMap operations = new EnumMap<>(Operation.class); + for (Operation operation : Operation.values()) { + String prefix = operation.name() + "."; + operations.put(operation, new OperationSnapshot(readLong(properties, prefix + "attempts"), + readLong(properties, prefix + "submitted"), readLong(properties, prefix + "completed"), + readLong(properties, prefix + "asyncAccepted"), readLong(properties, prefix + "skipped"), + readLong(properties, prefix + "races"), readLong(properties, prefix + "failed"), + readLong(properties, prefix + "outcomeUnknown"), + readLong(properties, prefix + "totalNanos"))); + } + return new Snapshot(workerId, maintenanceCreates, Map.copyOf(operations)); + } + + static RunSnapshot aggregate(Collection snapshots) { + EnumMap totals = new EnumMap<>(Operation.class); + for (Operation operation : Operation.values()) { + totals.put(operation, new Counts()); + } + long maintenanceCreates = 0; + int workers = 0; + for (Snapshot snapshot : snapshots) { + workers++; + maintenanceCreates += snapshot.maintenanceCreates(); + for (Operation operation : Operation.values()) { + OperationSnapshot current = snapshot.operations().get(operation); + Counts total = totals.get(operation); + total.attempts += current.attempts(); + total.submitted += current.submitted(); + total.completed += current.completed(); + total.asyncAccepted += current.asyncAccepted(); + total.skipped += current.skipped(); + total.races += current.races(); + total.failed += current.failed(); + total.outcomeUnknown += current.outcomeUnknown(); + total.totalNanos += current.totalNanos(); + } + } + EnumMap operations = new EnumMap<>(Operation.class); + for (Operation operation : Operation.values()) { + operations.put(operation, totals.get(operation).snapshot()); + } + return new RunSnapshot(workers, maintenanceCreates, Map.copyOf(operations)); + } + + private static long readLong(Properties properties, String key) { + return Long.parseLong(properties.getProperty(key, "0")); + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptions.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptions.java new file mode 100644 index 00000000..bd90ba94 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptions.java @@ -0,0 +1,134 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import org.apache.accumulo.core.data.ResourceGroupId; +import org.apache.accumulo.testing.cli.ClientOpts; + +import com.beust.jcommander.Parameter; + +class ManagerStressOptions extends ClientOpts { + + @Parameter(names = "--duration", converter = TimeConverter.class, + description = "duration of the stress run (for example 10m or 1h)") + long duration = 5 * 60 * 1000L; + + @Parameter(names = "--clients", description = "number of local worker JVMs") + int clients = 2; + + @Parameter(names = "--tables", description = "target number of managed tables") + int tables = 4; + + @Parameter(names = "--namespace", required = true, + description = "new namespace used for all tables created by this run") + String namespace; + + @Parameter(names = "--compactor-resource-group", required = true, + description = "resource group used by the external Compactor simulators") + String compactorResourceGroup; + + @Parameter(names = "--table-prefix", description = "prefix for tables created by this run") + String tablePrefix; + + @Parameter(names = "--hdfs-dir", description = "shared HDFS directory for bulk import files") + String hdfsDir; + + @Parameter(names = "--create-weight", description = "relative weight for create-table tasks") + double createWeight = 1; + + @Parameter(names = "--delete-weight", description = "relative weight for delete-table tasks") + double deleteWeight = 1; + + @Parameter(names = "--split-weight", description = "relative weight for split tasks") + double splitWeight = 1; + + @Parameter(names = "--merge-weight", description = "relative weight for merge tasks") + double mergeWeight = 1; + + @Parameter(names = "--availability-weight", + description = "relative weight for tablet availability changes") + double availabilityWeight = 1; + + @Parameter(names = "--compact-weight", description = "relative weight for compaction tasks") + double compactWeight = 1; + + @Parameter(names = "--bulk-import-weight", description = "relative weight for bulk import tasks") + double bulkImportWeight = 1; + + @Parameter(names = "--seed", description = "random seed (defaults to a per-run random seed)") + long seed = System.nanoTime(); + + @Parameter(names = "--internal-worker-id", hidden = true) + int workerId = -1; + + @Parameter(names = "--internal-state-dir", hidden = true) + String stateDir; + + @Parameter(names = "--internal-run-id", hidden = true) + String runId; + + void validate(boolean coordinator) { + if (duration <= 0) { + throw new IllegalArgumentException("--duration must be positive"); + } + if (clients <= 0) { + throw new IllegalArgumentException("--clients must be positive"); + } + if (tables <= 0) { + throw new IllegalArgumentException("--tables must be positive"); + } + if (namespace == null || namespace.isBlank()) { + throw new IllegalArgumentException("--namespace is required"); + } + if (!namespace.matches("[A-Za-z0-9_]+")) { + throw new IllegalArgumentException("--namespace may contain only letters, digits, and _"); + } + if (compactorResourceGroup == null || compactorResourceGroup.isBlank()) { + throw new IllegalArgumentException("--compactor-resource-group is required"); + } + ResourceGroupId.of(compactorResourceGroup); + if (coordinator && workerId >= 0) { + throw new IllegalArgumentException("--internal-worker-id is reserved for worker JVMs"); + } + if (!coordinator && (workerId < 0 || stateDir == null || runId == null)) { + throw new IllegalArgumentException("Missing internal worker configuration"); + } + if (tablePrefix == null || tablePrefix.isBlank()) { + if (!coordinator) { + throw new IllegalArgumentException("Worker table prefix was not supplied"); + } + } else if (!tablePrefix.matches("[A-Za-z0-9_]+")) { + throw new IllegalArgumentException("--table-prefix may contain only letters, digits, and _"); + } + double totalWeight = createWeight + deleteWeight + splitWeight + mergeWeight + + availabilityWeight + compactWeight + bulkImportWeight; + for (double weight : new double[] {createWeight, deleteWeight, splitWeight, mergeWeight, + availabilityWeight, compactWeight, bulkImportWeight}) { + if (!Double.isFinite(weight) || weight < 0) { + throw new IllegalArgumentException("Operation weights must be finite and non-negative"); + } + } + if (!Double.isFinite(totalWeight) || totalWeight <= 0) { + throw new IllegalArgumentException("At least one operation weight must be positive"); + } + if (bulkImportWeight > 0 && (hdfsDir == null || hdfsDir.isBlank())) { + throw new IllegalArgumentException("--hdfs-dir is required when bulk import is enabled"); + } + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressRows.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressRows.java new file mode 100644 index 00000000..03024b27 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressRows.java @@ -0,0 +1,201 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import java.math.BigInteger; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.List; +import java.util.Random; +import java.util.SortedSet; +import java.util.TreeSet; + +import org.apache.accumulo.core.client.admin.NewTableConfiguration; +import org.apache.accumulo.core.client.admin.TabletAvailability; +import org.apache.accumulo.core.data.RowRange; +import org.apache.hadoop.io.Text; + +/** Row keys and tablet ranges used by the Manager stress workload. */ +final class ManagerStressRows { + + static final int MAX_TABLETS = 1000; + static final int MIN_INITIAL_SPLITS = 10; + static final int MAX_INITIAL_SPLITS = MAX_TABLETS - 1; + + private static final BigInteger ROW_SPACE = BigInteger.ONE.shiftLeft(64); + + record MergeRange(Text start, Text end, int tabletCount) { + } + + record AvailabilityChange(int firstTablet, int tabletCount, TabletAvailability availability, + RowRange rowRange) { + } + + private ManagerStressRows() {} + + static NewTableConfiguration newTableConfiguration(Random random) { + int splitRange = MAX_INITIAL_SPLITS - MIN_INITIAL_SPLITS + 1; + int splitCount = MIN_INITIAL_SPLITS + random.nextInt(splitRange); + return new NewTableConfiguration().withInitialTabletAvailability(TabletAvailability.UNHOSTED) + .withSplits(initialSplits(splitCount)); + } + + static SortedSet initialSplits(int splitCount) { + if (splitCount < MIN_INITIAL_SPLITS || splitCount > MAX_INITIAL_SPLITS) { + throw new IllegalArgumentException("Initial split count must be between " + MIN_INITIAL_SPLITS + + " and " + MAX_INITIAL_SPLITS); + } + TreeSet splits = new TreeSet<>(); + BigInteger denominator = BigInteger.valueOf((long) splitCount + 1); + for (int i = 1; i <= splitCount; i++) { + BigInteger row = ROW_SPACE.multiply(BigInteger.valueOf(i)).divide(denominator); + splits.add(row(row)); + } + return splits; + } + + static Text randomRow(Random random) { + String hex = Long.toUnsignedString(random.nextLong(), 16); + return new Text("0".repeat(16 - hex.length()) + hex); + } + + static SortedSet additionalSplits(Collection currentSplits, Random random) { + TreeSet existing = new TreeSet<>(currentSplits); + int currentTablets = existing.size() + 1; + int available = MAX_TABLETS - currentTablets; + TreeSet additions = new TreeSet<>(); + if (available <= 0) { + return additions; + } + + int requested = 1 + random.nextInt(available); + while (additions.size() < requested) { + Text candidate = randomRow(random); + if (!existing.contains(candidate)) { + additions.add(candidate); + } + } + return additions; + } + + static MergeRange selectMergeRange(Collection currentSplits, int targetTabletCount, + Random random) { + List splits = new ArrayList<>(new TreeSet<>(currentSplits)); + int totalTablets = splits.size() + 1; + if (totalTablets < 2) { + return null; + } + int mergeTablets = Math.max(2, Math.min(Math.min(targetTabletCount, 100), totalTablets)); + int firstTablet = random.nextInt(totalTablets - mergeTablets + 1); + return mergeRange(splits, firstTablet, mergeTablets); + } + + static MergeRange mergeRange(Collection currentSplits, int firstTablet, int tabletCount) { + return mergeRange(new ArrayList<>(new TreeSet<>(currentSplits)), firstTablet, tabletCount); + } + + static MergeRange mergeRange(List sortedSplits, int firstTablet, int tabletCount) { + int totalTablets = sortedSplits.size() + 1; + if (tabletCount < 2 || firstTablet < 0 || firstTablet + tabletCount > totalTablets) { + throw new IllegalArgumentException("Invalid tablet range"); + } + Text start = firstTablet == 0 ? null : sortedSplits.get(firstTablet - 1); + int lastTablet = firstTablet + tabletCount - 1; + Text end = lastTablet == sortedSplits.size() ? null : sortedSplits.get(lastTablet); + return new MergeRange(start, end, tabletCount); + } + + static List randomAvailabilityChanges(List sortedSplits, + List currentAvailability, Random random) { + int totalTablets = sortedSplits.size() + 1; + if (currentAvailability.size() != totalTablets) { + throw new IllegalArgumentException("Availability must be supplied for every tablet"); + } + + List tabletIndices = new ArrayList<>(totalTablets); + for (int i = 0; i < totalTablets; i++) { + tabletIndices.add(i); + } + Collections.shuffle(tabletIndices, random); + int selectedCount = 1 + random.nextInt(totalTablets); + tabletIndices = new ArrayList<>(tabletIndices.subList(0, selectedCount)); + tabletIndices.sort(Integer::compareTo); + + List changes = new ArrayList<>(); + int runStart = tabletIndices.get(0); + int runEnd = runStart; + TabletAvailability runAvailability = + chooseDifferentAvailability(currentAvailability.get(runStart), random); + for (int i = 1; i < tabletIndices.size(); i++) { + int tablet = tabletIndices.get(i); + TabletAvailability nextAvailability = + chooseDifferentAvailability(currentAvailability.get(tablet), random); + if (tablet == runEnd + 1 && nextAvailability == runAvailability) { + runEnd = tablet; + } else { + changes.add(availabilityChange(sortedSplits, runStart, runEnd, runAvailability)); + runStart = tablet; + runEnd = tablet; + runAvailability = nextAvailability; + } + } + changes.add(availabilityChange(sortedSplits, runStart, runEnd, runAvailability)); + return List.copyOf(changes); + } + + static RowRange tabletRowRange(List sortedSplits, int firstTablet, int lastTablet) { + int totalTablets = sortedSplits.size() + 1; + if (firstTablet < 0 || lastTablet < firstTablet || lastTablet >= totalTablets) { + throw new IllegalArgumentException("Invalid tablet range"); + } + if (firstTablet == 0 && lastTablet == totalTablets - 1) { + return RowRange.all(); + } + if (lastTablet == totalTablets - 1) { + return RowRange.greaterThan(sortedSplits.get(firstTablet - 1)); + } + if (firstTablet == 0) { + return RowRange.atMost(sortedSplits.get(lastTablet)); + } + return RowRange.openClosed(sortedSplits.get(firstTablet - 1), sortedSplits.get(lastTablet)); + } + + private static AvailabilityChange availabilityChange(List sortedSplits, int firstTablet, + int lastTablet, TabletAvailability availability) { + return new AvailabilityChange(firstTablet, lastTablet - firstTablet + 1, availability, + tabletRowRange(sortedSplits, firstTablet, lastTablet)); + } + + private static TabletAvailability chooseDifferentAvailability(TabletAvailability current, + Random random) { + List candidates = new ArrayList<>(); + for (TabletAvailability availability : TabletAvailability.values()) { + if (availability != current) { + candidates.add(availability); + } + } + return candidates.get(random.nextInt(candidates.size())); + } + + private static Text row(BigInteger value) { + String hex = value.toString(16); + return new Text("0".repeat(16 - hex.length()) + hex); + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressWorker.java b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressWorker.java new file mode 100644 index 00000000..942e1922 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/ManagerStressWorker.java @@ -0,0 +1,420 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static java.nio.charset.StandardCharsets.UTF_8; + +import java.io.IOException; +import java.nio.file.Files; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Random; +import java.util.TreeSet; +import java.util.UUID; + +import org.apache.accumulo.core.client.Accumulo; +import org.apache.accumulo.core.client.AccumuloClient; +import org.apache.accumulo.core.client.TableExistsException; +import org.apache.accumulo.core.client.TableNotFoundException; +import org.apache.accumulo.core.client.TableOfflineException; +import org.apache.accumulo.core.client.admin.CompactionConfig; +import org.apache.accumulo.core.client.admin.TabletAvailability; +import org.apache.accumulo.core.client.admin.TabletInformation; +import org.apache.accumulo.core.client.rfile.RFile; +import org.apache.accumulo.core.client.rfile.RFileWriter; +import org.apache.accumulo.core.data.Key; +import org.apache.accumulo.core.data.RowRange; +import org.apache.accumulo.core.data.Value; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.FSDataInputStream; +import org.apache.hadoop.fs.FSDataOutputStream; +import org.apache.hadoop.fs.FileSystem; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.io.IOUtils; +import org.apache.hadoop.io.Text; +import org.apache.thrift.transport.TTransportException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +final class ManagerStressWorker { + + private static final Logger log = LoggerFactory.getLogger(ManagerStressWorker.class); + + private static class SubmissionTracker { + private final ManagerStressMetrics.Counts counts; + private boolean submitted; + + SubmissionTracker(ManagerStressMetrics.Counts counts) { + this.counts = counts; + } + + void markSubmitted() { + if (!submitted) { + counts.submitted++; + submitted = true; + } + } + } + + private final ManagerStressOptions options; + private final ManagerStressMetrics metrics = new ManagerStressMetrics(); + private java.nio.file.Path stateDirectory; + private long lastMetricsSnapshotNanos; + private Path workerBulkDir; + private Path seedFile; + private FileSystem fileSystem; + + ManagerStressWorker(ManagerStressOptions options) { + this.options = options; + } + + void run() throws Exception { + stateDirectory = java.nio.file.Path.of(options.stateDir); + TablePool tablePool = new TablePool(stateDirectory, options.namespace, options.tablePrefix, + options.tables, options.seed); + Random random = new Random(options.seed + (long) options.workerId * 0x9e3779b97f4a7c15L); + + try (AccumuloClient client = Accumulo.newClient().from(options.getClientProps()).build()) { + if (options.bulkImportWeight > 0) { + prepareBulkFile(random); + } + writeMetricsSnapshot(); + Files.createFile(stateDirectory.resolve("ready-" + options.workerId)); + long deadline = awaitStart(stateDirectory); + log.info("Worker {} started workload", options.workerId); + + while (System.currentTimeMillis() < deadline) { + try { + metrics.addMaintenanceCreates(replenish(tablePool, client)); + } catch (Exception e) { + log.warn("Worker {} could not replenish the table pool", options.workerId, e); + } + + List activeSlots = + tablePool.snapshot().stream().filter(slot -> slot.tableName() != null).toList(); + if (activeSlots.isEmpty()) { + Thread.sleep(100); + continue; + } + + Operation operation = Operation.choose(options, random); + ManagerStressMetrics.Counts operationCounts = metrics.counts(operation); + operationCounts.attempts++; + SubmissionTracker submission = new SubmissionTracker(operationCounts); + long start = System.nanoTime(); + try { + if (execute(operation, activeSlots, tablePool, client, random, submission)) { + if (operation.isAsynchronous()) { + operationCounts.asyncAccepted++; + } else { + operationCounts.completed++; + } + } else { + operationCounts.skipped++; + } + } catch (Exception e) { + if (isExpectedRace(e)) { + operationCounts.races++; + log.debug("Worker {} observed a table lifecycle race during {}", options.workerId, + operation, e); + } else { + operationCounts.failed++; + if (submission.submitted && isOutcomeUnknown(e)) { + operationCounts.outcomeUnknown++; + } + log.warn("Worker {} failed {} operation", options.workerId, operation, e); + } + } finally { + operationCounts.totalNanos += System.nanoTime() - start; + writeMetricsSnapshotIfDue(); + } + } + try { + metrics.addMaintenanceCreates(replenish(tablePool, client)); + } catch (Exception e) { + log.warn("Worker {} could not restore the table pool at shutdown", options.workerId, e); + } + } finally { + try { + if (fileSystem != null) { + try { + if (workerBulkDir != null) { + fileSystem.delete(workerBulkDir, true); + } + } catch (IOException e) { + log.warn("Worker {} could not remove its HDFS work directory {}", options.workerId, + workerBulkDir, e); + } finally { + fileSystem.close(); + } + } + } finally { + report(); + writeMetricsSnapshot(); + } + } + } + + private int replenish(TablePool pool, AccumuloClient client) throws Exception { + int created = pool.replenish(client); + if (created > 0) { + log.info("Worker {} replenished {} deleted table slot(s)", options.workerId, created); + } + return created; + } + + private boolean execute(Operation operation, List activeSlots, TablePool pool, + AccumuloClient client, Random random, SubmissionTracker submission) throws Exception { + TablePool.Slot slot = activeSlots.get(random.nextInt(activeSlots.size())); + String tableName = slot.tableName(); + return switch (operation) { + case CREATE -> { + if (!pool.replace(client, slot.index(), submission::markSubmitted)) { + throw new TableNotFoundException(null, tableName, + "all table slots are being updated by other workers"); + } + log.debug("Worker {} created a replacement table for slot {}", options.workerId, + slot.index()); + yield true; + } + case DELETE -> { + boolean deleted = pool.delete(client, slot.index(), submission::markSubmitted); + if (!deleted) { + throw new TableNotFoundException(null, tableName, "table was deleted by another worker"); + } + log.debug("Worker {} deleted table {}", options.workerId, tableName); + yield true; + } + case SPLIT -> split(client, pool, slot, random, submission); + case MERGE -> merge(client, pool, slot, random, submission); + case AVAILABILITY -> changeAvailability(client, pool, slot, random, submission); + case COMPACT -> { + submission.markSubmitted(); + client.tableOperations().compact(tableName, + new CompactionConfig().setFlush(true).setWait(false)); + log.debug("Worker {} submitted a compaction for table {}", options.workerId, tableName); + yield true; + } + case BULK_IMPORT -> { + bulkImport(client, tableName, submission); + yield true; + } + }; + } + + private boolean split(AccumuloClient client, TablePool pool, TablePool.Slot selected, + Random random, SubmissionTracker submission) throws Exception { + TablePool.Slot reservation = pool.lockTable(selected.index()); + if (reservation == null) { + return false; + } + try { + var tableOps = client.tableOperations(); + var currentSplits = new TreeSet<>(tableOps.listSplits(reservation.tableName())); + var newSplits = ManagerStressRows.additionalSplits(currentSplits, random); + if (newSplits.isEmpty()) { + log.debug("Worker {} skipped split for {} at the {}-tablet limit", options.workerId, + reservation.tableName(), ManagerStressRows.MAX_TABLETS); + return false; + } + submission.markSubmitted(); + tableOps.addSplits(reservation.tableName(), newSplits); + log.debug("Worker {} added {} split(s) to {}", options.workerId, newSplits.size(), + reservation.tableName()); + return true; + } finally { + pool.unlockTable(reservation); + } + } + + private boolean merge(AccumuloClient client, TablePool pool, TablePool.Slot selected, + Random random, SubmissionTracker submission) throws Exception { + TablePool.Slot reservation = pool.lockTable(selected.index()); + if (reservation == null) { + return false; + } + try { + var tableOps = client.tableOperations(); + var splits = tableOps.listSplits(reservation.tableName()); + int targetTablets = random.nextInt(100) + 1; + ManagerStressRows.MergeRange range = + ManagerStressRows.selectMergeRange(splits, targetTablets, random); + if (range == null) { + log.debug("Worker {} skipped merge for {} because it has fewer than two tablets", + options.workerId, reservation.tableName()); + return false; + } + submission.markSubmitted(); + tableOps.merge(reservation.tableName(), range.start(), range.end()); + log.debug("Worker {} merged {} tablets in {} over ({}, {}]", options.workerId, + range.tabletCount(), reservation.tableName(), range.start(), range.end()); + return true; + } finally { + pool.unlockTable(reservation); + } + } + + private boolean changeAvailability(AccumuloClient client, TablePool pool, TablePool.Slot selected, + Random random, SubmissionTracker submission) throws Exception { + TablePool.Slot reservation = pool.lockTable(selected.index()); + if (reservation == null) { + return false; + } + try { + var tableOps = client.tableOperations(); + List splits = new ArrayList<>(tableOps.listSplits(reservation.tableName())); + List currentAvailability = + new ArrayList<>(Collections.nCopies(splits.size() + 1, null)); + try (var tabletInformation = tableOps.getTabletInformation(reservation.tableName(), + List.of(RowRange.all()), TabletInformation.Field.AVAILABILITY)) { + var iterator = tabletInformation.iterator(); + while (iterator.hasNext()) { + var tablet = iterator.next(); + Text endRow = tablet.getTabletId().getEndRow(); + int tabletIndex = + endRow == null ? splits.size() : Collections.binarySearch(splits, endRow); + if (tabletIndex < 0 || tabletIndex >= currentAvailability.size() + || currentAvailability.set(tabletIndex, tablet.getTabletAvailability()) != null) { + throw new IllegalStateException( + "Could not map tablet availability to split boundaries for " + + reservation.tableName()); + } + } + } + if (currentAvailability.contains(null)) { + throw new IllegalStateException( + "Could not obtain availability for every tablet in " + reservation.tableName()); + } + + List changes = + ManagerStressRows.randomAvailabilityChanges(splits, currentAvailability, random); + int changedTablets = 0; + for (ManagerStressRows.AvailabilityChange change : changes) { + submission.markSubmitted(); + tableOps.setTabletAvailability(reservation.tableName(), change.rowRange(), + change.availability()); + changedTablets += change.tabletCount(); + } + log.debug("Worker {} submitted {} availability range change(s) for {} tablet(s) in {}", + options.workerId, changes.size(), changedTablets, reservation.tableName()); + return !changes.isEmpty(); + } finally { + pool.unlockTable(reservation); + } + } + + private void bulkImport(AccumuloClient client, String tableName, SubmissionTracker submission) + throws Exception { + Path importDir = + new Path(workerBulkDir, "import-" + UUID.randomUUID().toString().replace("-", "")); + Path stagedFile = new Path(importDir, "one-entry.rf"); + fileSystem.mkdirs(importDir); + try (FSDataInputStream in = fileSystem.open(seedFile); + FSDataOutputStream out = fileSystem.create(stagedFile, false)) { + IOUtils.copyBytes(in, out, new Configuration(), false); + } + try { + submission.markSubmitted(); + client.tableOperations().importDirectory(importDir.toString()).to(tableName).tableTime(true) + .load(); + log.debug("Worker {} bulk imported data into table {}", options.workerId, tableName); + } finally { + fileSystem.delete(importDir, true); + } + } + + private void prepareBulkFile(Random random) throws IOException { + Configuration configuration = new Configuration(); + Path configuredDir = new Path(options.hdfsDir); + fileSystem = configuredDir.getFileSystem(configuration); + workerBulkDir = new Path(new Path(configuredDir, options.runId), "worker-" + options.workerId); + if (!fileSystem.mkdirs(workerBulkDir) && !fileSystem.exists(workerBulkDir)) { + throw new IOException("Could not create HDFS work directory " + workerBulkDir); + } + seedFile = new Path(workerBulkDir, "seed.rf"); + Key key = new Key(ManagerStressRows.randomRow(random), new Text("cf"), new Text("cq")); + Value value = new Value(("manager-stress-" + options.runId).getBytes(UTF_8)); + try (RFileWriter writer = + RFile.newWriter().to(seedFile.toString()).withFileSystem(fileSystem).build()) { + writer.startDefaultLocalityGroup(); + writer.append(key, value); + } + log.info("Worker {} created one-entry RFile at {}", options.workerId, seedFile); + } + + private long awaitStart(java.nio.file.Path stateDir) throws Exception { + java.nio.file.Path startFile = stateDir.resolve("start"); + while (!Files.exists(startFile)) { + Thread.sleep(100); + } + return Long.parseLong(Files.readString(startFile, UTF_8)); + } + + private boolean isExpectedRace(Throwable error) { + for (Throwable current = error; current != null; current = current.getCause()) { + if (current instanceof TableNotFoundException || current instanceof TableExistsException + || current instanceof TableOfflineException) { + return true; + } + } + return false; + } + + private boolean isOutcomeUnknown(Throwable error) { + for (Throwable current = error; current != null; current = current.getCause()) { + if (current instanceof TTransportException || current instanceof InterruptedException) { + return true; + } + } + return false; + } + + private void writeMetricsSnapshotIfDue() { + if (System.nanoTime() - lastMetricsSnapshotNanos >= 1_000_000_000L) { + writeMetricsSnapshot(); + } + } + + private void writeMetricsSnapshot() { + try { + metrics.writeSnapshot(stateDirectory, options.workerId); + lastMetricsSnapshotNanos = System.nanoTime(); + } catch (IOException e) { + log.warn("Worker {} could not write its metrics snapshot", options.workerId, e); + } + } + + private void report() { + for (Operation operation : Operation.values()) { + ManagerStressMetrics.Counts result = metrics.counts(operation); + if (result.attempts > 0) { + double averageMillis = (double) result.totalNanos / result.attempts / 1_000_000; + log.info( + "Worker {} {}: attempts={}, submitted={}, completed={}, asyncAccepted={}, skipped={}, " + + "races={}, failed={}, outcomeUnknown={}, avgMs={}", + options.workerId, operation, result.attempts, result.submitted, result.completed, + result.asyncAccepted, result.skipped, result.races, result.failed, + result.outcomeUnknown, String.format("%.3f", averageMillis)); + } + } + log.info("Worker {} replacement-table maintenance creates={}", options.workerId, + metrics.maintenanceCreates()); + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/Operation.java b/src/main/java/org/apache/accumulo/testing/manager/stress/Operation.java new file mode 100644 index 00000000..433a3e8f --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/Operation.java @@ -0,0 +1,77 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import java.util.random.RandomGenerator; + +enum Operation { + CREATE, DELETE, SPLIT, MERGE, AVAILABILITY, COMPACT, BULK_IMPORT; + + boolean isAsynchronous() { + return this == AVAILABILITY || this == COMPACT; + } + + static Operation choose(ManagerStressOptions options, RandomGenerator random) { + double total = + options.createWeight + options.deleteWeight + options.splitWeight + options.mergeWeight + + options.availabilityWeight + options.compactWeight + options.bulkImportWeight; + return choose(options, random.nextDouble() * total); + } + + static Operation choose(ManagerStressOptions options, double selection) { + if ((selection -= options.createWeight) < 0) { + return CREATE; + } + if ((selection -= options.deleteWeight) < 0) { + return DELETE; + } + if ((selection -= options.splitWeight) < 0) { + return SPLIT; + } + if ((selection -= options.mergeWeight) < 0) { + return MERGE; + } + if ((selection -= options.availabilityWeight) < 0) { + return AVAILABILITY; + } + if ((selection -= options.compactWeight) < 0) { + return COMPACT; + } + if (options.bulkImportWeight > 0) { + return BULK_IMPORT; + } + // Guard floating point rounding at the final positive-weight boundary. + if (options.availabilityWeight > 0) { + return AVAILABILITY; + } + if (options.compactWeight > 0) { + return COMPACT; + } + if (options.mergeWeight > 0) { + return MERGE; + } + if (options.splitWeight > 0) { + return SPLIT; + } + if (options.deleteWeight > 0) { + return DELETE; + } + return CREATE; + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/TablePool.java b/src/main/java/org/apache/accumulo/testing/manager/stress/TablePool.java new file mode 100644 index 00000000..5f9d8e12 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/TablePool.java @@ -0,0 +1,294 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static java.nio.file.StandardCopyOption.ATOMIC_MOVE; +import static java.nio.file.StandardCopyOption.REPLACE_EXISTING; +import static java.nio.file.StandardOpenOption.CREATE; +import static java.nio.file.StandardOpenOption.WRITE; + +import java.io.IOException; +import java.nio.channels.FileChannel; +import java.nio.channels.FileLock; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.Random; + +import org.apache.accumulo.core.client.AccumuloClient; +import org.apache.accumulo.core.client.TableNotFoundException; +import org.apache.accumulo.core.client.admin.NewTableConfiguration; + +/** Shared, locally locked table-pool state for coordinator and worker JVMs. */ +final class TablePool { + + record Slot(int index, long generation, String tableName, boolean busy) { + } + + @FunctionalInterface + private interface LockedFunction { + T apply(List slots) throws Exception; + } + + private final Path stateDir; + private final Path stateFile; + private final Path lockFile; + private final String namespace; + private final String tablePrefix; + private final int tableCount; + private final long splitSeed; + + TablePool(Path stateDir, String namespace, String tablePrefix, int tableCount, long splitSeed) { + this.stateDir = stateDir; + this.stateFile = stateDir.resolve("tables.state"); + this.lockFile = stateDir.resolve("tables.lock"); + this.namespace = namespace; + this.tablePrefix = tablePrefix; + this.tableCount = tableCount; + this.splitSeed = splitSeed; + } + + void initialize(AccumuloClient client) throws Exception { + Files.createDirectories(stateDir); + List slots = new ArrayList<>(tableCount); + try { + for (int i = 0; i < tableCount; i++) { + String name = tableName(i, 1); + if (client.tableOperations().exists(name)) { + throw new IllegalStateException("Refusing to use existing table " + name); + } + createTable(client, name, i, 1); + slots.add(new Slot(i, 1, name, false)); + } + writeSlots(slots); + } catch (Exception e) { + for (Slot slot : slots) { + try { + client.tableOperations().delete(slot.tableName()); + } catch (Exception cleanupError) { + e.addSuppressed(cleanupError); + } + } + throw e; + } + } + + List snapshot() throws Exception { + return withLock(slots -> List.copyOf(slots)); + } + + int replenish(AccumuloClient client) throws Exception { + int created = 0; + for (int i = 0; i < tableCount; i++) { + Slot reserved = reserve(i, true, true); + if (reserved != null) { + createInVacantSlot(client, reserved); + created++; + } + } + return created; + } + + boolean delete(AccumuloClient client, int slotIndex, Runnable markSubmitted) throws Exception { + Slot reserved = reserve(slotIndex, false, false); + if (reserved == null) { + return false; + } + try { + markSubmitted.run(); + client.tableOperations().delete(reserved.tableName()); + } catch (TableNotFoundException e) { + updateSlot(new Slot(reserved.index(), reserved.generation(), null, false)); + return false; + } catch (Exception e) { + updateSlot(new Slot(reserved.index(), reserved.generation(), reserved.tableName(), false)); + throw e; + } + updateSlot(new Slot(reserved.index(), reserved.generation(), null, false)); + return true; + } + + boolean replace(AccumuloClient client, int slotIndex, Runnable markSubmitted) throws Exception { + Slot reserved = reserve(slotIndex, false, true); + if (reserved == null) { + return false; + } + String replacementName = tableName(reserved.index(), reserved.generation()); + boolean replacementCreated = false; + boolean oldTableDeleted = false; + try { + markSubmitted.run(); + createTable(client, replacementName, reserved.index(), reserved.generation()); + replacementCreated = true; + if (reserved.tableName() != null) { + try { + client.tableOperations().delete(reserved.tableName()); + oldTableDeleted = true; + } catch (TableNotFoundException e) { + // A competing client completed the deletion first. + oldTableDeleted = true; + } + } + updateSlot(new Slot(reserved.index(), reserved.generation(), replacementName, false)); + return true; + } catch (Exception e) { + if (replacementCreated) { + try { + client.tableOperations().delete(replacementName); + } catch (Exception cleanupError) { + e.addSuppressed(cleanupError); + } + } + String previousName = oldTableDeleted ? null : reserved.tableName(); + try { + updateSlot(new Slot(reserved.index(), reserved.generation(), previousName, false)); + } catch (Exception updateError) { + e.addSuppressed(updateError); + } + throw e; + } + } + + private void createInVacantSlot(AccumuloClient client, Slot reserved) throws Exception { + String replacementName = tableName(reserved.index(), reserved.generation()); + boolean created = false; + try { + createTable(client, replacementName, reserved.index(), reserved.generation()); + created = true; + updateSlot(new Slot(reserved.index(), reserved.generation(), replacementName, false)); + } catch (Exception e) { + if (created) { + try { + client.tableOperations().delete(replacementName); + } catch (Exception cleanupError) { + e.addSuppressed(cleanupError); + } + } + try { + updateSlot(new Slot(reserved.index(), reserved.generation(), null, false)); + } catch (Exception updateError) { + e.addSuppressed(updateError); + } + throw e; + } + } + + Slot lockTable(int preferredIndex) throws Exception { + return reserve(preferredIndex, false, false); + } + + void unlockTable(Slot reservation) throws Exception { + withLock(slots -> { + Slot current = slots.get(reservation.index()); + if (!current.busy() || current.generation() != reservation.generation() + || !java.util.Objects.equals(current.tableName(), reservation.tableName())) { + throw new IllegalStateException( + "Table slot reservation changed while in use: " + reservation.index()); + } + slots.set(reservation.index(), + new Slot(reservation.index(), reservation.generation(), reservation.tableName(), false)); + writeSlots(slots); + return null; + }); + } + + private void createTable(AccumuloClient client, String tableName, int slot, long generation) + throws Exception { + long seed = splitSeed ^ ((long) slot * 0x9e3779b97f4a7c15L) ^ Long.rotateLeft(generation, 32); + NewTableConfiguration config = ManagerStressRows.newTableConfiguration(new Random(seed)); + client.tableOperations().create(tableName, config); + } + + private Slot reserve(int preferredIndex, boolean vacant, boolean incrementGeneration) + throws Exception { + return withLock(slots -> { + for (int offset = 0; offset < slots.size(); offset++) { + int index = (preferredIndex + offset) % slots.size(); + Slot current = slots.get(index); + if (!current.busy() && (vacant == (current.tableName() == null))) { + long generation = current.generation() + (incrementGeneration ? 1 : 0); + Slot reserved = new Slot(index, generation, current.tableName(), true); + slots.set(index, reserved); + writeSlots(slots); + return reserved; + } + } + return null; + }); + } + + private void updateSlot(Slot updated) throws Exception { + withLock(slots -> { + slots.set(updated.index(), updated); + writeSlots(slots); + return null; + }); + } + + String tableName(int slot, long generation) { + return namespace + "." + tablePrefix + "_" + slot + "_" + generation; + } + + private T withLock(LockedFunction function) throws Exception { + try (FileChannel channel = FileChannel.open(lockFile, CREATE, WRITE); + FileLock ignored = channel.lock()) { + return function.apply(readSlots()); + } + } + + private List readSlots() throws IOException { + List lines = Files.readAllLines(stateFile, UTF_8); + if (lines.size() != tableCount) { + throw new IOException("Invalid table-pool state in " + stateFile); + } + List slots = new ArrayList<>(tableCount); + for (int i = 0; i < lines.size(); i++) { + String[] values = lines.get(i).split("\\t", -1); + if (values.length != 3) { + throw new IOException("Invalid table-pool state line: " + lines.get(i)); + } + long generation = Long.parseLong(values[0]); + String name = values[1].isEmpty() ? null : values[1]; + slots.add(new Slot(i, generation, name, Boolean.parseBoolean(values[2]))); + } + return slots; + } + + private void writeSlots(List slots) throws IOException { + Path tempFile = Files.createTempFile(stateDir, "tables", ".tmp"); + List< + String> lines = + slots.stream() + .map(slot -> slot.generation() + "\t" + + (slot.tableName() == null ? "" : slot.tableName()) + "\t" + slot.busy()) + .toList(); + try { + Files.write(tempFile, lines, UTF_8); + try { + Files.move(tempFile, stateFile, ATOMIC_MOVE, REPLACE_EXISTING); + } catch (UnsupportedOperationException | IOException e) { + Files.move(tempFile, stateFile, REPLACE_EXISTING); + } + } finally { + Files.deleteIfExists(tempFile); + } + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcome.java b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcome.java new file mode 100644 index 00000000..da6b05f1 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcome.java @@ -0,0 +1,46 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress.compactor; + +import java.util.random.RandomGenerator; + +enum CompactionOutcome { + SUCCESS, FAILURE, CANCELLATION; + + static CompactionOutcome choose(StressCompactorOptions options, RandomGenerator random) { + double total = options.successWeight + options.failureWeight + options.cancellationWeight; + return choose(options, random.nextDouble() * total); + } + + static CompactionOutcome choose(StressCompactorOptions options, double selection) { + if ((selection -= options.successWeight) < 0) { + return SUCCESS; + } + if ((selection -= options.failureWeight) < 0) { + return FAILURE; + } + if (options.cancellationWeight > 0) { + return CANCELLATION; + } + if (options.failureWeight > 0) { + return FAILURE; + } + return SUCCESS; + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/SimulatedCompactor.java b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/SimulatedCompactor.java new file mode 100644 index 00000000..e190abab --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/SimulatedCompactor.java @@ -0,0 +1,485 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress.compactor; + +import static java.util.concurrent.TimeUnit.MILLISECONDS; + +import java.util.List; +import java.util.UUID; +import java.util.concurrent.Semaphore; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.atomic.LongAdder; +import java.util.random.RandomGenerator; + +import org.apache.accumulo.core.cli.ServerOpts; +import org.apache.accumulo.core.client.Accumulo; +import org.apache.accumulo.core.client.AccumuloClient; +import org.apache.accumulo.core.client.AccumuloSecurityException; +import org.apache.accumulo.core.client.admin.servers.ServerId.Type; +import org.apache.accumulo.core.clientImpl.ClientContext; +import org.apache.accumulo.core.clientImpl.thrift.SecurityErrorCode; +import org.apache.accumulo.core.clientImpl.thrift.TInfo; +import org.apache.accumulo.core.clientImpl.thrift.ThriftSecurityException; +import org.apache.accumulo.core.compaction.thrift.CompactorService; +import org.apache.accumulo.core.compaction.thrift.TCompactionState; +import org.apache.accumulo.core.compaction.thrift.TCompactionStatusUpdate; +import org.apache.accumulo.core.compaction.thrift.TExternalCompaction; +import org.apache.accumulo.core.compaction.thrift.TNextCompactionJob; +import org.apache.accumulo.core.compaction.thrift.UnknownCompactionIdException; +import org.apache.accumulo.core.conf.Property; +import org.apache.accumulo.core.conf.SiteConfiguration; +import org.apache.accumulo.core.data.ResourceGroupId; +import org.apache.accumulo.core.lock.ServiceLock; +import org.apache.accumulo.core.lock.ServiceLock.LockLossReason; +import org.apache.accumulo.core.lock.ServiceLock.LockWatcher; +import org.apache.accumulo.core.lock.ServiceLockData; +import org.apache.accumulo.core.lock.ServiceLockData.ServiceDescriptor; +import org.apache.accumulo.core.lock.ServiceLockData.ServiceDescriptors; +import org.apache.accumulo.core.lock.ServiceLockData.ThriftService; +import org.apache.accumulo.core.lock.ServiceLockPaths.ServiceLockPath; +import org.apache.accumulo.core.lock.ServiceLockSupport; +import org.apache.accumulo.core.metadata.schema.ExternalCompactionId; +import org.apache.accumulo.core.metrics.MetricsInfo; +import org.apache.accumulo.core.rpc.ThriftUtil; +import org.apache.accumulo.core.securityImpl.thrift.TCredentials; +import org.apache.accumulo.core.tabletserver.thrift.ActiveCompaction; +import org.apache.accumulo.core.tabletserver.thrift.TCompactionStats; +import org.apache.accumulo.core.tabletserver.thrift.TExternalCompactionJob; +import org.apache.accumulo.core.trace.TraceUtil; +import org.apache.accumulo.core.util.compaction.ExternalCompactionUtil; +import org.apache.accumulo.server.AbstractServer; +import org.apache.accumulo.server.ServerContext; +import org.apache.accumulo.server.client.ClientServiceHandler; +import org.apache.accumulo.server.rpc.ServerAddress; +import org.apache.accumulo.server.rpc.TServerUtils; +import org.apache.accumulo.server.rpc.ThriftProcessorTypes; +import org.apache.thrift.TException; +import org.apache.zookeeper.KeeperException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.common.net.HostAndPort; + +/** A registered Compactor identity which polls the Manager without doing file compaction. */ +final class SimulatedCompactor extends AbstractServer implements CompactorService.Iface { + + private static final Logger log = LoggerFactory.getLogger(SimulatedCompactor.class); + + static final class Counters { + final LongAdder jobRequests = new LongAdder(); + final LongAdder idleResponses = new LongAdder(); + final LongAdder jobsReceived = new LongAdder(); + final LongAdder successes = new LongAdder(); + final LongAdder failures = new LongAdder(); + final LongAdder cancellations = new LongAdder(); + final LongAdder rpcFailures = new LongAdder(); + final LongAdder outcomeReportFailures = new LongAdder(); + } + + private final StressCompactorOptions options; + private final ServerContext serverContext; + private final int instanceNumber; + private final UUID serverUuid = UUID.randomUUID(); + private final ResourceGroupId group; + private final ClientContext clientContext; + private final RandomGenerator random; + private final AtomicBoolean stopping = new AtomicBoolean(); + private final AtomicBoolean closed = new AtomicBoolean(); + private final AtomicBoolean cancelRequested = new AtomicBoolean(); + private final AtomicReference currentId = new AtomicReference<>(); + private final AtomicReference currentCompaction = new AtomicReference<>(); + private final Semaphore wakeups = new Semaphore(0); + private final Counters counters = new Counters(); + + private ServerAddress serviceAddress; + private ServiceLock serviceLock; + private Thread worker; + private volatile String advertisedAddress; + + SimulatedCompactor(StressCompactorOptions options, int instanceNumber) throws Exception { + super(Type.COMPACTOR, new ServerOpts(), ServerContext::new, serverArgs(options)); + this.options = options; + this.instanceNumber = instanceNumber; + this.group = getResourceGroup(); + this.serverContext = getContext(); + ClientContext newClientContext; + try { + newClientContext = createClientContext(options.getClientProps()); + } catch (Exception e) { + try { + super.close(); + } catch (Exception cleanupError) { + e.addSuppressed(cleanupError); + } + throw e; + } + this.clientContext = newClientContext; + String configuredGroup = + serverContext.getSiteConfiguration().get(Property.COMPACTOR_GROUP_NAME); + if (!group.canonical().equals(configuredGroup) + || !clientContext.getInstanceID().equals(serverContext.getInstanceID())) { + try { + clientContext.close(); + } finally { + super.close(); + } + throw new IllegalStateException("Simulated Compactor configuration mismatch: group=" + + configuredGroup + ", expectedGroup=" + group.canonical() + + "; client and server contexts must also use the same Accumulo instance"); + } + this.random = new java.util.Random(options.seed + instanceNumber * 0x9e3779b97f4a7c15L); + } + + private static String[] serverArgs(StressCompactorOptions options) { + String compactorGroup = ResourceGroupId.of(options.resourceGroup).canonical(); + return new String[] {"-o", Property.COMPACTOR_GROUP_NAME.getKey() + "=" + compactorGroup}; + } + + @Override + protected String getResourceGroupPropertyValue(SiteConfiguration siteConfig) { + return siteConfig.get(Property.COMPACTOR_GROUP_NAME); + } + + @Override + public ServiceLock getLock() { + return serviceLock; + } + + Counters getCounters() { + return counters; + } + + void start() throws Exception { + var processor = ThriftProcessorTypes.getCompactorTProcessor(this, + new ClientServiceHandler(serverContext), this, serverContext); + serviceAddress = createThriftServer(processor); + serviceAddress.startThriftServer("SimulatedCompactorService-" + instanceNumber); + updateAdvertiseAddress(serviceAddress.address); + advertisedAddress = serviceAddress.address.toString(); + + ServiceLockPath lockPath = + serverContext.getServerPaths().createCompactorPath(group, serviceAddress.address); + ServiceLockSupport.createNonHaServiceLockPath(Type.COMPACTOR, + serverContext.getZooSession().asReaderWriter(), lockPath); + serviceLock = new ServiceLock(serverContext.getZooSession(), lockPath, serverUuid); + ServiceDescriptors descriptors = new ServiceDescriptors(); + descriptors.addService( + new ServiceDescriptor(serverUuid, ThriftService.CLIENT, advertisedAddress, group)); + descriptors.addService( + new ServiceDescriptor(serverUuid, ThriftService.COMPACTOR, advertisedAddress, group)); + LockWatcher lockWatcher = new LockWatcher() { + @Override + public void lostLock(LockLossReason reason) { + log.warn("Simulated Compactor {} lost its service registration: {}", instanceNumber, + reason); + stopping.set(true); + wakeups.release(); + } + + @Override + public void unableToMonitorLockNode(Exception e) { + log.warn("Unable to monitor simulated Compactor {} service registration", instanceNumber, + e); + stopping.set(true); + wakeups.release(); + } + }; + if (!serviceLock.tryLock(lockWatcher, new ServiceLockData(descriptors))) { + serviceAddress.server.stop(); + throw new IllegalStateException( + "Could not register simulated Compactor at " + advertisedAddress); + } + + log.info("Started simulated Compactor {} at {} in resource group {}", instanceNumber, + advertisedAddress, group); + } + + void startPolling() { + worker = new Thread(this, "CompactorSimulator-" + instanceNumber); + worker.setDaemon(false); + worker.start(); + } + + @Override + public void run() { + runLoop(); + } + + private ServerAddress createThriftServer(org.apache.thrift.TProcessor processor) { + var conf = serverContext.getConfiguration(); + var address = options.addresses(); + long threadTimeout = + Math.max(10_000, conf.getTimeInMillis(Property.COMPACTOR_MINTHREADS_TIMEOUT)); + return TServerUtils.createThriftServer(conf, serverContext.getThriftServerType(), processor, + serverContext.getInstanceID(), "SimulatedCompactor-" + instanceNumber, 1, threadTimeout, + conf.getTimeInMillis(Property.COMPACTOR_THREADCHECK), + conf.getAsBytes(Property.RPC_MAX_MESSAGE_SIZE), serverContext.getServerSslParams(), + serverContext.getSaslParams(), serverContext.getClientTimeoutInMillis(), + conf.getCount(Property.RPC_BACKLOG), serverContext.getMetricsInfo(), + address.toArray(HostAndPort[]::new)); + } + + private static ClientContext createClientContext(java.util.Properties clientProperties) { + AccumuloClient client = Accumulo.newClient().from(clientProperties).build(); + if (client instanceof ClientContext context) { + return context; + } + try { + client.close(); + } catch (Exception e) { + log.warn("Could not close unsupported client implementation", e); + } + throw new IllegalStateException("Expected the Accumulo client builder to return ClientContext"); + } + + private void runLoop() { + var metricsInfo = serverContext.getMetricsInfo(); + metricsInfo.addMetricsProducers(this); + try { + metricsInfo.init(MetricsInfo.serviceTags(serverContext.getInstanceName(), + getApplicationName(), serviceAddress.address, group)); + } catch (IllegalStateException e) { + // Catch this because it is thrown due to many Compactor instances in the same VM. + } + + long minWait = + serverContext.getConfiguration().getTimeInMillis(Property.COMPACTOR_MIN_JOB_WAIT_TIME); + long maxWait = + serverContext.getConfiguration().getTimeInMillis(Property.COMPACTOR_MAX_JOB_WAIT_TIME); + while (!stopping.get()) { + String externalId = ExternalCompactionId.generate(UUID.randomUUID()).canonical(); + boolean assignedJob = false; + cancelRequested.set(false); + currentId.set(externalId); + counters.jobRequests.increment(); + try { + TNextCompactionJob next = requestJob(externalId); + TExternalCompactionJob job = next.getJob(); + if (!job.isSetExternalCompactionId()) { + currentId.compareAndSet(externalId, null); + counters.idleResponses.increment(); + waitForWork(calculateWait(next.getCompactorCount(), minWait, maxWait)); + continue; + } + + counters.jobsReceived.increment(); + assignedJob = true; + if (!externalId.equals(job.getExternalCompactionId())) { + throw new IllegalStateException("Manager returned compaction ID " + + job.getExternalCompactionId() + " for request " + externalId); + } + TExternalCompaction running = new TExternalCompaction(); + running.setCompactor(advertisedAddress); + running.setGroupName(group.canonical()); + long startTime = System.currentTimeMillis(); + running.setStartTime(startTime); + running.putToUpdates(startTime, new TCompactionStatusUpdate(TCompactionState.STARTED, + "Simulated compaction started", -1, -1, -1, 0)); + running.setJob(job); + currentCompaction.set(running); + + CompactionOutcome outcome = CompactionOutcome.choose(options, random); + if (cancelRequested.get()) { + outcome = CompactionOutcome.CANCELLATION; + } + reportOutcome(job, outcome); + } catch (Exception e) { + counters.rpcFailures.increment(); + if (assignedJob) { + counters.outcomeReportFailures.increment(); + } + log.warn("Simulated Compactor {} failed a Manager RPC", instanceNumber, e); + waitForWork(minWait); + } finally { + currentCompaction.set(null); + currentId.set(null); + cancelRequested.set(false); + } + } + } + + private TNextCompactionJob requestJob(String externalId) throws Exception { + var coordinator = ExternalCompactionUtil.getCoordinatorClient(clientContext); + try { + return coordinator.getCompactionJob(TraceUtil.traceInfo(), clientContext.rpcCreds(), + group.canonical(), advertisedAddress, externalId); + } finally { + ThriftUtil.returnClient(coordinator, clientContext); + } + } + + private void reportOutcome(TExternalCompactionJob job, CompactionOutcome outcome) + throws Exception { + var coordinator = ExternalCompactionUtil.getCoordinatorClient(clientContext); + try { + switch (outcome) { + case SUCCESS -> { + TCompactionStats stats = new TCompactionStats(); + stats.setEntriesRead(0); + stats.setEntriesWritten(0); + stats.setFileSize(0); + coordinator.compactionCompleted(TraceUtil.traceInfo(), clientContext.rpcCreds(), + job.getExternalCompactionId(), job.getExtent(), stats, group.canonical(), + advertisedAddress); + counters.successes.increment(); + } + case FAILURE -> { + coordinator.compactionFailed(TraceUtil.traceInfo(), clientContext.rpcCreds(), + job.getExternalCompactionId(), job.getExtent(), "simulated compaction failure", + TCompactionState.FAILED, group.canonical(), advertisedAddress); + counters.failures.increment(); + } + case CANCELLATION -> { + coordinator.compactionFailed(TraceUtil.traceInfo(), clientContext.rpcCreds(), + job.getExternalCompactionId(), job.getExtent(), "simulated compaction cancellation", + TCompactionState.CANCELLED, group.canonical(), advertisedAddress); + counters.cancellations.increment(); + } + } + TCompactionState state = switch (outcome) { + case SUCCESS -> TCompactionState.SUCCEEDED; + case FAILURE -> TCompactionState.FAILED; + case CANCELLATION -> TCompactionState.CANCELLED; + }; + log.info("Simulated compaction {} finished with state {}", job.getExternalCompactionId(), + state); + } finally { + ThriftUtil.returnClient(coordinator, clientContext); + } + } + + private static long calculateWait(int compactorCount, long minWait, long maxWait) { + long wait = Math.max(minWait, (long) compactorCount * minWait / 3); + wait = Math.min(maxWait, wait); + return (long) (wait * (0.9 + 0.2 * ThreadLocalRandom.current().nextDouble())); + } + + private void waitForWork(long millis) { + if (stopping.get()) { + return; + } + try { + wakeups.tryAcquire(Math.max(1, millis), MILLISECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + stopping.set(true); + } + } + + private void requireSystemAction(TCredentials credentials) throws ThriftSecurityException { + if (!serverContext.getSecurityOperation().canPerformSystemActions(credentials)) { + throw new AccumuloSecurityException(credentials.getPrincipal(), + SecurityErrorCode.PERMISSION_DENIED).asThriftException(); + } + } + + @Override + public void gracefulShutdown(TCredentials credentials) { + super.gracefulShutdown(credentials); + if (isShutdownRequested()) { + stopping.set(true); + wakeups.release(); + } + } + + @Override + public void close() { + if (!closed.compareAndSet(false, true)) { + return; + } + stopping.set(true); + wakeups.release(); + if (worker != null) { + worker.interrupt(); + try { + worker.join(10_000); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + if (serviceLock != null) { + try { + serviceLock.unlock(); + } catch (KeeperException | InterruptedException e) { + if (e instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } + log.warn("Could not unregister simulated Compactor {}", advertisedAddress, e); + } + } + if (serviceAddress != null) { + serviceAddress.server.stop(); + } + try { + clientContext.close(); + } catch (Exception e) { + log.warn("Could not close client context for simulated Compactor {}", instanceNumber, e); + } + super.close(); + log.info( + "Compactor {} totals: polls={}, idle={}, jobs={}, success={}, failure={}, cancelled={}, " + + "rpcFailures={}, outcomeReportFailures={}", + instanceNumber, counters.jobRequests.sum(), counters.idleResponses.sum(), + counters.jobsReceived.sum(), counters.successes.sum(), counters.failures.sum(), + counters.cancellations.sum(), counters.rpcFailures.sum(), + counters.outcomeReportFailures.sum()); + } + + @Override + public TExternalCompaction getRunningCompaction(TInfo tinfo, TCredentials credentials) + throws ThriftSecurityException { + requireSystemAction(credentials); + TExternalCompaction current = currentCompaction.get(); + return current == null ? new TExternalCompaction() : new TExternalCompaction(current); + } + + @Override + public String getRunningCompactionId(TInfo tinfo, TCredentials credentials) + throws ThriftSecurityException { + requireSystemAction(credentials); + String id = currentId.get(); + return id == null ? "" : id; + } + + @Override + public List getActiveCompactions(TInfo tinfo, TCredentials credentials) + throws ThriftSecurityException { + requireSystemAction(credentials); + return List.of(); + } + + public void wake(TInfo tinfo, TCredentials credentials) throws ThriftSecurityException { + requireSystemAction(credentials); + wakeups.release(); + } + + @Override + public void cancel(TInfo tinfo, TCredentials credentials, String externalCompactionId) + throws TException { + requireSystemAction(credentials); + if (externalCompactionId.equals(currentId.get())) { + cancelRequested.set(true); + } else { + throw new UnknownCompactionIdException(); + } + } + +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactor.java b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactor.java new file mode 100644 index 00000000..b0435f94 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactor.java @@ -0,0 +1,105 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress.compactor; + +import java.net.InetAddress; +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** Starts multiple registered Compactor simulators as threads in this JVM. */ +public class StressCompactor { + + private static final Logger log = LoggerFactory.getLogger(StressCompactor.class); + + public static void main(String[] args) throws Exception { + if (args.length == 0) { + args = new String[] {"--help"}; + } + StressCompactorOptions options = new StressCompactorOptions(); + options.parseArgs(StressCompactor.class.getName(), args); + if (options.host == null || options.host.isBlank()) { + options.host = InetAddress.getLocalHost().getCanonicalHostName(); + } + options.validate(); + + List compactors = new ArrayList<>(options.compactorsPerJvm); + AtomicBoolean stopping = new AtomicBoolean(); + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + if (stopping.compareAndSet(false, true)) { + stopCompactors(compactors); + } + }, "compactor-stress-shutdown")); + + try { + for (int i = 0; i < options.compactorsPerJvm; i++) { + SimulatedCompactor compactor = new SimulatedCompactor(options, i); + compactors.add(compactor); + compactor.start(); + } + + long deadline = Math.addExact(System.currentTimeMillis(), options.duration); + compactors.forEach(SimulatedCompactor::startPolling); + log.info("Started {} Compactor simulator(s) in JVM {} for {} ms", compactors.size(), + ProcessHandle.current().pid(), options.duration); + while (System.currentTimeMillis() < deadline) { + long remaining = deadline - System.currentTimeMillis(); + Thread.sleep(Math.min(remaining, Duration.ofSeconds(1).toMillis())); + } + } finally { + if (stopping.compareAndSet(false, true)) { + stopCompactors(compactors); + } + } + } + + private static void stopCompactors(List compactors) { + for (int i = compactors.size() - 1; i >= 0; i--) { + compactors.get(i).close(); + } + long polls = 0; + long idle = 0; + long jobs = 0; + long successes = 0; + long failures = 0; + long cancellations = 0; + long rpcFailures = 0; + long outcomeReportFailures = 0; + for (SimulatedCompactor compactor : compactors) { + SimulatedCompactor.Counters counters = compactor.getCounters(); + polls += counters.jobRequests.sum(); + idle += counters.idleResponses.sum(); + jobs += counters.jobsReceived.sum(); + successes += counters.successes.sum(); + failures += counters.failures.sum(); + cancellations += counters.cancellations.sum(); + rpcFailures += counters.rpcFailures.sum(); + outcomeReportFailures += counters.outcomeReportFailures.sum(); + } + log.info( + "Compactor JVM totals: simulators={}, polls={}, idle={}, jobs={}, success={}, failure={}, " + + "cancelled={}, rpcFailures={}, outcomeReportFailures={}", + compactors.size(), polls, idle, jobs, successes, failures, cancellations, rpcFailures, + outcomeReportFailures); + } +} diff --git a/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactorOptions.java b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactorOptions.java new file mode 100644 index 00000000..609444b6 --- /dev/null +++ b/src/main/java/org/apache/accumulo/testing/manager/stress/compactor/StressCompactorOptions.java @@ -0,0 +1,112 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress.compactor; + +import java.util.ArrayList; +import java.util.List; + +import org.apache.accumulo.testing.cli.ClientOpts; + +import com.beust.jcommander.Parameter; +import com.google.common.net.HostAndPort; + +class StressCompactorOptions extends ClientOpts { + + @Parameter(names = "--duration", converter = TimeConverter.class, + description = "duration of this simulator JVM (for example 10m or 1h)") + long duration = 5 * 60 * 1000L; + + @Parameter(names = "--compactors-per-jvm", + description = "number of simulated Compactor instances in this JVM") + int compactorsPerJvm = 1; + + @Parameter(names = "--resource-group", + description = "Compactor resource group to register and request jobs from") + String resourceGroup = "default"; + + @Parameter(names = "--host", + description = "local host/interface reachable by the Manager for callbacks") + String host; + + @Parameter(names = "--port-range", + description = "Compactor callback port or inclusive range; 0 selects an ephemeral port") + String portRange = "0"; + + @Parameter(names = "--success-weight", + description = "relative weight for successful zero-output completions") + double successWeight = 1; + + @Parameter(names = "--failure-weight", description = "relative weight for failed jobs") + double failureWeight = 1; + + @Parameter(names = "--cancellation-weight", description = "relative weight for cancelled jobs") + double cancellationWeight = 1; + + @Parameter(names = "--seed", description = "random seed for outcome selection") + long seed = System.nanoTime(); + + void validate() { + if (duration <= 0) { + throw new IllegalArgumentException("--duration must be positive"); + } + if (compactorsPerJvm <= 0) { + throw new IllegalArgumentException("--compactors-per-jvm must be positive"); + } + if (resourceGroup == null || resourceGroup.isBlank()) { + throw new IllegalArgumentException("--resource-group must not be blank"); + } + if (host == null || host.isBlank()) { + throw new IllegalArgumentException( + "--host must be a host/interface reachable by the Manager"); + } + double totalWeight = successWeight + failureWeight + cancellationWeight; + for (double weight : new double[] {successWeight, failureWeight, cancellationWeight}) { + if (!Double.isFinite(weight) || weight < 0) { + throw new IllegalArgumentException("Outcome weights must be finite and non-negative"); + } + } + if (!Double.isFinite(totalWeight) || totalWeight <= 0) { + throw new IllegalArgumentException("At least one outcome weight must be positive"); + } + if (portRange == null || !portRange.matches("(?:0|[1-9][0-9]{0,4}(?:-[1-9][0-9]{0,4})?)")) { + throw new IllegalArgumentException("--port-range must be a port, 0, or inclusive port range"); + } + String[] ports = portRange.split("-"); + int first = Integer.parseInt(ports[0]); + int last = ports.length == 1 ? first : Integer.parseInt(ports[1]); + if (first > 65535 || last > 65535 || last < first || first == 0 && last != 0) { + throw new IllegalArgumentException("Invalid --port-range: " + portRange); + } + if (last - first > 4095) { + throw new IllegalArgumentException("--port-range cannot contain more than 4096 ports"); + } + } + + List addresses() { + String[] ports = portRange.split("-"); + int first = Integer.parseInt(ports[0]); + int last = ports.length == 1 ? first : Integer.parseInt(ports[1]); + List addresses = new ArrayList<>(last - first + 1); + for (int port = first; port <= last; port++) { + addresses.add(HostAndPort.fromParts(host, port)); + } + return addresses; + } + +} diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressCleanupTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressCleanupTest.java new file mode 100644 index 00000000..6f80b8d2 --- /dev/null +++ b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressCleanupTest.java @@ -0,0 +1,62 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertFalse; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class ManagerStressCleanupTest { + + @TempDir + Path temporaryDirectory; + + @Test + void controlDirectoryCleanupIsIdempotent() throws IOException { + Path controlDirectory = temporaryDirectory.resolve("manager-stress-control"); + Path nestedFile = controlDirectory.resolve("nested").resolve("snapshot.properties"); + Files.createDirectories(nestedFile.getParent()); + Files.writeString(nestedFile, "metrics=true"); + + ManagerStress.deleteRecursivelyIfExists(controlDirectory); + + assertFalse(Files.exists(controlDirectory)); + assertDoesNotThrow(() -> ManagerStress.deleteRecursivelyIfExists(controlDirectory)); + } + + @Test + void hdfsBulkImportRunDirectoryCleanupIsRecursiveAndIdempotent() throws IOException { + Path runDirectory = temporaryDirectory.resolve("run-1"); + Path stagedImport = runDirectory.resolve("worker-0/import-1/data.rf"); + Files.createDirectories(stagedImport.getParent()); + Files.writeString(stagedImport, "bulk import data"); + + ManagerStress.deleteBulkImportRunDirectory(temporaryDirectory.toUri().toString(), "run-1"); + + assertFalse(Files.exists(runDirectory)); + assertDoesNotThrow(() -> ManagerStress + .deleteBulkImportRunDirectory(temporaryDirectory.toUri().toString(), "run-1")); + } +} diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetricsTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetricsTest.java new file mode 100644 index 00000000..c7e743ce --- /dev/null +++ b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressMetricsTest.java @@ -0,0 +1,72 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +import java.io.IOException; +import java.nio.file.Path; +import java.util.List; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class ManagerStressMetricsTest { + + @TempDir + Path temporaryDirectory; + + @Test + void snapshotsAreReadAndAggregatedAcrossWorkers() throws IOException { + ManagerStressMetrics workerOne = new ManagerStressMetrics(); + ManagerStressMetrics.Counts create = workerOne.counts(Operation.CREATE); + create.attempts = 4; + create.submitted = 3; + create.completed = 2; + create.skipped = 1; + create.failed = 1; + create.outcomeUnknown = 1; + create.totalNanos = 40; + workerOne.addMaintenanceCreates(5); + workerOne.writeSnapshot(temporaryDirectory, 0); + + ManagerStressMetrics workerTwo = new ManagerStressMetrics(); + ManagerStressMetrics.Counts compact = workerTwo.counts(Operation.COMPACT); + compact.attempts = 2; + compact.submitted = 2; + compact.asyncAccepted = 2; + compact.totalNanos = 60; + workerTwo.addMaintenanceCreates(7); + workerTwo.writeSnapshot(temporaryDirectory, 1); + + var first = + ManagerStressMetrics.readSnapshot(ManagerStressMetrics.snapshotPath(temporaryDirectory, 0)); + var second = + ManagerStressMetrics.readSnapshot(ManagerStressMetrics.snapshotPath(temporaryDirectory, 1)); + var totals = ManagerStressMetrics.aggregate(List.of(first, second)); + + assertEquals(2, totals.workers()); + assertEquals(12, totals.maintenanceCreates()); + assertEquals(4, totals.operations().get(Operation.CREATE).attempts()); + assertEquals(3, totals.operations().get(Operation.CREATE).submitted()); + assertEquals(2, totals.operations().get(Operation.CREATE).completed()); + assertEquals(1, totals.operations().get(Operation.CREATE).outcomeUnknown()); + assertEquals(2, totals.operations().get(Operation.COMPACT).asyncAccepted()); + } +} diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptionsTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptionsTest.java new file mode 100644 index 00000000..335534b6 --- /dev/null +++ b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressOptionsTest.java @@ -0,0 +1,73 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +class ManagerStressOptionsTest { + + @org.junit.jupiter.api.Test + void requiresNamespace() { + ManagerStressOptions options = validOptions(); + options.namespace = null; + + assertThrows(IllegalArgumentException.class, () -> options.validate(true)); + } + + @org.junit.jupiter.api.Test + void validatesNamespaceCharacters() { + ManagerStressOptions options = validOptions(); + options.namespace = "invalid.namespace"; + + assertThrows(IllegalArgumentException.class, () -> options.validate(true)); + } + + @org.junit.jupiter.api.Test + void acceptsValidNamespace() { + assertDoesNotThrow(() -> validOptions().validate(true)); + } + + @org.junit.jupiter.api.Test + void requiresCompactorResourceGroup() { + ManagerStressOptions options = validOptions(); + options.compactorResourceGroup = null; + + assertThrows(IllegalArgumentException.class, () -> options.validate(true)); + } + + @org.junit.jupiter.api.Test + void configuresMgrstressCompactionPlannerForSelectedGroup() { + var properties = ManagerStress.compactionServiceProperties("stressGroup"); + + assertEquals("org.apache.accumulo.core.spi.compaction.RatioBasedCompactionPlanner", + properties.get("compaction.service.mgrstress.planner")); + assertEquals("[{\"group\":\"stressGroup\"}]", + properties.get("compaction.service.mgrstress.planner.opts.groups")); + } + + private static ManagerStressOptions validOptions() { + ManagerStressOptions options = new ManagerStressOptions(); + options.namespace = "stress_namespace_1"; + options.compactorResourceGroup = "default"; + options.hdfsDir = "/tmp/manager-stress-test"; + return options; + } +} diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressRowsTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressRowsTest.java new file mode 100644 index 00000000..4d55e1bd --- /dev/null +++ b/src/test/java/org/apache/accumulo/testing/manager/stress/ManagerStressRowsTest.java @@ -0,0 +1,119 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; +import java.util.Random; +import java.util.TreeSet; + +import org.apache.accumulo.core.client.admin.TabletAvailability; +import org.apache.accumulo.core.data.RowRange; +import org.apache.hadoop.io.Text; + +class ManagerStressRowsTest { + + @org.junit.jupiter.api.Test + void createsUnhostedTablesWithBoundedInitialSplitCounts() { + var config = ManagerStressRows.newTableConfiguration(new Random(7)); + + assertEquals(TabletAvailability.UNHOSTED, config.getInitialTabletAvailability()); + assertTrue(config.getSplits().size() >= ManagerStressRows.MIN_INITIAL_SPLITS); + assertTrue(config.getSplits().size() <= ManagerStressRows.MAX_INITIAL_SPLITS); + assertEquals(config.getSplits().size(), new TreeSet<>(config.getSplits()).size()); + config.getSplits().forEach(split -> assertTrue(split.toString().matches("[0-9a-f]{16}"))); + } + + @org.junit.jupiter.api.Test + void addedSplitsRespectTheThousandTabletLimit() { + var additions = ManagerStressRows.additionalSplits(List.of(), new Random(5)); + assertFalse(additions.isEmpty()); + assertTrue(additions.size() + 1 <= ManagerStressRows.MAX_TABLETS); + + var nearlyFull = ManagerStressRows.initialSplits(998); + var finalSplit = ManagerStressRows.additionalSplits(nearlyFull, new Random(11)); + assertEquals(1, finalSplit.size()); + assertTrue(finalSplit.stream().noneMatch(nearlyFull::contains)); + + var full = ManagerStressRows.initialSplits(ManagerStressRows.MAX_INITIAL_SPLITS); + assertTrue(ManagerStressRows.additionalSplits(full, new Random(13)).isEmpty()); + } + + @org.junit.jupiter.api.Test + void mergeRangesUseTheAdjacentSplitBoundaries() { + List splits = List.of(new Text("10"), new Text("20"), new Text("30"), new Text("40")); + + var first = ManagerStressRows.mergeRange(splits, 0, 2); + assertNull(first.start()); + assertEquals(new Text("20"), first.end()); + + var middle = ManagerStressRows.mergeRange(splits, 1, 2); + assertEquals(new Text("10"), middle.start()); + assertEquals(new Text("30"), middle.end()); + + var last = ManagerStressRows.mergeRange(splits, 3, 2); + assertEquals(new Text("30"), last.start()); + assertNull(last.end()); + + assertEquals(2, ManagerStressRows.selectMergeRange(splits, 1, new Random(17)).tabletCount()); + assertEquals(5, ManagerStressRows.selectMergeRange(splits, 100, new Random(19)).tabletCount()); + } + + @org.junit.jupiter.api.Test + void availabilityRangesMapExactlyToTabletBoundaries() { + List splits = List.of(new Text("10"), new Text("20"), new Text("30"), new Text("40")); + + assertEquals(RowRange.atMost(new Text("10")), ManagerStressRows.tabletRowRange(splits, 0, 0)); + assertEquals(RowRange.atMost(new Text("30")), ManagerStressRows.tabletRowRange(splits, 0, 2)); + assertEquals(RowRange.openClosed(new Text("10"), new Text("30")), + ManagerStressRows.tabletRowRange(splits, 1, 2)); + assertEquals(RowRange.greaterThan(new Text("30")), + ManagerStressRows.tabletRowRange(splits, 3, 4)); + assertEquals(RowRange.all(), ManagerStressRows.tabletRowRange(splits, 0, 4)); + } + + @org.junit.jupiter.api.Test + void availabilityChangesSelectUniqueTabletsAndUseDifferentStates() { + List splits = List.of(new Text("10"), new Text("20"), new Text("30"), new Text("40")); + List current = + List.of(TabletAvailability.ONDEMAND, TabletAvailability.ONDEMAND, + TabletAvailability.ONDEMAND, TabletAvailability.ONDEMAND, TabletAvailability.ONDEMAND); + + var changes = ManagerStressRows.randomAvailabilityChanges(splits, current, new Random(43)); + int selectedCount = + changes.stream().mapToInt(ManagerStressRows.AvailabilityChange::tabletCount).sum(); + + assertTrue(selectedCount >= 1); + assertTrue(selectedCount <= current.size()); + for (int i = 0; i < changes.size(); i++) { + var change = changes.get(i); + assertTrue(change.availability() != TabletAvailability.ONDEMAND); + assertEquals(ManagerStressRows.tabletRowRange(splits, change.firstTablet(), + change.firstTablet() + change.tabletCount() - 1), change.rowRange()); + if (i > 0) { + var previous = changes.get(i - 1); + assertTrue(previous.firstTablet() + previous.tabletCount() <= change.firstTablet()); + } + } + } +} diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/OperationTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/OperationTest.java new file mode 100644 index 00000000..0521a293 --- /dev/null +++ b/src/test/java/org/apache/accumulo/testing/manager/stress/OperationTest.java @@ -0,0 +1,56 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class OperationTest { + + @org.junit.jupiter.api.Test + void choosesOperationByRelativeWeight() { + ManagerStressOptions options = new ManagerStressOptions(); + options.createWeight = 0; + options.deleteWeight = 2; + options.splitWeight = 0; + options.mergeWeight = 3; + options.availabilityWeight = 0; + options.compactWeight = 0; + options.bulkImportWeight = 0; + + assertEquals(Operation.DELETE, Operation.choose(options, 0)); + assertEquals(Operation.DELETE, Operation.choose(options, 1.99)); + assertEquals(Operation.MERGE, Operation.choose(options, 2)); + assertEquals(Operation.MERGE, Operation.choose(options, 4.99)); + } + + @org.junit.jupiter.api.Test + void choosesAvailabilityByRelativeWeight() { + ManagerStressOptions options = new ManagerStressOptions(); + options.createWeight = 0; + options.deleteWeight = 0; + options.splitWeight = 0; + options.mergeWeight = 0; + options.availabilityWeight = 4; + options.compactWeight = 0; + options.bulkImportWeight = 0; + + assertEquals(Operation.AVAILABILITY, Operation.choose(options, 0)); + assertEquals(Operation.AVAILABILITY, Operation.choose(options, 3.99)); + } +} diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/TablePoolTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/TablePoolTest.java new file mode 100644 index 00000000..3144220a --- /dev/null +++ b/src/test/java/org/apache/accumulo/testing/manager/stress/TablePoolTest.java @@ -0,0 +1,33 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +import java.nio.file.Path; + +class TablePoolTest { + + @org.junit.jupiter.api.Test + void qualifiesGeneratedTableNamesWithTheNamespace() { + TablePool pool = new TablePool(Path.of("unused"), "stress_namespace", "stress", 2, 17); + + assertEquals("stress_namespace.stress_3_7", pool.tableName(3, 7)); + } +} diff --git a/src/test/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcomeTest.java b/src/test/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcomeTest.java new file mode 100644 index 00000000..dab5495c --- /dev/null +++ b/src/test/java/org/apache/accumulo/testing/manager/stress/compactor/CompactionOutcomeTest.java @@ -0,0 +1,37 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.accumulo.testing.manager.stress.compactor; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class CompactionOutcomeTest { + + @org.junit.jupiter.api.Test + void choosesConfiguredRelativeWeights() { + StressCompactorOptions options = new StressCompactorOptions(); + options.successWeight = 0; + options.failureWeight = 2; + options.cancellationWeight = 3; + + assertEquals(CompactionOutcome.FAILURE, CompactionOutcome.choose(options, 0)); + assertEquals(CompactionOutcome.FAILURE, CompactionOutcome.choose(options, 1.99)); + assertEquals(CompactionOutcome.CANCELLATION, CompactionOutcome.choose(options, 2)); + assertEquals(CompactionOutcome.CANCELLATION, CompactionOutcome.choose(options, 4.99)); + } +}