-
Notifications
You must be signed in to change notification settings - Fork 151
add Elbencho S3 benchmark with async parallel SSH #359
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
Show all changes
4 commits
Select commit
Hold shift + click to select a range
5784b8d
add a RemoteExecutor abstraction for pdsh-free parallel SSH
KenanUK aba5115
add Elbencho S3 benchmark driven by the Workloads pipeline
KenanUK 765b688
add Elbencho S3 user documentation
KenanUK 3c7c851
elbencho: adapt to monitoring package refactor
KenanUK File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,162 @@ | ||
| """Elbencho S3 benchmark, driven by the shared Workloads pipeline. | ||
|
|
||
| The generated commands are fanned out to the client nodes through the pdsh-free | ||
| ``RemoteExecutor``.""" | ||
|
|
||
| import logging | ||
|
gitkenan marked this conversation as resolved.
|
||
| import os | ||
|
|
||
| import yaml | ||
|
|
||
| from monitoring.monitoring_factory import MonitoringFactory | ||
| import settings | ||
| from remote.async_ssh import AsyncSSHExecutor | ||
| from remote.remote_executor import RemoteExecutor | ||
|
|
||
| from .benchmark import Benchmark | ||
|
|
||
| logger = logging.getLogger("cbt") | ||
|
|
||
|
|
||
| class Elbencho(Benchmark): | ||
|
|
||
| def __init__(self, archive_dir: str, cluster, config: dict) -> None: | ||
| # auth comes in from the YAML as a nested dict; flatten it to | ||
| # strings now so the Workloads pipeline (which stringifies everything) | ||
| # doesn't mangle it. | ||
| self.auth = config.get("auth", {}) | ||
| config["s3_auth_config"] = self.auth.get("config", "") | ||
| config["s3_session_token"] = self.auth.get("s3_session_token", "") | ||
| config.pop("auth", None) | ||
|
|
||
| super().__init__(archive_dir, cluster, config) | ||
|
|
||
| self.cmd_path = config.get("cmd_path", "/usr/local/bin/elbencho") | ||
|
|
||
| # RemoteExecutor is the ABC which allows us to easily | ||
| # swap out AsyncIO as the fan-out tool later if needed. | ||
| self._remote: RemoteExecutor = AsyncSSHExecutor() | ||
|
|
||
| self.base_run_dir = self.run_dir | ||
|
|
||
| workloads = config.get("workloads", {}) | ||
| if not isinstance(workloads, dict): | ||
| raise ValueError(f"workloads must be a dict, got {type(workloads).__name__}") | ||
| self._validate_workloads(workloads) | ||
|
|
||
| for wl_name, wl_params in workloads.items(): | ||
| logger.info("Elbencho workload '%s': %s", wl_name, wl_params) | ||
|
|
||
| # ------------------------------------------------------------------ | ||
| # Lifecycle overrides (pdsh-free) | ||
| # ------------------------------------------------------------------ | ||
|
|
||
| def exists(self) -> bool: | ||
| if os.path.exists(self.archive_dir): | ||
| logger.info("Skipping existing Elbencho results in %s.", self.archive_dir) | ||
| return True | ||
| return False | ||
|
|
||
| def initialize(self) -> None: | ||
| super().initialize() | ||
|
|
||
| logger.info("Verifying elbencho binary is executable on all client nodes: %s", self.cmd_path) | ||
| self._remote.run_command_with_error_checking(settings.getnodes('clients'), f"test -x {self.cmd_path}") | ||
|
|
||
| self.cleandir() | ||
|
|
||
| if not os.path.exists(self.archive_dir): | ||
| os.makedirs(self.archive_dir) | ||
|
|
||
| def cleandir(self) -> None: | ||
| clients = settings.getnodes('clients') | ||
| self._remote.clean_remote_dir(clients, self.run_dir) | ||
| self._remote.make_remote_dir(clients, self.run_dir) | ||
|
|
||
| def dropcaches(self) -> None: | ||
| nodes = settings.getnodes('clients', 'osds') | ||
| self._remote.run_command(nodes, 'sync', continue_if_error=False) | ||
| self._remote.run_command( | ||
| nodes, | ||
| 'echo 3 | sudo tee /proc/sys/vm/drop_caches', | ||
| continue_if_error=False, | ||
| ) | ||
|
|
||
| def run(self) -> None: | ||
| if self.osd_ra and self.osd_ra_changed: | ||
| logger.info('Setting OSD Read Ahead to: %s', self.osd_ra) | ||
| self.cluster.set_osd_param('read_ahead_kb', self.osd_ra) | ||
|
|
||
| config_file = os.path.join(self.archive_dir, 'benchmark_config.yaml') | ||
| if not os.path.exists(self.archive_dir): | ||
| os.makedirs(self.archive_dir) | ||
| if not os.path.exists(config_file): | ||
| config_dict = dict(cluster=self.config) | ||
| with open(config_file, 'w') as fd: | ||
| yaml.dump(config_dict, fd, default_flow_style=False) | ||
|
|
||
| if not self._workloads.exist(): | ||
| logger.warning("Elbencho: no workloads defined — nothing to run.") | ||
| return | ||
|
|
||
| self.dropcaches() | ||
| # TODO: call super().run() once Benchmark.run() is executor-driven and | ||
| # is no longer using AsyncIO. | ||
| self._remote.make_remote_dir(settings.getnodes('clients'), self.run_dir) | ||
| self.cluster.dump_config(self.run_dir) | ||
|
|
||
| self._run_workloads() | ||
|
|
||
| self._remote.sync_files(settings.getnodes('clients'), self.run_dir, self.archive_dir) | ||
|
|
||
| def cleanup(self) -> None: | ||
| pass | ||
|
|
||
| # ------------------------------------------------------------------ | ||
| # Validation | ||
| # ------------------------------------------------------------------ | ||
|
|
||
| def _validate_workloads(self, workloads: dict) -> None: | ||
| """Fail early with a precise error rather than mid-run inside the pipeline.""" | ||
| for name, params in workloads.items(): | ||
| if not isinstance(params, dict): | ||
| raise ValueError(f"workload '{name}' must be a dict") | ||
| for field in ("mode", "s3_bucket"): | ||
| if field not in params: | ||
| raise ValueError(f"workload '{name}' missing required key '{field}'") | ||
| for field in ("threads", "iodepth"): | ||
| values = params.get(field) | ||
| if values is None: | ||
| continue | ||
| for value in (values if isinstance(values, list) else [values]): | ||
| try: | ||
| int(value) | ||
| except (TypeError, ValueError): | ||
| raise ValueError( | ||
| f"workload '{name}': {field} value {value!r} is not an integer" | ||
| ) | ||
|
|
||
| # ------------------------------------------------------------------ | ||
| # Run loop | ||
| # ------------------------------------------------------------------ | ||
|
|
||
| def _run_workloads(self) -> None: | ||
| clients = settings.getnodes("clients") | ||
|
|
||
| self._workloads.set_benchmark_type("elbencho") | ||
| self._workloads.set_executable(self.cmd_path) | ||
|
|
||
| for output_directory, commands in self._workloads.command_groups(): | ||
| live_commands = [cmd for cmd in commands if cmd] | ||
| if not live_commands: | ||
| continue | ||
|
|
||
| self._remote.make_remote_dir(settings.getnodes('clients'), output_directory) | ||
| logger.info("Elbencho: running %d command(s) → %s", len(live_commands), output_directory) | ||
| MonitoringFactory.start(output_directory) | ||
| for cmd in live_commands: | ||
| logger.debug("Elbencho cmd: %s", cmd) | ||
| self._remote.run_command(clients, cmd, continue_if_error=False) | ||
| MonitoringFactory.stop() | ||
|
|
||
| logger.info("Elbencho: all workloads complete.") | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,180 @@ | ||
| """Builds the elbencho command line for a single S3 workload instance. | ||
|
|
||
| It returns the full executable string that can be used to run a cli command. | ||
| It is instantiated by ``Workload._create_command_class`` as part of the | ||
| shared Workloads pipeline.""" | ||
|
|
||
| import os | ||
| import re | ||
| import shlex | ||
| from logging import Logger, getLogger | ||
|
|
||
| from cli_options import CliOptions | ||
| from command.command import Command | ||
|
|
||
| log: Logger = getLogger("cbt") | ||
|
|
||
| _BS_SUFFIXES = {"k": 1024, "m": 1024 ** 2, "g": 1024 ** 3} | ||
|
|
||
|
|
||
| class ElbenchoCommand(Command): | ||
| """A single elbencho S3 command line for one run cell.""" | ||
|
|
||
| _MODE_FLAGS = { | ||
| "write": ["--write"], | ||
| "read": ["--read"], | ||
| "readwrite": ["--write", "--read"], | ||
| "stat": ["--stat"], | ||
| "list": ["--s3listobjpar"], | ||
| } | ||
|
|
||
| _MODES_NO_BLOCKSIZE = {"stat", "list"} | ||
|
|
||
| def __init__(self, options: dict[str, str], workload_output_directory: str) -> None: | ||
| # Must be set before super().__init__() because the base constructor calls _parse_options(). | ||
| self._workload_output_directory: str = workload_output_directory | ||
| super().__init__(options) | ||
|
|
||
| @classmethod | ||
| def mode_is_supported(cls, mode: str) -> bool: | ||
| """Return True if mode is currently supported (stat/list are not yet).""" | ||
| return mode not in cls._MODES_NO_BLOCKSIZE | ||
|
|
||
| @staticmethod | ||
| def parse_blocksize_to_bytes(blocksize: str) -> int: | ||
| s = str(blocksize).strip().lower() | ||
| m = re.fullmatch(r"(\d+(?:\.\d+)?)([kmg]?)", s) | ||
| if not m: | ||
| raise ValueError(f"Unrecognised blocksize format: {blocksize!r}") | ||
| value, suffix = m.group(1), m.group(2) | ||
| return int(float(value) * _BS_SUFFIXES.get(suffix, 1)) | ||
|
|
||
| @staticmethod | ||
| def build_auth_flags(auth: dict[str, str]) -> list[str]: | ||
| """Return S3 auth flags for whatever credentials are present; empty auth yields no flags.""" | ||
| flags: list[str] = [] | ||
|
|
||
| config_str = auth.get("config", "") | ||
| if config_str: | ||
| pairs = dict( | ||
| kv.split("=", 1) | ||
| for kv in config_str.split(";") | ||
| if "=" in kv | ||
| ) | ||
| if "url" in pairs: | ||
| flags += ["--s3endpoints", pairs["url"]] | ||
| if "access_key" in pairs: | ||
| flags += ["--s3key", pairs["access_key"]] | ||
| if "secret_key" in pairs: | ||
| flags += ["--s3secret", pairs["secret_key"]] | ||
|
|
||
| token = auth.get("s3_session_token", "") | ||
| if token: | ||
| flags += ["--s3authtoken", token] | ||
|
|
||
| return flags | ||
|
|
||
| # ------------------------------------------------------------------ | ||
| # Command ABC implementation | ||
| # ------------------------------------------------------------------ | ||
|
|
||
| def _parse_options(self, options: dict[str, str]) -> CliOptions: | ||
| # Populate CliOptions with parsed options and defaults. | ||
| parsed_options: CliOptions = CliOptions() | ||
|
|
||
| parsed_options["mode"] = options.get("mode") | ||
| parsed_options["s3_bucket"] = options.get("s3_bucket") | ||
| parsed_options["s3_region"] = options.get("s3_region", "default") | ||
| parsed_options["threads"] = str(options.get("threads", 1)) | ||
| parsed_options["blocksize"] = str(options.get("blocksize", "4k")) | ||
| parsed_options["iodepth"] = str(options.get("iodepth", 1)) | ||
|
|
||
| # Optional value flags | ||
| for key in ("size", "num_objects", "num_dirs", "duration", "hosts"): | ||
| value = options.get(key) | ||
| parsed_options[key] = str(value) if value is not None else None | ||
|
|
||
| # Boolean presence flags: "true" when truthy, None otherwise. | ||
| for key in ("deldirs", "s3nompcheck", "mkdirs"): | ||
| parsed_options[key] = "true" if options.get(key) else None | ||
|
|
||
| parsed_options["s3_auth_config"] = options.get("s3_auth_config", "") | ||
| parsed_options["s3_session_token"] = options.get("s3_session_token", "") | ||
|
|
||
| return parsed_options | ||
|
|
||
| def _parse_global_options(self, options: dict[str, str]) -> CliOptions: | ||
| return CliOptions(options) | ||
|
|
||
| def _generate_output_directory_path(self) -> str: | ||
| # {base}/elbencho/{mode}_{blocksize_bytes}/threads-{NNN}/iodepth-{NNN} | ||
| options = self._options | ||
| mode = str(options["mode"]) | ||
| blocksize = str(options["blocksize"]) | ||
| threads = int(str(options["threads"])) | ||
| iodepth = int(str(options["iodepth"])) | ||
| return os.path.join( | ||
| self._workload_output_directory, | ||
| self.benchmark, | ||
| f"{mode}_{self.parse_blocksize_to_bytes(blocksize)}", | ||
| f"threads-{threads:03d}", | ||
| f"iodepth-{iodepth:03d}", | ||
| ) | ||
|
|
||
| @property | ||
| def benchmark(self) -> str: | ||
| return "elbencho" | ||
|
|
||
| def _generate_full_command(self) -> str: | ||
| if self._executable is None: | ||
| return "" | ||
|
|
||
| options = self._options | ||
| mode = str(options["mode"]) | ||
|
|
||
| if not self.mode_is_supported(mode): | ||
| log.warning( | ||
| "Elbencho: mode '%s' is not yet supported by the formatter. " | ||
| "Skipping run (blocksize=%s, threads=%s, iodepth=%s).", | ||
| mode, options["blocksize"], options["threads"], options["iodepth"], | ||
| ) | ||
| return "" | ||
|
|
||
| mode_flags = self._MODE_FLAGS.get(mode) | ||
| if mode_flags is None: | ||
| raise ValueError(f"Unknown elbencho mode: {mode!r}") | ||
|
|
||
| cmd_parts: list[str] = [self._executable] | ||
| cmd_parts += mode_flags | ||
| cmd_parts += ["--threads", str(options["threads"])] | ||
| cmd_parts += ["--block", str(options["blocksize"])] | ||
| cmd_parts += ["--iodepth", str(options["iodepth"])] | ||
|
|
||
| if options["size"] is not None: | ||
| cmd_parts += ["--size", str(options["size"])] | ||
| if options["num_objects"] is not None: | ||
| cmd_parts += ["--files", str(options["num_objects"])] | ||
| if options["num_dirs"] is not None: | ||
| cmd_parts += ["--dirs", str(options["num_dirs"])] | ||
| if options["duration"] is not None: | ||
| cmd_parts += ["--timelimit", str(options["duration"])] | ||
| if options["deldirs"]: | ||
| cmd_parts += ["--deldirs"] | ||
| if options["s3nompcheck"]: | ||
| cmd_parts += ["--s3nompcheck"] | ||
| if options["hosts"] is not None: | ||
| cmd_parts += ["--hosts", str(options["hosts"])] | ||
|
|
||
| cmd_parts += self.build_auth_flags({ | ||
| "config": options["s3_auth_config"] or "", | ||
| "s3_session_token": options["s3_session_token"] or "", | ||
| }) | ||
| cmd_parts += ["--s3region", str(options["s3_region"])] | ||
| cmd_parts += ["--resfile", os.path.join(self._generate_output_directory_path(), "result.csv")] | ||
|
|
||
| if options["mkdirs"]: | ||
| cmd_parts += ["--mkdirs"] | ||
| cmd_parts += [f"s3://{options['s3_bucket']}"] | ||
|
|
||
| # Safely quote arguments for remote shell execution. | ||
| return shlex.join(cmd_parts) |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.