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)); + } +}