Enh/default single output #135

Merged
fraimondo merged 6 commits from enh/default_single_output into main 2022-11-21 19:27:31 +00:00
15 changed files with 78 additions and 62 deletions

View file

@ -100,4 +100,6 @@ Bugs
- Fix a bug in which CLI command would not work using elements with more than one field (by `Fede Raimondo`_).
API changes
~~~~~~~~~~~
~~~~~~~~~~~
- Change the ``single_output`` default parameter in storage classes from ``False`` to ``True`` (:gh:`134` by `Fede Raimondo`_).

View file

@ -53,7 +53,6 @@ with tempfile.TemporaryDirectory() as tmpdir:
storage = {
"kind": "SQLiteFeatureStorage",
"uri": f"{tmpdir}/test.db",
"single_output": False,
}
# Run the defined junifer feature extraction pipeline
run(
@ -66,10 +65,7 @@ with tempfile.TemporaryDirectory() as tmpdir:
# Collect extracted features data
collect(storage=storage)
# Create storage object to read in extracted features
db = SQLiteFeatureStorage(
uri=storage["uri"],
single_output=True, # as we ran collect, we have single output now
)
db = SQLiteFeatureStorage(uri=storage["uri"])
# Read extracted features
df_vbm = db.read_df(feature_name="BOLD_Schaefer100x17_RSSETS")

View file

@ -83,7 +83,7 @@ with tempfile.TemporaryDirectory() as tmpdir:
# read in extracted features and add confounds and targets
# for julearn run cross validation
collect(storage)
db = SQLiteFeatureStorage(uri=storage["uri"], single_output=True)
db = SQLiteFeatureStorage(uri=storage["uri"])
df_vbm = db.read_df(feature_name="VBM_GM_Schaefer200x17_Mean")
oasis_subjects = [x[0] for x in df_vbm.index]

View file

@ -103,6 +103,8 @@ def run(
# Get storage engine to use
storage_params = storage.copy()
storage_kind = storage_params.pop("kind")
if "single_output" not in storage_params:
storage_params["single_output"] = False
storage_object = build(
step="storage",
name=storage_kind,
@ -137,6 +139,8 @@ def collect(storage: Dict) -> None:
storage_kind = storage_params.pop("kind")
logger.info(f"Collecting data using {storage_kind}")
logger.debug(f"\tStorage params: {storage_params}")
if "single_output" not in storage_params:
storage_params["single_output"] = False
storage_object = build(
step="storage",
name=storage_kind,
@ -155,7 +159,7 @@ def queue(
jobname: str = "junifer_job",
overwrite: bool = False,
elements: Union[str, List[Union[str, Tuple]], Tuple, None] = None,
**kwargs: Union[str, int, bool],
**kwargs: Union[str, int, bool, Dict, Tuple, List],
) -> None:
"""Queue a job to be executed later.
@ -429,7 +433,8 @@ def _queue_condor(
dag_file.write(f"JOB run{i_job} {submit_run_fname}\n")
dag_file.write(
f'VARS run{i_job} element="{str_elem} '
f'log_element="{log_elem}"\n\n')
f'log_element="{log_elem}"\n\n'
)
if collect is True:
dag_file.write(f"JOB collect {submit_collect_fname}\n")
dag_file.write("PARENT ")
@ -454,7 +459,7 @@ def _queue_slurm(
jobdir: Path,
yaml_config: Path,
elements: List[Union[str, Tuple]],
config: Dict
config: Dict,
) -> None:
"""Submit job to SLURM.

View file

@ -94,6 +94,7 @@ def test_run_multi_element(tmp_path: Path) -> None:
# Create storage
uri = outdir / "test.db"
storage["uri"] = uri # type: ignore
storage["single_output"] = False # type: ignore
# Run operations
run(
workdir=workdir,
@ -107,6 +108,39 @@ def test_run_multi_element(tmp_path: Path) -> None:
assert len(files) == 2
def test_run_multi_element_single_output(tmp_path: Path) -> None:
"""Test run function with multi element.
Parameters
----------
tmp_path : pathlib.Path
The path to the test directory.
"""
# Create working directory
workdir = tmp_path / "workdir_multi"
workdir.mkdir()
# Create output directory
outdir = tmp_path / "out"
outdir.mkdir()
# Create storage
uri = outdir / "test.db"
storage["uri"] = uri # type: ignore
storage["single_output"] = True # type: ignore
# Run operations
run(
workdir=workdir,
datagrabber=datagrabber,
markers=markers,
storage=storage,
elements=["sub-01", "sub-03"],
)
# Check files
files = list(outdir.glob("*.db"))
assert len(files) == 1
assert files[0].name == "test.db"
def test_run_and_collect(tmp_path: Path) -> None:
"""Test run and collect functions.
@ -125,6 +159,7 @@ def test_run_and_collect(tmp_path: Path) -> None:
# Create storage
uri = outdir / "test.db"
storage["uri"] = uri # type: ignore
storage["single_output"] = False # type: ignore
# Run operations
run(
workdir=workdir,

View file

@ -135,7 +135,7 @@ def test_marker_collection_storage(tmp_path: Path) -> None:
dg = OasisVBMTestingDatagrabber()
uri = tmp_path / "test_marker_collection_storage.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
mc = MarkerCollection(
markers=markers,
storage=storage,

View file

@ -79,9 +79,7 @@ def test_store(tmp_path: Path) -> None:
correlation_method="spearman",
)
uri = tmp_path / "test_crossparcellation.db"
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
out = crossparcellation.fit_transform(input_dict, storage=storage)

View file

@ -77,8 +77,6 @@ def test_store(tmp_path: Path) -> None:
ets_rss_marker = RSSETSMarker(parcellation=PARCELLATION)
# Create storage
storage = SQLiteFeatureStorage(
uri=str((tmp_path / "test.db").absolute()),
single_output=True,
)
uri=str((tmp_path / "test.db").absolute()))
# Store
ets_rss_marker.fit_transform(input=input_dict, storage=storage)

View file

@ -80,9 +80,7 @@ def test_FunctionalConnectivityParcels(tmp_path: Path) -> None:
uri = tmp_path / "test_fc_parcellation.db"
# Single storage, must be the uri
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
meta = {
"element": "test",
"version": "0.0.1",

View file

@ -64,9 +64,7 @@ def test_FunctionalConnectivitySpheres(tmp_path: Path) -> None:
uri = tmp_path / "test_fc_parcel.db"
# Single storage, must be the uri
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
meta = {
"element": "test",
"version": "0.0.1",

View file

@ -118,9 +118,7 @@ def test_SphereAggregation_storage(tmp_path: Path) -> None:
img = nib.load(vbm)
uri = tmp_path / "test_sphere_storage_3D.db"
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
meta = {
"element": "test",
"version": "0.0.1",

View file

@ -27,7 +27,7 @@ class BaseFeatureStorage(ABC):
storage_types : str or list of str
The available storage types for the class.
single_output : bool, optional
Whether to have single output (default False).
Whether to have single output (default True).
"""
@ -35,7 +35,7 @@ class BaseFeatureStorage(ABC):
self,
uri: Union[str, Path],
storage_types: Union[List[str], str],
single_output: bool = False,
single_output: bool = True,
) -> None:
self.uri = uri
if not isinstance(storage_types, list):

View file

@ -24,7 +24,7 @@ class PandasBaseFeatureStorage(BaseFeatureStorage):
uri : str or pathlib.Path
The path to the storage.
single_output : bool, optional
Whether to have single output (default False).
Whether to have single output (default True).
**kwargs
Keyword arguments passed to superclass.
@ -35,7 +35,7 @@ class PandasBaseFeatureStorage(BaseFeatureStorage):
"""
def __init__(
self, uri: Union[str, Path], single_output: bool = False, **kwargs
self, uri: Union[str, Path], single_output: bool = True, **kwargs
) -> None:
super().__init__(uri=uri, single_output=single_output, **kwargs)

View file

@ -37,7 +37,7 @@ class SQLiteFeatureStorage(PandasBaseFeatureStorage):
If True, will create only one file as specified in the `uri` and
store all the elements in the same file. This behaviour is only
suitable for non-parallel executions. SQLite does not support
concurrency (default False).
concurrency (default True).
upsert : {"ignore", "update"}, optional
Upsert mode. If "ignore" is used, the existing elements are ignored.
If "update", the existing elements are updated (default "update").
@ -53,7 +53,7 @@ class SQLiteFeatureStorage(PandasBaseFeatureStorage):
def __init__(
self,
uri: Union[str, Path],
single_output: bool = False,
single_output: bool = True,
upsert: str = "update",
**kwargs: str,
) -> None:
@ -604,14 +604,12 @@ class SQLiteFeatureStorage(PandasBaseFeatureStorage):
f"{self.uri.parent}/*{self.uri.name}" # type: ignore
)
# Create new instance
out_storage = SQLiteFeatureStorage(
uri=self.uri, single_output=True, upsert="ignore"
)
out_storage = SQLiteFeatureStorage(uri=self.uri, upsert="ignore")
# Glob files
files = self.uri.parent.glob(f"*{self.uri.name}") # type: ignore
for elem in tqdm(files, desc="file"):
logger.debug(f"Reading from {str(elem.absolute())}")
in_storage = SQLiteFeatureStorage(uri=elem, single_output=True)
in_storage = SQLiteFeatureStorage(uri=elem)
in_engine = in_storage.get_engine()
# Open "meta" table
t_meta_df = pd.read_sql(

View file

@ -94,9 +94,7 @@ def test_get_engine_single_output(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_single_output.db"
# Single storage, must be the uri
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
assert storage.single_output is True
engine = storage.get_engine()
assert engine.url.drivername == "sqlite"
@ -133,7 +131,7 @@ def test_get_engine_single_output_creation(tmp_path: Path) -> None:
# Path does not exist yet
assert not tocreate.exists()
uri = tocreate.absolute() / "test_single_output.db"
_ = SQLiteFeatureStorage(uri=uri, single_output=True, upsert="ignore")
_ = SQLiteFeatureStorage(uri=uri, upsert="ignore")
# Path exists now
assert tocreate.exists()
@ -149,9 +147,7 @@ def test_upsert_replace(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_upsert_replace.db"
# Single storage, must be the uri
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
# Metadata to store
meta = {"element": "test", "version": "0.0.1"}
# Save to database
@ -185,9 +181,7 @@ def test_upsert_ignore(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_upsert_ignore.db"
# Single storage, must be the uri
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
# Metadata to store
meta = {"element": "test", "version": "0.0.1"}
# Save to database
@ -224,7 +218,7 @@ def test_upsert_update(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_upsert_delete.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
# Metadata to store
meta = {"element": "test", "version": "0.0.1"}
# Save to database
@ -258,7 +252,7 @@ def test_upsert_invalid_option(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_upsert_invalid.db"
with pytest.raises(ValueError):
SQLiteFeatureStorage(uri=uri, single_output=True, upsert="wrong")
SQLiteFeatureStorage(uri=uri, upsert="wrong")
# TODO: can the tests be separated?
@ -272,9 +266,7 @@ def test_store_df_and_read_df(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_store_df_and_read_df.db"
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
# Metadata to store
meta = {
"element": "test",
@ -336,9 +328,7 @@ def test_store_metadata(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_metadata_store.db"
# Single storage, must be the uri
storage = SQLiteFeatureStorage(
uri=uri, single_output=True, upsert="ignore"
)
storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
# Metadata to store
meta = {"element": "test", "version": "0.0.1"}
# Store metadata
@ -356,7 +346,7 @@ def test_store_table(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_store_table.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
# Metadata to store
meta = {"element": "test", "version": "0.0.1", "marker": {"name": "fc"}}
# Data to store
@ -415,7 +405,7 @@ def test_store_matrix(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_store_table.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
# Metadata to store
meta = {"element": "test", "version": "0.0.1", "marker": {"name": "fc"}}
@ -446,7 +436,7 @@ def test_store_matrix(tmp_path: Path) -> None:
assert list(read_df.columns) == stored_names
# Store without row and column names
uri = tmp_path / "test_store_table_nonames.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix(data=data, meta=meta)
stored_names = [
f"r{i}~c{j}"
@ -478,7 +468,7 @@ def test_store_matrix(tmp_path: Path) -> None:
row_names = ["row1", "row2", "row3"]
col_names = ["col1", "col2", "col3"]
uri = tmp_path / "test_store_table_triu.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix(
data=data,
meta=meta,
@ -507,7 +497,7 @@ def test_store_matrix(tmp_path: Path) -> None:
# Store upper triangular matrix without diagonal
uri = tmp_path / "test_store_table_triu_nodiagonal.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix(
data=data,
meta=meta,
@ -537,7 +527,7 @@ def test_store_matrix(tmp_path: Path) -> None:
row_names = ["row1", "row2", "row3"]
col_names = ["col1", "col2", "col3"]
uri = tmp_path / "test_store_table_tril.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix(
data=data,
meta=meta,
@ -566,7 +556,7 @@ def test_store_matrix(tmp_path: Path) -> None:
# Store lower triangular matrix without diagonal
uri = tmp_path / "test_store_table_tril_nodiagonal.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True)
storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix(
data,
meta,
@ -702,7 +692,7 @@ def test_collect(tmp_path: Path) -> None:
"""
uri = tmp_path / "test_collect.db"
storage = SQLiteFeatureStorage(uri=uri)
storage = SQLiteFeatureStorage(uri=uri, single_output=False)
# Metadata for storage
meta1 = {
"element": {"subject": "test-01", "session": "ses-01"},