diff --git a/src/feasibility_data/build.py b/src/feasibility_data/build.py index 8f17cd6..2d4bbc4 100644 --- a/src/feasibility_data/build.py +++ b/src/feasibility_data/build.py @@ -14,6 +14,7 @@ BLD_REDCAP = BLD / "redcap" RAW_REDCAP = RAW / "redcap" +RAW_REDCAP_DATA = DirectoryNode(root_dir=RAW_REDCAP, pattern="*.csv.gz") RAW_MYFOOD24 = RAW / "myfood24" @@ -21,6 +22,7 @@ EVENT_METADATA_PATH = BLD_REDCAP / "event_metadata.json" REPEATING_FORMS_METADATA_PATH = BLD_REDCAP / "repeating_forms_metadata.json" FIELD_METADATA_PREPROCESSED_PATH = BLD_REDCAP / "field_metadata_preprocessed.json" +FORMS = BLD_REDCAP / "forms" def task_download_field_metadata( @@ -50,11 +52,7 @@ def task_download_repeating_forms_metadata( def task_download_raw_redcap_data( - raw_data_dir: Annotated[ - Path, - DirectoryNode(root_dir=RAW_REDCAP, pattern="*.csv.gz"), - Product, - ], + raw_data_dir: Annotated[Path, RAW_REDCAP_DATA, Product], ) -> None: """Download the latest data from all centers to `RAW_REDCAP/.csv.gz`.""" # TODO: Handle all centers @@ -77,6 +75,43 @@ def task_preprocess_field_metadata( common.json.write(field_metadata_preprocessed_path, field_metadata_preprocessed) +def task_split_forms( + forms_dir: Annotated[ + Path, + DirectoryNode(root_dir=FORMS, pattern="**/*.parquet"), + Product, + ], + raw_data_paths: Annotated[list[Path], RAW_REDCAP_DATA], + field_metadata_path: Path = FIELD_METADATA_PREPROCESSED_PATH, + event_metadata_path: Path = EVENT_METADATA_PATH, + repeating_forms_path: Path = REPEATING_FORMS_METADATA_PATH, +) -> None: + """Split each batch of raw data into one Parquet file per form. + + Written to `FORMS//.parquet`. + """ + form_to_fields = data.redcap.core.get_form_field_mapping( + common.json.read(field_metadata_path) + ) + form_to_events = data.redcap.core.get_form_event_mapping( + common.json.read(event_metadata_path) + ) + repeating_form_names = data.redcap.core.get_repeating_forms( + common.json.read(repeating_forms_path) + ) + for raw_data_path in raw_data_paths: + raw_data = data.redcap.core.read_raw(raw_data_path, form_to_fields) + forms = data.redcap.core.split_forms( + raw_data, + form_to_fields, + form_to_events, + repeating_form_names, + ) + + for form in forms: + data.redcap.core.write_form(form, forms_dir, raw_data_path) + + def task_download_myfood24_data( myfood24_raw_data_dir: Annotated[ Path, diff --git a/src/feasibility_data/data/redcap/__init__.py b/src/feasibility_data/data/redcap/__init__.py index 611980e..f6dc8f5 100644 --- a/src/feasibility_data/data/redcap/__init__.py +++ b/src/feasibility_data/data/redcap/__init__.py @@ -1,5 +1,5 @@ """REDCap data functions.""" -from . import raw +from . import core, raw -__all__ = ["raw"] +__all__ = ["core", "raw"] diff --git a/src/feasibility_data/data/redcap/core.py b/src/feasibility_data/data/redcap/core.py new file mode 100644 index 0000000..cd88188 --- /dev/null +++ b/src/feasibility_data/data/redcap/core.py @@ -0,0 +1,128 @@ +from collections import defaultdict +from dataclasses import dataclass +from operator import itemgetter +from pathlib import Path + +import polars as pl +import seedcase_soil as so + + +@dataclass +class Form: + """Class to hold the name and data of a form.""" + + name: str + data: pl.DataFrame + + +REDCAP_ID_COLS = [ + "record_id_s", + "redcap_event_name", + "redcap_repeat_instrument", + "redcap_repeat_instance", +] + + +def read_raw(raw_data_path: Path, form_to_fields: dict[str, list[str]]) -> pl.LazyFrame: + """Read the raw data into a LazyFrame with missing columns added.""" + raw_lf = pl.scan_csv(raw_data_path, infer_schema=False) + return _with_missing_columns(raw_lf, form_to_fields) + + +def get_form_field_mapping( + field_metadata: list[dict[str, str]], +) -> dict[str, list[str]]: + """Get a mapping from form name to field names in that form.""" + mapping: dict[str, list[str]] = defaultdict(list) + for field in field_metadata: + mapping[field["form_name"]].append(field["field_name"]) + + return mapping + + +def get_form_event_mapping( + event_metadata: list[dict[str, str]], +) -> dict[str, list[str]]: + """Get a mapping from form name to event names where the form is filled in.""" + mapping: dict[str, list[str]] = defaultdict(list) + for item in event_metadata: + mapping[item["form"]].append(item["unique_event_name"]) + + return mapping + + +def get_repeating_forms(repeating_forms: list[dict[str, str]]) -> set[str]: + """Get the set of repeating form names.""" + return set(so.fmap(repeating_forms, itemgetter("form_name"))) + + +def split_forms( + raw_lf: pl.LazyFrame, + form_to_fields: dict[str, list[str]], + form_to_events: dict[str, list[str]], + repeating_form_names: set[str], +) -> list[Form]: + """Split the raw data into one dataframe per form.""" + forms = so.fmap( + form_to_fields.items(), + lambda form_entry: _create_df_for_form( + form_entry, raw_lf, form_to_events, repeating_form_names + ), + ) + return so.keep(forms, lambda form: not form.data.is_empty()) + + +def write_form(form: Form, forms_dir: Path, raw_data_path: Path) -> None: + """Write the dataframe.""" + timestamp = raw_data_path.name.removesuffix(".csv.gz") + file_path = forms_dir / timestamp / f"{form.name}.parquet" + file_path.parent.mkdir(parents=True, exist_ok=True) + form.data.write_parquet(file_path) + + +def _with_missing_columns( + lf: pl.LazyFrame, form_to_fields: dict[str, list[str]] +) -> pl.LazyFrame: + """Add any missing metadata fields as columns in the dataframe.""" + metadata_fields = so.flat_fmap(form_to_fields.values(), lambda fields: fields) + data_fields = set(lf.collect_schema().names()) + missing_fields = so.keep(metadata_fields, lambda field: field not in data_fields) + return lf.with_columns(so.fmap(missing_fields, pl.lit(None).alias)) + + +def _create_df_for_form( + form_entry: tuple[str, list[str]], + raw_lf: pl.LazyFrame, + form_to_events: dict[str, list[str]], + repeating_form_names: set[str], +) -> Form: + form_name, field_names = form_entry + events = form_to_events.get(form_name, []) + is_repeating = form_name in repeating_form_names + content_fields = so.keep(field_names, lambda field: field not in REDCAP_ID_COLS) + + columns = [ + pl.col("record_id_s").alias("participant_id"), + pl.col("redcap_event_name").alias("event_id"), + *so.fmap(content_fields, pl.col), + # TODO: handle different centers + pl.lit("Copenhagen").alias("center"), + ] + + if is_repeating: + # Submissions for the same participant at the same event are told apart by + # `redcap_repeat_instance`. + columns.insert( + 2, pl.col("redcap_repeat_instance").cast(pl.String).alias("submission_id") + ) + + filters = [ + # Keep only rows coming from events where the form was filled in + pl.col("redcap_event_name").is_in(events), + # Keep only non-empty rows + pl.any_horizontal( + so.fmap(content_fields, lambda field: pl.col(field).is_not_null()) + ), + ] + + return Form(name=form_name, data=raw_lf.filter(filters).select(columns).collect()) diff --git a/uv.lock b/uv.lock index 900e025..0f050c8 100644 --- a/uv.lock +++ b/uv.lock @@ -418,7 +418,7 @@ wheels = [ [[package]] name = "cyclopts" -version = "4.22.4" +version = "4.22.5" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "attrs" }, @@ -426,9 +426,9 @@ dependencies = [ { name = "rich" }, { name = "rich-rst" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/3c/8b/fa4bfca58481ff7ef3d48ba706ccd4a7eaa1e27e7b1d9e10cbb3ae0f780f/cyclopts-4.22.4.tar.gz", hash = "sha256:d48c17e8d4a334b3f33b82920afabbf52c877a2e21edd55a172433c288bc7720", size = 194636, upload-time = "2026-08-02T14:04:17.836Z" } +sdist = { url = "https://files.pythonhosted.org/packages/be/05/689617b7e86503417c172f577d791524cb13b9697303d5d44409a971ba10/cyclopts-4.22.5.tar.gz", hash = "sha256:94044506317462cad90fb01a917dadce1f48a0915ba3605dc8d178dea1229e24", size = 195144, upload-time = "2026-08-04T13:53:00.303Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/97/49/f50d1ebb472902c952835d6141bfe920a16c402660a25c813151a861ca4c/cyclopts-4.22.4-py3-none-any.whl", hash = "sha256:90debd2468c5b33d7ca55ca64209a3bdf2667c11a2f75d973584198b58ca5e46", size = 234023, upload-time = "2026-08-02T14:04:16.282Z" }, + { url = "https://files.pythonhosted.org/packages/83/58/bcab9c33fb7a25a1f5970f357c5b19729bc81d50615d2f737b20c4255909/cyclopts-4.22.5-py3-none-any.whl", hash = "sha256:cf9ce285836053d156730ea4ea0ad0c75cf63beb3f3d8edf222a795bc57666ab", size = 234557, upload-time = "2026-08-04T13:52:58.509Z" }, ] [[package]] @@ -1515,11 +1515,11 @@ wheels = [ [[package]] name = "packaging" -version = "26.2" +version = "26.3" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/d7/f1/e7a6dd94a8d4a5626c03e4e99c87f241ba9e350cd9e6d75123f992427270/packaging-26.2.tar.gz", hash = "sha256:ff452ff5a3e828ce110190feff1178bb1f2ea2281fa2075aadb987c2fb221661", size = 228134, upload-time = "2026-04-24T20:15:23.917Z" } +sdist = { url = "https://files.pythonhosted.org/packages/7d/fa/3944b40b07da9ce895c0e6303a5ab7d53da063554f534556b134a54d6093/packaging-26.3.tar.gz", hash = "sha256:94edc256424af38762eb31306eed28beb9f0efc50a8837492c9d6fd6004aed79", size = 313412, upload-time = "2026-08-04T18:15:28.737Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/df/b2/87e62e8c3e2f4b32e5fe99e0b86d576da1312593b39f47d8ceef365e95ed/packaging-26.2-py3-none-any.whl", hash = "sha256:5fc45236b9446107ff2415ce77c807cee2862cb6fac22b8a73826d0693b0980e", size = 100195, upload-time = "2026-04-24T20:15:22.081Z" }, + { url = "https://files.pythonhosted.org/packages/63/34/ba1c580383c9eada3711951fef0795c80b829a078d72188184bcab9dd527/packaging-26.3-py3-none-any.whl", hash = "sha256:d7193f7c8e4e93f444fde0262bf90af30e16fa0ad0ad44cb553c87339b23cd1c", size = 129956, upload-time = "2026-08-04T18:15:27.159Z" }, ] [[package]]