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

@ -101,3 +101,5 @@ Bugs
API changes 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 = { storage = {
"kind": "SQLiteFeatureStorage", "kind": "SQLiteFeatureStorage",
"uri": f"{tmpdir}/test.db", "uri": f"{tmpdir}/test.db",
"single_output": False,
} }
# Run the defined junifer feature extraction pipeline # Run the defined junifer feature extraction pipeline
run( run(
@ -66,10 +65,7 @@ with tempfile.TemporaryDirectory() as tmpdir:
# Collect extracted features data # Collect extracted features data
collect(storage=storage) collect(storage=storage)
# Create storage object to read in extracted features # Create storage object to read in extracted features
db = SQLiteFeatureStorage( db = SQLiteFeatureStorage(uri=storage["uri"])
uri=storage["uri"],
single_output=True, # as we ran collect, we have single output now
)
# Read extracted features # Read extracted features
df_vbm = db.read_df(feature_name="BOLD_Schaefer100x17_RSSETS") 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 # read in extracted features and add confounds and targets
# for julearn run cross validation # for julearn run cross validation
collect(storage) 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") df_vbm = db.read_df(feature_name="VBM_GM_Schaefer200x17_Mean")
oasis_subjects = [x[0] for x in df_vbm.index] oasis_subjects = [x[0] for x in df_vbm.index]

View file

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

View file

@ -94,6 +94,7 @@ def test_run_multi_element(tmp_path: Path) -> None:
# Create storage # Create storage
uri = outdir / "test.db" uri = outdir / "test.db"
storage["uri"] = uri # type: ignore storage["uri"] = uri # type: ignore
storage["single_output"] = False # type: ignore
# Run operations # Run operations
run( run(
workdir=workdir, workdir=workdir,
@ -107,6 +108,39 @@ def test_run_multi_element(tmp_path: Path) -> None:
assert len(files) == 2 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: def test_run_and_collect(tmp_path: Path) -> None:
"""Test run and collect functions. """Test run and collect functions.
@ -125,6 +159,7 @@ def test_run_and_collect(tmp_path: Path) -> None:
# Create storage # Create storage
uri = outdir / "test.db" uri = outdir / "test.db"
storage["uri"] = uri # type: ignore storage["uri"] = uri # type: ignore
storage["single_output"] = False # type: ignore
# Run operations # Run operations
run( run(
workdir=workdir, workdir=workdir,

View file

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

View file

@ -79,9 +79,7 @@ def test_store(tmp_path: Path) -> None:
correlation_method="spearman", correlation_method="spearman",
) )
uri = tmp_path / "test_crossparcellation.db" uri = tmp_path / "test_crossparcellation.db"
storage = SQLiteFeatureStorage( storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
uri=uri, single_output=True, upsert="ignore"
)
out = crossparcellation.fit_transform(input_dict, storage=storage) 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) ets_rss_marker = RSSETSMarker(parcellation=PARCELLATION)
# Create storage # Create storage
storage = SQLiteFeatureStorage( storage = SQLiteFeatureStorage(
uri=str((tmp_path / "test.db").absolute()), uri=str((tmp_path / "test.db").absolute()))
single_output=True,
)
# Store # Store
ets_rss_marker.fit_transform(input=input_dict, storage=storage) 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" uri = tmp_path / "test_fc_parcellation.db"
# Single storage, must be the uri # Single storage, must be the uri
storage = SQLiteFeatureStorage( storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
uri=uri, single_output=True, upsert="ignore"
)
meta = { meta = {
"element": "test", "element": "test",
"version": "0.0.1", "version": "0.0.1",

View file

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

View file

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

View file

@ -27,7 +27,7 @@ class BaseFeatureStorage(ABC):
storage_types : str or list of str storage_types : str or list of str
The available storage types for the class. The available storage types for the class.
single_output : bool, optional 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, self,
uri: Union[str, Path], uri: Union[str, Path],
storage_types: Union[List[str], str], storage_types: Union[List[str], str],
single_output: bool = False, single_output: bool = True,
) -> None: ) -> None:
self.uri = uri self.uri = uri
if not isinstance(storage_types, list): if not isinstance(storage_types, list):

View file

@ -24,7 +24,7 @@ class PandasBaseFeatureStorage(BaseFeatureStorage):
uri : str or pathlib.Path uri : str or pathlib.Path
The path to the storage. The path to the storage.
single_output : bool, optional single_output : bool, optional
Whether to have single output (default False). Whether to have single output (default True).
**kwargs **kwargs
Keyword arguments passed to superclass. Keyword arguments passed to superclass.
@ -35,7 +35,7 @@ class PandasBaseFeatureStorage(BaseFeatureStorage):
""" """
def __init__( def __init__(
self, uri: Union[str, Path], single_output: bool = False, **kwargs self, uri: Union[str, Path], single_output: bool = True, **kwargs
) -> None: ) -> None:
super().__init__(uri=uri, single_output=single_output, **kwargs) 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 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 store all the elements in the same file. This behaviour is only
suitable for non-parallel executions. SQLite does not support suitable for non-parallel executions. SQLite does not support
concurrency (default False). concurrency (default True).
upsert : {"ignore", "update"}, optional upsert : {"ignore", "update"}, optional
Upsert mode. If "ignore" is used, the existing elements are ignored. Upsert mode. If "ignore" is used, the existing elements are ignored.
If "update", the existing elements are updated (default "update"). If "update", the existing elements are updated (default "update").
@ -53,7 +53,7 @@ class SQLiteFeatureStorage(PandasBaseFeatureStorage):
def __init__( def __init__(
self, self,
uri: Union[str, Path], uri: Union[str, Path],
single_output: bool = False, single_output: bool = True,
upsert: str = "update", upsert: str = "update",
**kwargs: str, **kwargs: str,
) -> None: ) -> None:
@ -604,14 +604,12 @@ class SQLiteFeatureStorage(PandasBaseFeatureStorage):
f"{self.uri.parent}/*{self.uri.name}" # type: ignore f"{self.uri.parent}/*{self.uri.name}" # type: ignore
) )
# Create new instance # Create new instance
out_storage = SQLiteFeatureStorage( out_storage = SQLiteFeatureStorage(uri=self.uri, upsert="ignore")
uri=self.uri, single_output=True, upsert="ignore"
)
# Glob files # Glob files
files = self.uri.parent.glob(f"*{self.uri.name}") # type: ignore files = self.uri.parent.glob(f"*{self.uri.name}") # type: ignore
for elem in tqdm(files, desc="file"): for elem in tqdm(files, desc="file"):
logger.debug(f"Reading from {str(elem.absolute())}") 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() in_engine = in_storage.get_engine()
# Open "meta" table # Open "meta" table
t_meta_df = pd.read_sql( 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" uri = tmp_path / "test_single_output.db"
# Single storage, must be the uri # Single storage, must be the uri
storage = SQLiteFeatureStorage( storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
uri=uri, single_output=True, upsert="ignore"
)
assert storage.single_output is True assert storage.single_output is True
engine = storage.get_engine() engine = storage.get_engine()
assert engine.url.drivername == "sqlite" 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 # Path does not exist yet
assert not tocreate.exists() assert not tocreate.exists()
uri = tocreate.absolute() / "test_single_output.db" uri = tocreate.absolute() / "test_single_output.db"
_ = SQLiteFeatureStorage(uri=uri, single_output=True, upsert="ignore") _ = SQLiteFeatureStorage(uri=uri, upsert="ignore")
# Path exists now # Path exists now
assert tocreate.exists() assert tocreate.exists()
@ -149,9 +147,7 @@ def test_upsert_replace(tmp_path: Path) -> None:
""" """
uri = tmp_path / "test_upsert_replace.db" uri = tmp_path / "test_upsert_replace.db"
# Single storage, must be the uri # Single storage, must be the uri
storage = SQLiteFeatureStorage( storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
uri=uri, single_output=True, upsert="ignore"
)
# Metadata to store # Metadata to store
meta = {"element": "test", "version": "0.0.1"} meta = {"element": "test", "version": "0.0.1"}
# Save to database # Save to database
@ -185,9 +181,7 @@ def test_upsert_ignore(tmp_path: Path) -> None:
""" """
uri = tmp_path / "test_upsert_ignore.db" uri = tmp_path / "test_upsert_ignore.db"
# Single storage, must be the uri # Single storage, must be the uri
storage = SQLiteFeatureStorage( storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
uri=uri, single_output=True, upsert="ignore"
)
# Metadata to store # Metadata to store
meta = {"element": "test", "version": "0.0.1"} meta = {"element": "test", "version": "0.0.1"}
# Save to database # Save to database
@ -224,7 +218,7 @@ def test_upsert_update(tmp_path: Path) -> None:
""" """
uri = tmp_path / "test_upsert_delete.db" uri = tmp_path / "test_upsert_delete.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True) storage = SQLiteFeatureStorage(uri=uri)
# Metadata to store # Metadata to store
meta = {"element": "test", "version": "0.0.1"} meta = {"element": "test", "version": "0.0.1"}
# Save to database # Save to database
@ -258,7 +252,7 @@ def test_upsert_invalid_option(tmp_path: Path) -> None:
""" """
uri = tmp_path / "test_upsert_invalid.db" uri = tmp_path / "test_upsert_invalid.db"
with pytest.raises(ValueError): with pytest.raises(ValueError):
SQLiteFeatureStorage(uri=uri, single_output=True, upsert="wrong") SQLiteFeatureStorage(uri=uri, upsert="wrong")
# TODO: can the tests be separated? # 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" uri = tmp_path / "test_store_df_and_read_df.db"
storage = SQLiteFeatureStorage( storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
uri=uri, single_output=True, upsert="ignore"
)
# Metadata to store # Metadata to store
meta = { meta = {
"element": "test", "element": "test",
@ -336,9 +328,7 @@ def test_store_metadata(tmp_path: Path) -> None:
""" """
uri = tmp_path / "test_metadata_store.db" uri = tmp_path / "test_metadata_store.db"
# Single storage, must be the uri # Single storage, must be the uri
storage = SQLiteFeatureStorage( storage = SQLiteFeatureStorage(uri=uri, upsert="ignore")
uri=uri, single_output=True, upsert="ignore"
)
# Metadata to store # Metadata to store
meta = {"element": "test", "version": "0.0.1"} meta = {"element": "test", "version": "0.0.1"}
# Store metadata # Store metadata
@ -356,7 +346,7 @@ def test_store_table(tmp_path: Path) -> None:
""" """
uri = tmp_path / "test_store_table.db" uri = tmp_path / "test_store_table.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True) storage = SQLiteFeatureStorage(uri=uri)
# Metadata to store # Metadata to store
meta = {"element": "test", "version": "0.0.1", "marker": {"name": "fc"}} meta = {"element": "test", "version": "0.0.1", "marker": {"name": "fc"}}
# Data to store # Data to store
@ -415,7 +405,7 @@ def test_store_matrix(tmp_path: Path) -> None:
""" """
uri = tmp_path / "test_store_table.db" uri = tmp_path / "test_store_table.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True) storage = SQLiteFeatureStorage(uri=uri)
# Metadata to store # Metadata to store
meta = {"element": "test", "version": "0.0.1", "marker": {"name": "fc"}} 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 assert list(read_df.columns) == stored_names
# Store without row and column names # Store without row and column names
uri = tmp_path / "test_store_table_nonames.db" 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) storage.store_matrix(data=data, meta=meta)
stored_names = [ stored_names = [
f"r{i}~c{j}" f"r{i}~c{j}"
@ -478,7 +468,7 @@ def test_store_matrix(tmp_path: Path) -> None:
row_names = ["row1", "row2", "row3"] row_names = ["row1", "row2", "row3"]
col_names = ["col1", "col2", "col3"] col_names = ["col1", "col2", "col3"]
uri = tmp_path / "test_store_table_triu.db" uri = tmp_path / "test_store_table_triu.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True) storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix( storage.store_matrix(
data=data, data=data,
meta=meta, meta=meta,
@ -507,7 +497,7 @@ def test_store_matrix(tmp_path: Path) -> None:
# Store upper triangular matrix without diagonal # Store upper triangular matrix without diagonal
uri = tmp_path / "test_store_table_triu_nodiagonal.db" uri = tmp_path / "test_store_table_triu_nodiagonal.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True) storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix( storage.store_matrix(
data=data, data=data,
meta=meta, meta=meta,
@ -537,7 +527,7 @@ def test_store_matrix(tmp_path: Path) -> None:
row_names = ["row1", "row2", "row3"] row_names = ["row1", "row2", "row3"]
col_names = ["col1", "col2", "col3"] col_names = ["col1", "col2", "col3"]
uri = tmp_path / "test_store_table_tril.db" uri = tmp_path / "test_store_table_tril.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True) storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix( storage.store_matrix(
data=data, data=data,
meta=meta, meta=meta,
@ -566,7 +556,7 @@ def test_store_matrix(tmp_path: Path) -> None:
# Store lower triangular matrix without diagonal # Store lower triangular matrix without diagonal
uri = tmp_path / "test_store_table_tril_nodiagonal.db" uri = tmp_path / "test_store_table_tril_nodiagonal.db"
storage = SQLiteFeatureStorage(uri=uri, single_output=True) storage = SQLiteFeatureStorage(uri=uri)
storage.store_matrix( storage.store_matrix(
data, data,
meta, meta,
@ -702,7 +692,7 @@ def test_collect(tmp_path: Path) -> None:
""" """
uri = tmp_path / "test_collect.db" uri = tmp_path / "test_collect.db"
storage = SQLiteFeatureStorage(uri=uri) storage = SQLiteFeatureStorage(uri=uri, single_output=False)
# Metadata for storage # Metadata for storage
meta1 = { meta1 = {
"element": {"subject": "test-01", "session": "ses-01"}, "element": {"subject": "test-01", "session": "ses-01"},