diff --git a/.github/workflows/transpile.yml b/.github/workflows/transpile.yml index 4d75074ea..72d7d8055 100644 --- a/.github/workflows/transpile.yml +++ b/.github/workflows/transpile.yml @@ -19,6 +19,8 @@ jobs: run: pip install -e . --no-deps - name: Run transpiler unit tests run: pytest tests/test_transpile.py -q + - name: Test line duration regressions + run: pytest tests/test_line_durations.py -q generated-up-to-date: name: Check cross-dialect concepts match BQ diff --git a/mimic-iii/concepts/durations/arterial_line_durations.sql b/mimic-iii/concepts/durations/arterial_line_durations.sql index 2f87ef307..5d6798c2c 100644 --- a/mimic-iii/concepts/durations/arterial_line_durations.sql +++ b/mimic-iii/concepts/durations/arterial_line_durations.sql @@ -30,6 +30,49 @@ with mv as -- , 228286 -- Intraosseous Device | None | 12 | Processes ) ) +-- Collapse overlapping or touching MetaVision intervals to time with a line present. +-- Use the running maximum: LAG(endtime) alone fails for nested intervals. +, mv_intervals as +( + select distinct icustay_id, starttime, endtime + from mv + where arterial_line = 1 +) +, mv_ordered as +( + select icustay_id, starttime, endtime + , MAX(endtime) OVER + ( + partition by icustay_id order by starttime, endtime + rows between unbounded preceding and 1 preceding + ) as prior_endtime + from mv_intervals +) +, mv_marked as +( + select icustay_id, starttime, endtime + , case when prior_endtime is null or starttime > prior_endtime + then 1 else 0 end as new_interval + from mv_ordered +) +, mv_grouped as +( + select icustay_id, starttime, endtime + , SUM(new_interval) OVER + ( + partition by icustay_id order by starttime, endtime + rows between unbounded preceding and current row + ) as interval_group + from mv_marked +) +, mv_dur as +( + select icustay_id + , MIN(starttime) as starttime + , MAX(endtime) as endtime + from mv_grouped + group by icustay_id, interval_group +) , cv_grp as ( -- group type+site @@ -168,11 +211,9 @@ select icustay_id , starttime, endtime, duration_hours from cv_dur UNION ALL ---TODO: collapse metavision durations if they overlap select icustay_id -- , ROW_NUMBER() over (PARTITION BY icustay_id ORDER BY starttime) as arterial_line_rownum , starttime, endtime , DATETIME_DIFF(endtime, starttime, HOUR) AS duration_hours -from mv -where arterial_line = 1 +from mv_dur order by icustay_id, starttime; diff --git a/mimic-iii/concepts/durations/central_line_durations.sql b/mimic-iii/concepts/durations/central_line_durations.sql index 3f837c196..48ee53b09 100644 --- a/mimic-iii/concepts/durations/central_line_durations.sql +++ b/mimic-iii/concepts/durations/central_line_durations.sql @@ -25,6 +25,49 @@ with mv as , 224270 -- Dialysis Catheter ) ) +-- Collapse overlapping or touching MetaVision intervals to time with a line present. +-- Use the running maximum: LAG(endtime) alone fails for nested intervals. +, mv_intervals as +( + select distinct icustay_id, starttime, endtime + from mv + where central_line = 1 +) +, mv_ordered as +( + select icustay_id, starttime, endtime + , MAX(endtime) OVER + ( + partition by icustay_id order by starttime, endtime + rows between unbounded preceding and 1 preceding + ) as prior_endtime + from mv_intervals +) +, mv_marked as +( + select icustay_id, starttime, endtime + , case when prior_endtime is null or starttime > prior_endtime + then 1 else 0 end as new_interval + from mv_ordered +) +, mv_grouped as +( + select icustay_id, starttime, endtime + , SUM(new_interval) OVER + ( + partition by icustay_id order by starttime, endtime + rows between unbounded preceding and current row + ) as interval_group + from mv_marked +) +, mv_dur as +( + select icustay_id + , MIN(starttime) as starttime + , MAX(endtime) as endtime + from mv_grouped + group by icustay_id, interval_group +) , cv_grp as ( -- group type+site @@ -163,11 +206,9 @@ select icustay_id , starttime, endtime, duration_hours from cv_dur UNION ALL ---TODO: collapse metavision durations if they overlap select icustay_id -- , ROW_NUMBER() over (PARTITION BY icustay_id ORDER BY starttime) as central_line_rownum , starttime, endtime , DATETIME_DIFF(endtime, starttime, HOUR) AS duration_hours -from mv -where central_line = 1 +from mv_dur order by icustay_id, starttime; diff --git a/mimic-iii/concepts_duckdb/durations/arterial_line_durations.sql b/mimic-iii/concepts_duckdb/durations/arterial_line_durations.sql index 5958efc71..5903a06fd 100644 --- a/mimic-iii/concepts_duckdb/durations/arterial_line_durations.sql +++ b/mimic-iii/concepts_duckdb/durations/arterial_line_durations.sql @@ -17,6 +17,52 @@ WITH mv AS ( FROM mimiciii.procedureevents_mv AS pe WHERE pe.itemid IN (224263, 224267, 224268, 225199, 225752, 225789, 224272) +), mv_intervals AS ( + SELECT DISTINCT + icustay_id, + starttime, + endtime + FROM mv + WHERE + arterial_line = 1 +), mv_ordered AS ( + SELECT + icustay_id, + starttime, + endtime, + MAX(endtime) OVER ( + PARTITION BY icustay_id + ORDER BY starttime NULLS FIRST, endtime NULLS FIRST + rows BETWEEN UNBOUNDED preceding AND 1 preceding + ) AS prior_endtime + FROM mv_intervals +), mv_marked AS ( + SELECT + icustay_id, + starttime, + endtime, + CASE WHEN prior_endtime IS NULL OR starttime > prior_endtime THEN 1 ELSE 0 END AS new_interval + FROM mv_ordered +), mv_grouped AS ( + SELECT + icustay_id, + starttime, + endtime, + SUM(new_interval) OVER ( + PARTITION BY icustay_id + ORDER BY starttime NULLS FIRST, endtime NULLS FIRST + rows BETWEEN UNBOUNDED preceding AND CURRENT ROW + ) AS interval_group + FROM mv_marked +), mv_dur AS ( + SELECT + icustay_id, + MIN(starttime) AS starttime, + MAX(endtime) AS endtime + FROM mv_grouped + GROUP BY + icustay_id, + interval_group ), cv_grp AS ( SELECT ce.icustay_id, @@ -136,9 +182,7 @@ SELECT starttime, endtime, DATE_DIFF('HOUR', starttime, endtime) AS duration_hours -FROM mv -WHERE - arterial_line = 1 +FROM mv_dur ORDER BY icustay_id NULLS FIRST, starttime NULLS FIRST \ No newline at end of file diff --git a/mimic-iii/concepts_duckdb/durations/central_line_durations.sql b/mimic-iii/concepts_duckdb/durations/central_line_durations.sql index 33da694ac..8f6297acb 100644 --- a/mimic-iii/concepts_duckdb/durations/central_line_durations.sql +++ b/mimic-iii/concepts_duckdb/durations/central_line_durations.sql @@ -27,6 +27,52 @@ WITH mv AS ( 227719, 224270 ) +), mv_intervals AS ( + SELECT DISTINCT + icustay_id, + starttime, + endtime + FROM mv + WHERE + central_line = 1 +), mv_ordered AS ( + SELECT + icustay_id, + starttime, + endtime, + MAX(endtime) OVER ( + PARTITION BY icustay_id + ORDER BY starttime NULLS FIRST, endtime NULLS FIRST + rows BETWEEN UNBOUNDED preceding AND 1 preceding + ) AS prior_endtime + FROM mv_intervals +), mv_marked AS ( + SELECT + icustay_id, + starttime, + endtime, + CASE WHEN prior_endtime IS NULL OR starttime > prior_endtime THEN 1 ELSE 0 END AS new_interval + FROM mv_ordered +), mv_grouped AS ( + SELECT + icustay_id, + starttime, + endtime, + SUM(new_interval) OVER ( + PARTITION BY icustay_id + ORDER BY starttime NULLS FIRST, endtime NULLS FIRST + rows BETWEEN UNBOUNDED preceding AND CURRENT ROW + ) AS interval_group + FROM mv_marked +), mv_dur AS ( + SELECT + icustay_id, + MIN(starttime) AS starttime, + MAX(endtime) AS endtime + FROM mv_grouped + GROUP BY + icustay_id, + interval_group ), cv_grp AS ( SELECT ce.icustay_id, @@ -234,9 +280,7 @@ SELECT starttime, endtime, DATE_DIFF('HOUR', starttime, endtime) AS duration_hours -FROM mv -WHERE - central_line = 1 +FROM mv_dur ORDER BY icustay_id NULLS FIRST, starttime NULLS FIRST \ No newline at end of file diff --git a/mimic-iii/concepts_postgres/durations/arterial_line_durations.sql b/mimic-iii/concepts_postgres/durations/arterial_line_durations.sql index cc40408b6..ed03abb0b 100644 --- a/mimic-iii/concepts_postgres/durations/arterial_line_durations.sql +++ b/mimic-iii/concepts_postgres/durations/arterial_line_durations.sql @@ -25,6 +25,52 @@ WITH mv AS ( 225789, /* Sheath */ 224272 /* IABP Line */ ) /* , 227719 -- AVA Line | None | 12 | Processes */ /* , 228286 -- Intraosseous Device | None | 12 | Processes */ +), mv_intervals /* Collapse overlapping or touching MetaVision intervals to time with a line present. */ /* Use the running maximum: LAG(endtime) alone fails for nested intervals. */ AS ( + SELECT DISTINCT + icustay_id, + starttime, + endtime + FROM mv + WHERE + arterial_line = 1 +), mv_ordered AS ( + SELECT + icustay_id, + starttime, + endtime, + MAX(endtime) OVER ( + PARTITION BY icustay_id + ORDER BY starttime NULLS FIRST, endtime NULLS FIRST + rows BETWEEN UNBOUNDED preceding AND 1 preceding + ) AS prior_endtime + FROM mv_intervals +), mv_marked AS ( + SELECT + icustay_id, + starttime, + endtime, + CASE WHEN prior_endtime IS NULL OR starttime > prior_endtime THEN 1 ELSE 0 END AS new_interval + FROM mv_ordered +), mv_grouped AS ( + SELECT + icustay_id, + starttime, + endtime, + SUM(new_interval) OVER ( + PARTITION BY icustay_id + ORDER BY starttime NULLS FIRST, endtime NULLS FIRST + rows BETWEEN UNBOUNDED preceding AND CURRENT ROW + ) AS interval_group + FROM mv_marked +), mv_dur AS ( + SELECT + icustay_id, + MIN(starttime) AS starttime, + MAX(endtime) AS endtime + FROM mv_grouped + GROUP BY + icustay_id, + interval_group ), cv_grp AS ( /* group type+site */ SELECT @@ -144,15 +190,12 @@ SELECT duration_hours FROM cv_dur UNION ALL -/* TODO: collapse metavision durations if they overlap */ SELECT icustay_id, /* , ROW_NUMBER() over (PARTITION BY icustay_id ORDER BY starttime) as arterial_line_rownum */ starttime, endtime, CAST(EXTRACT(EPOCH FROM DATE_TRUNC('hour', endtime) - DATE_TRUNC('hour', starttime)) / 3600 AS BIGINT) AS duration_hours -FROM mv -WHERE - arterial_line = 1 +FROM mv_dur ORDER BY icustay_id NULLS FIRST, starttime NULLS FIRST \ No newline at end of file diff --git a/mimic-iii/concepts_postgres/durations/central_line_durations.sql b/mimic-iii/concepts_postgres/durations/central_line_durations.sql index 905e17a0f..f207ab380 100644 --- a/mimic-iii/concepts_postgres/durations/central_line_durations.sql +++ b/mimic-iii/concepts_postgres/durations/central_line_durations.sql @@ -27,6 +27,52 @@ WITH mv AS ( 227719, /* AVA Line | None | 12 | Processes */ /* , 228286 -- Intraosseous Device | None | 12 | Processes */ 224270 /* Dialysis Catheter */ ) +), mv_intervals /* Collapse overlapping or touching MetaVision intervals to time with a line present. */ /* Use the running maximum: LAG(endtime) alone fails for nested intervals. */ AS ( + SELECT DISTINCT + icustay_id, + starttime, + endtime + FROM mv + WHERE + central_line = 1 +), mv_ordered AS ( + SELECT + icustay_id, + starttime, + endtime, + MAX(endtime) OVER ( + PARTITION BY icustay_id + ORDER BY starttime NULLS FIRST, endtime NULLS FIRST + rows BETWEEN UNBOUNDED preceding AND 1 preceding + ) AS prior_endtime + FROM mv_intervals +), mv_marked AS ( + SELECT + icustay_id, + starttime, + endtime, + CASE WHEN prior_endtime IS NULL OR starttime > prior_endtime THEN 1 ELSE 0 END AS new_interval + FROM mv_ordered +), mv_grouped AS ( + SELECT + icustay_id, + starttime, + endtime, + SUM(new_interval) OVER ( + PARTITION BY icustay_id + ORDER BY starttime NULLS FIRST, endtime NULLS FIRST + rows BETWEEN UNBOUNDED preceding AND CURRENT ROW + ) AS interval_group + FROM mv_marked +), mv_dur AS ( + SELECT + icustay_id, + MIN(starttime) AS starttime, + MAX(endtime) AS endtime + FROM mv_grouped + GROUP BY + icustay_id, + interval_group ), cv_grp AS ( /* group type+site */ SELECT @@ -234,15 +280,12 @@ SELECT duration_hours FROM cv_dur UNION ALL -/* TODO: collapse metavision durations if they overlap */ SELECT icustay_id, /* , ROW_NUMBER() over (PARTITION BY icustay_id ORDER BY starttime) as central_line_rownum */ starttime, endtime, CAST(EXTRACT(EPOCH FROM DATE_TRUNC('hour', endtime) - DATE_TRUNC('hour', starttime)) / 3600 AS BIGINT) AS duration_hours -FROM mv -WHERE - central_line = 1 +FROM mv_dur ORDER BY icustay_id NULLS FIRST, starttime NULLS FIRST \ No newline at end of file diff --git a/tests/test_line_durations.py b/tests/test_line_durations.py new file mode 100644 index 000000000..cb4223c9d --- /dev/null +++ b/tests/test_line_durations.py @@ -0,0 +1,64 @@ +"""Regression coverage for overlapping MetaVision line intervals (#371).""" +from datetime import datetime, timedelta +from pathlib import Path + +import pytest + +from mimic_utils.transpile import transpile_query + +duckdb = pytest.importorskip("duckdb") +ROOT = Path(__file__).resolve().parents[1] / "mimic-iii" + + +@pytest.mark.parametrize("line_type", ["arterial", "central"]) +@pytest.mark.parametrize("sql_source", ["bigquery", "committed_duckdb"]) +def test_line_intervals_are_unioned_without_bridging_gaps(line_type, sql_source): + origin = datetime(2150, 1, 1) + itemid = 225752 if line_type == "arterial" else 224263 + location = "Invasive Arterial" if line_type == "arterial" else "Invasive Venous" + with duckdb.connect() as con: + con.execute("CREATE SCHEMA mimiciii; CREATE SCHEMA mimiciii_derived") + con.execute(""" + CREATE TABLE mimiciii.procedureevents_mv ( + icustay_id BIGINT, starttime TIMESTAMP, endtime TIMESTAMP, + itemid INTEGER, locationcategory VARCHAR + ); + CREATE TABLE mimiciii.chartevents ( + icustay_id BIGINT, charttime TIMESTAMP, itemid INTEGER, value VARCHAR + ) + """) + # Nested rows must not break the enclosing interval. Duplicates, + # overlap chains, and touching intervals belong to the same island. + # The 14-to-16 gap and another patient's interval must stay separate. + for stay, start, end in [ + (1, 0, 10), (1, 0, 10), (1, 0, 1), (1, 2, 3), (1, 4, 5), (1, 4, 5), + (1, 9, 12), (1, 12, 14), (1, 16, 18), (2, 1, 6), + ]: + con.execute("INSERT INTO mimiciii.procedureevents_mv VALUES (?, ?, ?, ?, ?)", + [stay, origin + timedelta(hours=start), origin + timedelta(hours=end), itemid, location]) + # An unrelated procedure cannot bridge the gap. + con.execute("INSERT INTO mimiciii.procedureevents_mv VALUES (?, ?, ?, ?, ?)", + [1, origin + timedelta(hours=14), origin + timedelta(hours=16), 999999, location]) + # Preserve the existing CareVue extraction, including a documentation gap. + for hour in (0, 2, 20, 22): + con.execute("INSERT INTO mimiciii.chartevents VALUES (?, ?, 229, ?)", + [3, origin + timedelta(hours=hour), "A-Line" if line_type == "arterial" else "Multi-lumen"]) + + name = f"{line_type}_line_durations" + if sql_source == "bigquery": + query = (ROOT / "concepts" / "durations" / f"{name}.sql").read_text() + con.execute( + f"CREATE TABLE mimiciii_derived.{name} AS " + + transpile_query(query, "bigquery", "duckdb", {"mimiciii_clinical": "mimiciii"}) + ) + else: + con.execute((ROOT / "concepts_duckdb" / "durations" / f"{name}.sql").read_text()) + actual = con.execute(f""" + SELECT icustay_id, starttime, endtime, duration_hours + FROM mimiciii_derived.{name} ORDER BY icustay_id, starttime + """).fetchall() + expected = [ + (stay, origin + timedelta(hours=start), origin + timedelta(hours=end), end - start) + for stay, start, end in [(1, 0, 14), (1, 16, 18), (2, 1, 6), (3, 0, 2), (3, 20, 22)] + ] + assert actual == expected