[ENH]: Allow for local queue using GNU parallel #306

Merged
synchon merged 18 commits from feature/local-queue into main 2024-03-28 09:26:23 +00:00
7 changed files with 469 additions and 6 deletions

View file

@ -0,0 +1 @@
Add support for local ``junifer queue`` via GNU Parallel by `Synchon Mandal`_

View file

@ -18,7 +18,7 @@ from ..pipeline.registry import build
from ..preprocess.base import BasePreprocessor
from ..storage.base import BaseFeatureStorage
from ..utils import logger, raise_error
from .queue_context import HTCondorAdapter
from .queue_context import GnuParallelLocalAdapter, HTCondorAdapter
from .utils import yaml
@ -220,7 +220,7 @@ def queue(
----------
config : dict
The configuration to be used for queueing the job.
kind : {"HTCondor"}
kind : {"HTCondor", "GNUParallelLocal"}
The kind of job queue system to use.
jobname : str, optional
The name of the job (default "junifer_job").
@ -239,7 +239,7 @@ def queue(
if the ``jobdir`` exists and ``overwrite = False``.
fraimondo commented 2024-03-26 09:52:04 +00:00 (Migrated from github.com)

Why are we changing this? Is this absolutely necessary? If someone updates junifer and then continues with a pre-existing yaml, it will have a weird junifer_jobs directory.

Why are we changing this? Is this absolutely necessary? If someone updates junifer and then continues with a pre-existing yaml, it will have a weird junifer_jobs directory.
synchon commented 2024-03-26 12:21:07 +00:00 (Migrated from github.com)

If an user tries out different queueing system but with the same job name (which will most likely be the case), then this is the simplest way imo to handle the assets for each system.

If an user tries out different queueing system but with the same job name (which will most likely be the case), then this is the simplest way imo to handle the assets for each system.
fraimondo commented 2024-03-26 13:31:44 +00:00 (Migrated from github.com)

From what I know, the junifer queue command does not allow to use an existing jobdir.

From what I know, the junifer queue command does not allow to use an existing `jobdir`.
synchon commented 2024-03-26 13:34:09 +00:00 (Migrated from github.com)

I don't quite understand how that can be a problem as the jobdir includes the queue kind dir.

I don't quite understand how that can be a problem as the `jobdir` includes the queue kind dir.
fraimondo commented 2024-03-26 13:36:49 +00:00 (Migrated from github.com)

Updating junifer creates an inconsistent junifer_jobs directory.

Updating junifer creates an inconsistent `junifer_jobs` directory.
synchon commented 2024-03-26 13:38:48 +00:00 (Migrated from github.com)

We can add it to the release notes and / or issue a note in the documentation. Do you have an alternative for this?

We can add it to the release notes and / or issue a note in the documentation. Do you have an alternative for this?
fraimondo commented 2024-03-26 15:17:12 +00:00 (Migrated from github.com)

My take is that we don't need the kind.lower() part, as you won't be able to run junifer queue if the jobname directory exists. It will be, technically, a different job with a different name.

I don't see any usecase to complicate ourselves like this.

My take is that we don't need the `kind.lower()` part, as you won't be able to run `junifer queue` if the `jobname` directory exists. It will be, technically, a different job with a different name. I don't see any usecase to complicate ourselves like this.
synchon commented 2024-03-26 15:32:48 +00:00 (Migrated from github.com)

So if a user wants to run a YAML using GNU Parallel (for testing) on juseless for example and then switch to HTCondor (for complete run), you say that either the user changes the job name, overwrites it or creates a new YAML?

So if a user wants to run a YAML using GNU Parallel (for testing) on juseless for example and then switch to HTCondor (for complete run), you say that either the user changes the job name, overwrites it or creates a new YAML?
fraimondo commented 2024-03-27 15:06:48 +00:00 (Migrated from github.com)

The user needs to change the YAML anyways (the queue section).

Usually the testing is done with junifer run. You don't need to test the queueing mechanism unless you are doing some advanced things. If you want to test the non-interactive run, even the GNU parallel tool is not the way. You should just run junifer queue with --element to queue only one element.

The user needs to change the YAML anyways (the `queue` section). Usually the testing is done with `junifer run`. You don't need to test the _queueing mechanism_ unless you are doing some advanced things. If you want to test the non-interactive run, even the GNU parallel tool is not the way. You should just run `junifer queue` with `--element` to queue only one element.
synchon commented 2024-03-27 15:17:05 +00:00 (Migrated from github.com)

Fair enough.

Fair enough.
"""
valid_kind = ["HTCondor"]
valid_kind = ["HTCondor", "GNUParallelLocal"]
if kind not in valid_kind:
raise_error(
f"Invalid value for `kind`: {kind}, "
@ -311,6 +311,14 @@ def queue(
elements=elements,
**kwargs, # type: ignore
)
elif kind == "GNUParallelLocal":
adapter = GnuParallelLocalAdapter(
job_name=jobname,
job_dir=jobdir,
yaml_config_path=yaml_config,
elements=elements,
**kwargs, # type: ignore
)
adapter.prepare() # type: ignore
logger.info("Queue done")

View file

@ -5,3 +5,4 @@
from .queue_context_adapter import QueueContextAdapter
from .htcondor_adapter import HTCondorAdapter
from .gnu_parallel_local_adapter import GnuParallelLocalAdapter

View file

@ -0,0 +1,258 @@
"""Define concrete class for generating GNU Parallel (local) assets."""
fraimondo commented 2024-03-26 09:56:46 +00:00 (Migrated from github.com)

While this syntax is fine for small jobs, it will create a huge command line and does not scale.

A better option is to use the --arg-file argument and give a file with the arguments to use.

While this syntax is fine for small jobs, it will create a huge command line and does not scale. A better option is to use the `--arg-file` argument and give a file with the arguments to use.
fraimondo commented 2024-03-26 09:58:53 +00:00 (Migrated from github.com)

We can directly add this to the command (see my comment above)

We can directly add this to the command (see my comment above)
fraimondo commented 2024-03-26 10:05:13 +00:00 (Migrated from github.com)

Some comments about the parallel arguments.
--bar: ok, shows a bar.
--halt: disagree, we should run as much as we can. If one subject fails, continue.
--resume: add this so the user sees the full command that can be stopped and continued (with ctrl+C)
--resume-failed: add this also, for the same reason. It will not have any effect on the first run. It will have an effect on the subsequent runs.
--joblob: perfect

Some other to consider:
--output-as-files: to redirect output to a file and keep the stdout clean
--delay: prevent having N-jobs doing the same IO operation at the beginning, which will definitely create a bottleneck and maybe a failure.

Some comments about the `parallel` arguments. `--bar`: ok, shows a bar. `--halt`: disagree, we should run as much as we can. If one subject fails, continue. `--resume`: add this so the user sees the full command that can be stopped and continued (with ctrl+C) `--resume-failed`: add this also, for the same reason. It will not have any effect on the first run. It will have an effect on the subsequent runs. `--joblob`: perfect Some other to consider: `--output-as-files`: to redirect output to a file and keep the stdout clean `--delay`: prevent having N-jobs doing the same IO operation at the beginning, which will definitely create a bottleneck and maybe a failure.
synchon commented 2024-03-26 12:35:22 +00:00 (Migrated from github.com)

Agreed. Any reason for preferring --output-as-files over --results? The latter stores the stdout, stderr and seq value in a structured way.

Agreed. Any reason for preferring `--output-as-files` over `--results`? The latter stores the stdout, stderr and seq value in a structured way.
fraimondo commented 2024-03-26 13:33:40 +00:00 (Migrated from github.com)

Did not know of the existence of --results. Pick what you think it's the best option.

Did not know of the existence of `--results`. Pick what you think it's the best option.
# Authors: Synchon Mandal <s.mandal@fz-juelich.de>
# License: AGPL
import shutil
import textwrap
from pathlib import Path
from typing import Dict, List, Optional, Tuple, Union
from ...utils import logger, make_executable, raise_error, run_ext_cmd
from .queue_context_adapter import QueueContextAdapter
__all__ = ["GnuParallelLocalAdapter"]
class GnuParallelLocalAdapter(QueueContextAdapter):
"""Class for generating commands for GNU Parallel (local).
Parameters
----------
job_name : str
The job name.
job_dir : pathlib.Path
The path to the job directory.
yaml_config_path : pathlib.Path
The path to the YAML config file.
elements : list of str or tuple
Element(s) to process. Will be used to index the DataGrabber.
pre_run : str or None, optional
Extra shell commands to source before the run (default None).
pre_collect : str or None, optional
Extra bash commands to source before the collect (default None).
env : dict, optional
The Python environment configuration. If None, will run without a
virtual environment of any kind (default None).
verbose : str, optional
The level of verbosity (default "info").
submit : bool, optional
Whether to submit the jobs (default False).
Raises
------
ValueError
If``env`` is invalid.
See Also
--------
QueueContextAdapter :
The base class for QueueContext.
HTCondorAdapter :
The concrete class for queueing via HTCondor.
"""
def __init__(
self,
job_name: str,
job_dir: Path,
yaml_config_path: Path,
elements: List[Union[str, Tuple]],
pre_run: Optional[str] = None,
pre_collect: Optional[str] = None,
env: Optional[Dict[str, str]] = None,
verbose: str = "info",
submit: bool = False,
) -> None:
"""Initialize the class."""
self._job_name = job_name
self._job_dir = job_dir
self._yaml_config_path = yaml_config_path
self._elements = elements
self._pre_run = pre_run
self._pre_collect = pre_collect
self._check_env(env)
self._verbose = verbose
self._submit = submit
self._log_dir = self._job_dir / "logs"
self._pre_run_path = self._job_dir / "pre_run.sh"
self._pre_collect_path = self._job_dir / "pre_collect.sh"
self._run_path = self._job_dir / f"run_{self._job_name}.sh"
self._collect_path = self._job_dir / f"collect_{self._job_name}.sh"
self._run_joblog_path = self._job_dir / f"run_{self._job_name}_joblog"
self._elements_file_path = self._job_dir / "elements"
def _check_env(self, env: Optional[Dict[str, str]]) -> None:
"""Check value of env parameter on init.
Parameters
----------
env : dict or None
The value of env parameter.
Raises
------
ValueError
If ``env.kind`` is invalid.
"""
# Set env related variables
if env is None:
env = {"kind": "local"}
# Check env kind
valid_env_kinds = ["conda", "venv", "local"]
if env["kind"] not in valid_env_kinds:
raise_error(
f"Invalid value for `env.kind`: {env['kind']}, "
f"must be one of {valid_env_kinds}"
)
else:
# Set variables
if env["kind"] == "local":
# No virtual environment
self._executable = "junifer"
self._arguments = ""
else:
self._executable = f"run_{env['kind']}.sh"
self._arguments = f"{env['name']} junifer"
self._exec_path = self._job_dir / self._executable
def elements(self) -> str:
"""Return elements to run."""
elements_to_run = []
for element in self._elements:
# Stringify elements if tuple for operation
str_element = (
",".join(element) if isinstance(element, tuple) else element
)
elements_to_run.append(str_element)
return "\n".join(elements_to_run)
def pre_run(self) -> str:
"""Return pre-run commands."""
fixed = (
"#!/usr/bin/env bash\n\n"
"# This script is auto-generated by junifer.\n\n"
"# Force datalad to run in non-interactive mode\n"
"DATALAD_UI_INTERACTIVE=false\n"
)
var = self._pre_run or ""
return fixed + "\n" + var
def run(self) -> str:
"""Return run commands."""
return (
"#!/usr/bin/env bash\n\n"
"# This script is auto-generated by junifer.\n\n"
"# Run pre_run.sh\n"
f"sh {self._pre_run_path.resolve()!s}\n\n"
"# Run `junifer run` using `parallel`\n"
"parallel --bar --resume --resume-failed "
f"--joblog {self._run_joblog_path} "
"--delay 60 " # wait 1 min before next job is spawned
f"--results {self._log_dir} "
f"--arg-file {self._elements_file_path.resolve()!s} "
f"{self._job_dir.resolve()!s}/{self._executable} "
f"{self._arguments} run "
f"{self._yaml_config_path.resolve()!s} "
f"--verbose {self._verbose} "
f"--element"
)
def pre_collect(self) -> str:
"""Return pre-collect commands."""
fixed = (
"#!/usr/bin/env bash\n\n"
"# This script is auto-generated by junifer.\n"
)
var = self._pre_collect or ""
return fixed + "\n" + var
def collect(self) -> str:
"""Return collect commands."""
return (
"#!/usr/bin/env bash\n\n"
"# This script is auto-generated by junifer.\n\n"
"# Run pre_collect.sh\n"
f"sh {self._pre_collect_path.resolve()!s}\n\n"
"# Run `junifer collect`\n"
f"{self._job_dir.resolve()!s}/{self._executable} "
f"{self._arguments} collect "
f"{self._yaml_config_path.resolve()!s} "
f"--verbose {self._verbose}"
)
def prepare(self) -> None:
"""Prepare assets for submission."""
logger.info("Preparing for local queue via GNU parallel")
# Copy executable if not local
if hasattr(self, "_exec_path"):
logger.info(
f"Copying {self._executable} to "
f"{self._exec_path.resolve()!s}"
)
shutil.copy(
src=Path(__file__).parent.parent / "res" / self._executable,
dst=self._exec_path,
)
make_executable(self._exec_path)
# Create elements file
logger.info(
f"Writing {self._elements_file_path.name} to "
f"{self._elements_file_path.resolve()!s}"
)
self._elements_file_path.touch()
self._elements_file_path.write_text(textwrap.dedent(self.elements()))
# Create pre run
logger.info(
f"Writing {self._pre_run_path.name} to "
f"{self._job_dir.resolve()!s}"
)
self._pre_run_path.touch()
self._pre_run_path.write_text(textwrap.dedent(self.pre_run()))
make_executable(self._pre_run_path)
# Create run
logger.info(
f"Writing {self._run_path.name} to " f"{self._job_dir.resolve()!s}"
)
self._run_path.touch()
self._run_path.write_text(textwrap.dedent(self.run()))
make_executable(self._run_path)
# Create pre collect
logger.info(
f"Writing {self._pre_collect_path.name} to "
f"{self._job_dir.resolve()!s}"
)
self._pre_collect_path.touch()
self._pre_collect_path.write_text(textwrap.dedent(self.pre_collect()))
make_executable(self._pre_collect_path)
# Create collect
logger.info(
f"Writing {self._collect_path.name} to "
f"{self._job_dir.resolve()!s}"
)
self._collect_path.touch()
self._collect_path.write_text(textwrap.dedent(self.collect()))
make_executable(self._collect_path)
# Submit if required
run_cmd = f"sh {self._run_path.resolve()!s}"
collect_cmd = f"sh {self._collect_path.resolve()!s}"
if self._submit:
logger.info(
"Shell scripts created, the following will be run:\n"
f"{run_cmd}\n"
"After successful completion of the previous step, run:\n"
f"{collect_cmd}"
)
run_ext_cmd(name=f"{self._run_path.resolve()!s}", cmd=[run_cmd])
else:
logger.info(
"Shell scripts created, to start, run:\n"
f"{run_cmd}\n"
"After successful completion of the previous step, run:\n"
f"{collect_cmd}"
)

View file

@ -66,7 +66,10 @@ class HTCondorAdapter(QueueContextAdapter):
See Also
--------
QueueContextAdapter : The base class for QueueContext.
QueueContextAdapter :
The base class for QueueContext.
GnuParallelLocalAdapter :
The concrete class for queueing via GNU Parallel (local).
"""

View file

@ -0,0 +1,192 @@
"""Provide tests for GnuParallelLocalAdapter."""
# Authors: Synchon Mandal <s.mandal@fz-juelich.de>
# License: AGPL
import logging
from pathlib import Path
from typing import Dict, List, Optional, Tuple, Union
import pytest
from junifer.api.queue_context import GnuParallelLocalAdapter
def test_GnuParallelLocalAdapter_env_error() -> None:
"""Test error for invalid env kind."""
with pytest.raises(ValueError, match="Invalid value for `env.kind`"):
GnuParallelLocalAdapter(
job_name="check_env",
job_dir=Path("."),
yaml_config_path=Path("."),
elements=["sub01"],
env={"kind": "jambalaya"},
)
@pytest.mark.parametrize(
"elements, expected_text",
[
(["sub01", "sub02"], "sub01\nsub02"),
([("sub01", "ses01"), ("sub02", "ses01")], "sub01,ses01\nsub02,ses01"),
],
)
def test_GnuParallelLocalAdapter_elements(
elements: List[Union[str, Tuple]],
expected_text: str,
) -> None:
"""Test GnuParallelLocalAdapter elements().
Parameters
----------
elements : list of str or tuple
The parametrized elements.
expected_text : str
The parametrized expected text.
"""
adapter = GnuParallelLocalAdapter(
job_name="test_elements",
job_dir=Path("."),
yaml_config_path=Path("."),
elements=elements,
)
assert expected_text in adapter.elements()
@pytest.mark.parametrize(
"pre_run, expected_text",
[
(None, "# Force datalad"),
("# Check this out\n", "# Check this out"),
],
)
def test_GnuParallelLocalAdapter_pre_run(
pre_run: Optional[str], expected_text: str
) -> None:
"""Test GnuParallelLocalAdapter pre_run().
Parameters
----------
pre_run : str or None
The parametrized pre run text.
expected_text : str
The parametrized expected text.
"""
adapter = GnuParallelLocalAdapter(
job_name="test_pre_run",
job_dir=Path("."),
yaml_config_path=Path("."),
elements=["sub01"],
pre_run=pre_run,
)
assert expected_text in adapter.pre_run()
@pytest.mark.parametrize(
"pre_collect, expected_text",
[
(None, "# This script"),
("# Check this out\n", "# Check this out"),
],
)
def test_GnuParallelLocalAdapter_pre_collect(
pre_collect: Optional[str],
expected_text: str,
) -> None:
"""Test GnuParallelLocalAdapter pre_collect().
Parameters
----------
pre_collect : str or None
The parametrized pre collect text.
expected_text : str
The parametrized expected text.
"""
adapter = GnuParallelLocalAdapter(
job_name="test_pre_collect",
job_dir=Path("."),
yaml_config_path=Path("."),
elements=["sub01"],
pre_collect=pre_collect,
)
assert expected_text in adapter.pre_collect()
def test_GnuParallelLocalAdapter_run() -> None:
"""Test HTCondorAdapter run()."""
adapter = GnuParallelLocalAdapter(
job_name="test_run",
job_dir=Path("."),
yaml_config_path=Path("."),
elements=["sub01"],
)
assert "run" in adapter.run()
def test_GnuParallelLocalAdapter_collect() -> None:
"""Test HTCondorAdapter collect()."""
adapter = GnuParallelLocalAdapter(
job_name="test_run_collect",
job_dir=Path("."),
yaml_config_path=Path("."),
elements=["sub01"],
)
assert "collect" in adapter.collect()
@pytest.mark.parametrize(
"env",
[
{"kind": "conda", "name": "junifer"},
{"kind": "venv", "name": "./junifer"},
],
)
def test_GnuParallelLocalAdapter_prepare(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
caplog: pytest.LogCaptureFixture,
env: Dict[str, str],
) -> None:
"""Test GnuParallelLocalAdapter prepare().
Parameters
----------
tmp_path : pathlib.Path
The path to the test directory.
monkeypatch : pytest.MonkeyPatch
The pytest.MonkeyPatch object.
caplog : pytest.LogCaptureFixture
The pytest.LogCaptureFixture object.
env : dict
The parametrized Python environment config.
"""
with monkeypatch.context() as m:
m.chdir(tmp_path)
with caplog.at_level(logging.DEBUG):
adapter = GnuParallelLocalAdapter(
job_name="test_prepare",
job_dir=tmp_path,
yaml_config_path=tmp_path / "config.yaml",
elements=["sub01"],
env=env,
)
adapter.prepare()
assert "GNU parallel" in caplog.text
assert f"Copying run_{env['kind']}" in caplog.text
assert "Writing pre_run.sh" in caplog.text
assert "Writing run_test_prepare.sh" in caplog.text
assert "Writing pre_collect.sh" in caplog.text
assert "Writing collect_test_prepare.sh" in caplog.text
assert "Shell scripts created" in caplog.text
assert adapter._exec_path.stat().st_size != 0
assert adapter._elements_file_path.stat().st_size != 0
assert adapter._pre_run_path.stat().st_size != 0
assert adapter._run_path.stat().st_size != 0
assert adapter._pre_collect_path.stat().st_size != 0
assert adapter._collect_path.stat().st_size != 0

View file

@ -570,14 +570,14 @@ def test_reset_queue(
},
"mem": "8G",
},
kind="HTCondor",
kind="GNUParallelLocal",
jobname=job_name,
)
# Reset operation
reset(
config={
"storage": storage,
"queue": {"jobname": job_name},
"queue": {"kind": "GNUParallelLocal", "jobname": job_name},
}
)