mmcdermott-meds-transforms-374
`reduce_over` can begin reading a mapper output as soon as its path exists, even when the parquet content is still being written. This can cause a parquet-format error instead of waiting for the input to become readable and then producing the reduced result.
The reducer must wait until every input path is a file containing a complete, readable parquet. It must re-check inputs every `polling_time` seconds and must not sleep past a bounded deadline. The `reduce_over` interface in `MEDS_transforms.mapreduce.reducer` gains the keyword argument `max_poll_time: float | None = None`; `None` preserves indefinite waiting. If `max_poll_time` is provided and is not strictly greater than `polling_time`, it must raise `ValueError` immediately with a message containing `"must be greater than"`. When the deadline expires, it must raise `TimeoutError`. Its message must include `"present but unreadable: <paths>"` for still-existing inputs that never became valid parquet and `"missing: <paths>"` for inputs that never appeared, including each group only when non-empty. A missing path must not be reported as unreadable, and an existing invalid or incomplete path must not be reported as missing.
`write_df` in `MEDS_transforms.dataframe.write_fn` must publish parquet output atomically, so a reader never sees a partially written final file. The tests patch `MEDS_transforms.dataframe.write_fn.os.replace` with an `OSError` side effect and expect that to be what makes the publish step fail, and afterwards they check that no staging file such as `<name>.parquet.tmp` remains next to the output. If parquet writing or publishing fails, the staging path must be removed, ignoring an absent path, and the original exception must be re-raised; a failed operation must leave neither a staging file nor a newly published final file. `write_df` must continue to create missing parent directories and support both eager and lazy Polars data frames.
Hidden tests · 6 fail-to-pass, 3 pass-to-passrun after the agent submits, in a clean verifier
Test patch · 239 lines
diff --git a/tests/test_reducer_race.py b/tests/test_reducer_race.py
new file mode 100644
index 00000000..1f74a7ea
--- /dev/null
+++ b/tests/test_reducer_race.py
@@ -0,0 +1,149 @@
+"""Regression test for https://github.com/mmcdermott/MEDS_transforms/issues/373.
+
+``reduce_over`` used to poll for input readiness with ``fp.is_file()``, which returns True as soon
+as the mapper creates the output file — well before ``df.write_parquet(...)`` has flushed a valid
+parquet footer. The reducer then read a partial file and raised ``polars.exceptions.ComputeError:
+parquet: File out of specification``. These tests lock in the fix (atomic ``write_df`` publish plus
+a completeness-aware readiness poll).
+"""
+
+from __future__ import annotations
+
+import threading
+import time
+from typing import TYPE_CHECKING
+
+import polars as pl
+import pytest
+from polars.testing import assert_frame_equal
+
+from MEDS_transforms.dataframe import read_df, write_df
+from MEDS_transforms.mapreduce.reducer import reduce_over
+
+if TYPE_CHECKING:
+ from pathlib import Path
+
+# Window during which the reducer must observe the empty-but-existing parquet file. On a machine
+# where this isn't enough, the whole test suite is already in trouble.
+_RACE_WINDOW_SECONDS = 0.3
+
+
+def _publish_before_complete_write(
+ df: pl.DataFrame,
+ fp: Path,
+ touched: threading.Event,
+ race_window: float,
+) -> None:
+ """Simulate a publish-before-complete-write interleaving.
+
+ Creates ``fp`` (so ``is_file()`` returns True), signals ``touched``, then holds
+ ``race_window`` seconds before writing real parquet content. This is a synthetic interleaving,
+ not a literal model of ``df.write_parquet``, but it exposes the same contract violation: the
+ reducer treats file existence as publication.
+ """
+ fp.touch()
+ touched.set()
+ time.sleep(race_window)
+ write_df(df, fp)
+
+
+def _reduce_fn(*dfs: pl.LazyFrame | pl.DataFrame) -> pl.LazyFrame | pl.DataFrame:
+ return pl.concat(dfs, how="vertical")
+
+
+def test_reduce_over_waits_for_complete_parquet(tmp_path: Path) -> None:
+ """Reducer should wait for valid parquet, not just file existence."""
+ in_fps = [tmp_path / f"in_{i}.parquet" for i in range(2)]
+ out_fp = tmp_path / "out.parquet"
+
+ df0 = pl.DataFrame({"a": [1, 2], "b": [3, 4]})
+ df1 = pl.DataFrame({"a": [5, 6], "b": [7, 8]})
+
+ write_df(df0, in_fps[0])
+
+ touched = threading.Event()
+ slow_writer = threading.Thread(
+ target=_publish_before_complete_write,
+ args=(df1, in_fps[1], touched, _RACE_WINDOW_SECONDS),
+ )
+ slow_writer.start()
+ try:
+ # Explicit handshake: proceed only once the partial file exists.
+ assert touched.wait(timeout=5.0), "writer thread never created the partial file"
+ reduce_over(
+ in_fps=in_fps,
+ out_fp=out_fp,
+ read_fn=read_df,
+ write_fn=write_df,
+ reduce_fn=_reduce_fn,
+ polling_time=0.005,
+ )
+ finally:
+ slow_writer.join(timeout=5.0)
+
+ # ``read_df`` returns a LazyFrame; collect before comparing.
+ result = read_df(out_fp).collect().sort("a")
+ expected = pl.concat([df0, df1], how="vertical").sort("a")
+ # check_dtypes=False because ``reduce_over`` calls ``shrink_dtype`` on numeric columns.
+ assert_frame_equal(result, expected, check_dtypes=False)
+
+
+def test_reduce_over_times_out_on_permanently_invalid_input(tmp_path: Path) -> None:
+ """A permanently invalid (empty) parquet should time out, not hang forever."""
+ in_fps = [tmp_path / "ok.parquet", tmp_path / "broken.parquet"]
+ out_fp = tmp_path / "out.parquet"
+
+ write_df(pl.DataFrame({"a": [1]}), in_fps[0])
+ in_fps[1].touch() # exists but invalid parquet, never fixed
+
+ with pytest.raises(TimeoutError, match=r"present but unreadable.*broken\.parquet") as excinfo:
+ reduce_over(
+ in_fps=in_fps,
+ out_fp=out_fp,
+ read_fn=read_df,
+ write_fn=write_df,
+ reduce_fn=_reduce_fn,
+ polling_time=0.01,
+ max_poll_time=0.3,
+ )
+ # The "missing" branch should NOT fire here — the file exists (just isn't valid parquet).
+ assert "missing:" not in str(excinfo.value)
+
+
+def test_reduce_over_times_out_on_missing_input(tmp_path: Path) -> None:
+ """An input path that never appears should time out with a ``missing:`` message."""
+ in_fps = [tmp_path / "ok.parquet", tmp_path / "never_created.parquet"]
+ out_fp = tmp_path / "out.parquet"
+
+ write_df(pl.DataFrame({"a": [1]}), in_fps[0])
+ # Deliberately do not create in_fps[1].
+
+ with pytest.raises(TimeoutError, match=r"missing:.*never_created\.parquet") as excinfo:
+ reduce_over(
+ in_fps=in_fps,
+ out_fp=out_fp,
+ read_fn=read_df,
+ write_fn=write_df,
+ reduce_fn=_reduce_fn,
+ polling_time=0.01,
+ max_poll_time=0.3,
+ )
+ # The "present but unreadable" branch should NOT fire here — the file never existed.
+ assert "present but unreadable" not in str(excinfo.value)
+
+
+def test_reduce_over_rejects_max_poll_time_not_larger_than_polling_time(tmp_path: Path) -> None:
+ """``max_poll_time`` must exceed ``polling_time`` to avoid spurious timeouts on the first poll."""
+ in_fps = [tmp_path / "a.parquet"]
+ out_fp = tmp_path / "out.parquet"
+
+ with pytest.raises(ValueError, match="must be greater than"):
+ reduce_over(
+ in_fps=in_fps,
+ out_fp=out_fp,
+ read_fn=read_df,
+ write_fn=write_df,
+ reduce_fn=_reduce_fn,
+ polling_time=1.0,
+ max_poll_time=1.0,
+ )
diff --git a/tests/test_write_df_atomic.py b/tests/test_write_df_atomic.py
new file mode 100644
index 00000000..a1ce7379
--- /dev/null
+++ b/tests/test_write_df_atomic.py
@@ -0,0 +1,78 @@
+"""Tests for the atomic-publish behavior of :func:`MEDS_transforms.dataframe.write_df`."""
+
+from __future__ import annotations
+
+from typing import TYPE_CHECKING
+from unittest.mock import patch
+
+import polars as pl
+import pytest
+
+from MEDS_transforms.dataframe import write_df
+
+if TYPE_CHECKING:
+ from pathlib import Path
+
+
+def test_write_df_publishes_final_file(tmp_path: Path) -> None:
+ """Happy path: the final file exists, the staging .tmp file does not."""
+ out_fp = tmp_path / "out.parquet"
+ write_df(pl.DataFrame({"a": [1, 2, 3]}), out_fp)
+
+ assert out_fp.is_file()
+ assert not out_fp.with_suffix(".parquet.tmp").exists()
+
+
+def test_write_df_creates_parent_directory(tmp_path: Path) -> None:
+ """Missing parent dirs should be created by write_df (covers the mkdir branch)."""
+ out_fp = tmp_path / "nested" / "deeper" / "out.parquet"
+ assert not out_fp.parent.exists()
+ write_df(pl.DataFrame({"a": [1]}), out_fp)
+ assert out_fp.is_file()
+
+
+def test_write_df_collects_lazyframe(tmp_path: Path) -> None:
+ """A LazyFrame should be collected and written as parquet."""
+ out_fp = tmp_path / "out.parquet"
+ write_df(pl.LazyFrame({"a": [1, 2, 3]}), out_fp)
+ assert out_fp.is_file()
+ assert pl.read_parquet(out_fp)["a"].to_list() == [1, 2, 3]
+
+
+def test_write_df_cleans_up_tmp_file_on_rename_failure(tmp_path: Path) -> None:
+ """If ``os.replace`` raises after the tmp parquet is written, the tmp file is unlinked and the original
+ exception propagates.
+
+ The final ``out_fp`` should remain untouched.
+ """
+ out_fp = tmp_path / "out.parquet"
+ tmp_fp = out_fp.with_suffix(".parquet.tmp")
+
+ with (
+ patch("MEDS_transforms.dataframe.write_fn.os.replace", side_effect=OSError("boom")),
+ pytest.raises(OSError, match="boom"),
+ ):
+ write_df(pl.DataFrame({"a": [1]}), out_fp)
+
+ assert not tmp_fp.exists(), "staging .tmp file should be cleaned up on failure"
+ as
… [863 more characters]Reference fix · 2 files, +61 −4the upstream merge, used only for grading calibration
The agent could not see this: the repository holds one commit and the sandbox has no network. Leak audit.
src/MEDS_transforms/dataframe/write_fn.py, src/MEDS_transforms/mapreduce/reducer.py
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py
index 76cd21e1..8229fb59 100644
--- a/src/MEDS_transforms/dataframe/write_fn.py
+++ b/src/MEDS_transforms/dataframe/write_fn.py
@@ -1,3 +1,4 @@
+import os
from collections.abc import Callable
from pathlib import Path
@@ -9,8 +10,20 @@
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """Atomically write a dataframe, either lazy or eager, to a parquet file.
+
+ Content is staged at ``<out_fp>.tmp`` and then moved into place with
+ ``os.replace`` so concurrent readers never observe a partial parquet file at
+ the final path. The rename is atomic on POSIX and Windows as long as
+ ``out_fp`` and its staging path share a filesystem.
+ """
if isinstance(df, pl.LazyFrame):
df = df.collect()
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ tmp_fp = out_fp.with_suffix(out_fp.suffix + ".tmp")
+ try:
+ df.write_parquet(tmp_fp, use_pyarrow=True)
+ os.replace(tmp_fp, out_fp)
+ except BaseException:
+ tmp_fp.unlink(missing_ok=True)
+ raise
diff --git a/src/MEDS_transforms/mapreduce/reducer.py b/src/MEDS_transforms/mapreduce/reducer.py
index 3b29cebd..0d13aa83 100644
--- a/src/MEDS_transforms/mapreduce/reducer.py
+++ b/src/MEDS_transforms/mapreduce/reducer.py
@@ -8,6 +8,7 @@
import polars as pl
from ..dataframe import DF_T, READ_FN_T, WRITE_FN_T
+from .rwlock import default_file_checker
logger = logging.getLogger(__name__)
@@ -28,6 +29,7 @@ def reduce_over(
merge_fn: REDUCE_FN_T | None = None,
do_overwrite: bool = False,
polling_time: float = 0.1,
+ max_poll_time: float | None = None,
):
"""Performs a reduction operation on a list of input file paths, with optional merging to existing data.
@@ -35,6 +37,11 @@ def reduce_over(
in_fps: List of input file paths containing data over which the reduction should be performed.
out_fp: Output file path where the reduced data will be saved.
polling_time: Time in seconds to wait between checks for file readiness.
+ max_poll_time: Optional maximum total seconds to wait for input files to become readable.
+ Defaults to ``None`` (wait indefinitely). When set, ``TimeoutError`` is raised if any
+ input file exists but stays invalid past the deadline — guards against hangs when a
+ mapper is permanently stuck or produced a corrupt file. Callers should pick a value
+ meaningfully larger than ``polling_time``.
read_fn: Function to read data from the input file paths.
write_fn: Function to write the reduced data to the output file path.
reduce_fn: Function to perform the reduction operation on the data. It should take two dataframe
@@ -45,6 +52,7 @@ def reduce_over(
Raises:
FileExistsError: If the output file already exists.
+ TimeoutError: If input files remain unreadable after ``max_poll_time`` seconds (only when set).
Examples:
>>> def reduce_fn(*dfs: pl.DataFrame) -> pl.DataFrame:
@@ -203,9 +211,45 @@ def reduce_over(
if out_fp.is_file() and not do_overwrite:
raise FileExistsError(f"Output file already exists: {out_fp.resolve()!s}")
- while not all(fp.is_file() for fp in in_fps):
+ if max_poll_time is not None and max_poll_time <= polling_time:
+ raise ValueError(
+ f"max_poll_time ({max_poll_time}s) must be greater than polling_time ({polling_time}s); "
+ "otherwise legitimately slow mappers will trigger a TimeoutError on the first recheck."
+ )
+
+ deadline = None if max_poll_time is None else time.monotonic() + max_poll_time
+ ready: set[Path] = set()
+
+ def _is_ready(fp: Path) -> bool:
+ # ``is_file`` is a cheap stat; only if it passes do we open the parquet to check completeness.
+ return fp.is_file() and default_file_checker(fp)
+
+ while True:
+ ready.update(fp for fp in in_fps if fp not in ready and _is_ready(fp))
+ if len(ready) == len(in_fps):
+ break
+
+ if deadline is None:
+ sleep_for = polling_time
+ else:
+ remaining = deadline - time.monotonic()
+ if remaining <= 0:
+ pending = [fp for fp in in_fps if fp not in ready]
+ stuck = [fp for fp in pending if fp.exists()]
+ missing = [fp for fp in pending if not fp.exists()]
+ parts = []
+ if stuck:
+ parts.append(f"present but unreadable: {', '.join(str(fp) for fp in stuck)}")
+ if missing:
+ parts.append(f"missing: {', '.join(str(fp) for fp in missing)}")
+ raise TimeoutError(
+ f"Timed out after {max_poll_time}s waiting for reduction inputs — " + "; ".join(parts)
+ )
+ # Don't oversleep past the deadline; if remaining < polling_time, wake up sooner.
+ sleep_for = min(polling_time, remaining)
+
logger.info("Waiting to begin reduction for all files to be written...")
- time.sleep(polling_time)
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])
GPT-5.6 Luna · attempt 3 · solved · 11 steps · $0.013
Submitted diff · 106 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py
index 76cd21e..f90958c 100644
--- a/src/MEDS_transforms/dataframe/write_fn.py
+++ b/src/MEDS_transforms/dataframe/write_fn.py
@@ -1,3 +1,4 @@
+import os
from collections.abc import Callable
from pathlib import Path
@@ -13,4 +14,24 @@ def write_df(df: DF_T, out_fp: Path) -> None:
if isinstance(df, pl.LazyFrame):
df = df.collect()
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ # Write beside the destination and publish it in one filesystem operation. In
+ # particular, creating the destination first would allow readers to observe a
+ # parquet file without its footer while it is being written.
+ staging_fp = out_fp.with_name(f"{out_fp.name}.tmp")
+ final_existed = out_fp.exists()
+ try:
+ df.write_parquet(staging_fp, use_pyarrow=True)
+ os.replace(staging_fp, out_fp)
+ except BaseException:
+ # Do not leave either an incomplete staging file or a destination which
+ # may have been created by an unusual failing replace implementation.
+ try:
+ staging_fp.unlink()
+ except OSError:
+ pass
+ if not final_existed:
+ try:
+ out_fp.unlink()
+ except OSError:
+ pass
+ raise
diff --git a/src/MEDS_transforms/mapreduce/reducer.py b/src/MEDS_transforms/mapreduce/reducer.py
index 3b29ceb..472df4e 100644
--- a/src/MEDS_transforms/mapreduce/reducer.py
+++ b/src/MEDS_transforms/mapreduce/reducer.py
@@ -28,6 +28,7 @@ def reduce_over(
merge_fn: REDUCE_FN_T | None = None,
do_overwrite: bool = False,
polling_time: float = 0.1,
+ max_poll_time: float | None = None,
):
"""Performs a reduction operation on a list of input file paths, with optional merging to existing data.
@@ -35,6 +36,7 @@ def reduce_over(
in_fps: List of input file paths containing data over which the reduction should be performed.
out_fp: Output file path where the reduced data will be saved.
polling_time: Time in seconds to wait between checks for file readiness.
+ max_poll_time: Maximum time to wait for readable inputs. ``None`` waits indefinitely.
read_fn: Function to read data from the input file paths.
write_fn: Function to write the reduced data to the output file path.
reduce_fn: Function to perform the reduction operation on the data. It should take two dataframe
@@ -200,12 +202,48 @@ def reduce_over(
└─────┴─────┘
"""
+ if max_poll_time is not None and max_poll_time <= polling_time:
+ raise ValueError("max_poll_time must be greater than polling_time")
+
if out_fp.is_file() and not do_overwrite:
raise FileExistsError(f"Output file already exists: {out_fp.resolve()!s}")
- while not all(fp.is_file() for fp in in_fps):
- logger.info("Waiting to begin reduction for all files to be written...")
- time.sleep(polling_time)
+ # Existence is not sufficient here: a producer may have created the path and
+ # still be writing its parquet footer. Read each file eagerly as a readiness
+ # probe, independently of read_fn (which is commonly lazy).
+ deadline = None if max_poll_time is None else time.monotonic() + max_poll_time
+ while True:
+ missing = [fp for fp in in_fps if not fp.exists()]
+ unreadable = []
+ for fp in in_fps:
+ if fp in missing:
+ continue
+ try:
+ if not fp.is_file():
+ raise ValueError("input is not a file")
+ pl.read_parquet(fp)
+ except Exception:
+ unreadable.append(fp)
+
+ if not missing and not unreadable:
+ break
+
+ if deadline is not None:
+ remaining = deadline - time.monotonic()
+ if remaining <= 0:
+ details = []
+ if unreadable:
+ details.append(
+ "present but unreadable: " + ", ".join(str(fp) for fp in unreadable)
+ )
+ if missing:
+ details.append("missing: " + ", ".join(str(fp) for fp in missing))
+ raise TimeoutError("Timed out waiting for inputs (" + "; ".join(details) + ")")
+ sleep_for = min(polling_time, remaining)
+ else:
+ sleep_for = polling_time
+ logger.info("Waiting to begin reduction for all inputs to be readable...")
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])$ ls && grep -R "def reduce_over\|def write_df" -n . | head -20 && find . -maxdepth 3 -type f | head
output · 26 lines
CITATION.cff LICENSE README.md conftest.py docs example mkdocs.yml pyproject.toml src static tests uv.lock
14 more lines
./src/MEDS_transforms/dataframe/write_fn.py:11:def write_df(df: DF_T, out_fp: Path) -> None: ./src/MEDS_transforms/mapreduce/reducer.py:21:def reduce_over( ./src/MEDS_transforms/mapreduce/reducer.py:58: >>> def write_dfs(in_fps: list[Path], dfs: list[pl.DataFrame] = dfs, delay_per: float = 0): ./src/MEDS_transforms/mapreduce/shard_iteration.py:103: >>> def write_dfs( ./.codecov.yml ./.editorconfig ./.git/COMMIT_EDITMSG ./.git/config ./.git/description ./.git/HEAD ./.git/hooks/applypatch-msg.sample ./.git/hooks/commit-msg.sample ./.git/hooks/fsmonitor-watchman.sample ./.git/hooks/post-update.sample
Plan: inspect reducer and dataframe writing implementations plus their tests and dependency/API conventions; reproduce the race/failure behavior if feasible; implement bounded polling with parquet readability validation and atomic staging publication; then run focused tests and review the diff.
$ cat pyproject.toml | sed -n '1,180p'; grep -R "read_parquet" -n src/MEDS_transforms | head -30
output · 42 lines
[build-system]
requires = ["setuptools>=64", "setuptools-scm>=8.0", "wheel"]
build-backend = "setuptools.build_meta"
[project]
name = "MEDS_transforms"
dynamic = ["version"]
authors = [
{ name="Matthew McDermott", email="mattmcdermott8@gmail.com" },
]
description = "A framework for compiling simple, mapreduce style pipelines over MEDS datasets."
readme = "README.md"30 more lines
requires-python = ">=3.11"
classifiers = [
"Programming Language :: Python :: 3",
"License :: OSI Approved :: MIT License",
"Operating System :: OS Independent",
]
dependencies = [
"polars>=1.35,<2",
"pyarrow",
"hydra-core",
"numpy",
"meds~=0.4.0",
"filelock",
"pytest",
"meds_testing_helpers~=0.3.0",
"pretty-print-directory>=0.1.3",
"yaml_to_disk",
]
[dependency-groups]
dev = ["pre-commit<4", "ruff", "pytest-cov", "pip"]
[tool.setuptools_scm]
[project.optional-dependencies]
local_paral
…[3842 chars omitted]…
t_vocabulary_indices/fit_vocabulary_indices.py:260: ... pl.read_parquet(output_dir / "codes.parquet")
src/MEDS_transforms/stages/fit_vocabulary_indices/fit_vocabulary_indices.py:323: code_metadata = pl.read_parquet(metadata_input_dir / "codes.parquet", use_pyarrow=True)
src/MEDS_transforms/stages/reshard_to_split/reshard_to_split.py:276: partial(pl.read_parquet, use_pyarrow=True),$ find tests -type f | sort | grep -E 'reduc|write|dataframe' && grep -R "reduce_over" -n tests || true; sed -n '1,100p' src/MEDS_transforms/dataframe/__init__.py; grep -R "write_df" -n src tests | head -30
output · 18 lines
from .read_fn import READ_FN_T, read_and_filter_fntr, read_df from .types import DF_T from .write_fn import WRITE_FN_T, write_df __all__ = ["DF_T", "READ_FN_T", "WRITE_FN_T", "read_and_filter_fntr", "read_df", "write_df"] src/MEDS_transforms/dataframe/write_fn.py:11:def write_df(df: DF_T, out_fp: Path) -> None: src/MEDS_transforms/dataframe/__init__.py:3:from .write_fn import WRITE_FN_T, write_df src/MEDS_transforms/dataframe/__init__.py:5:__all__ = ["DF_T", "READ_FN_T", "WRITE_FN_T", "read_and_filter_fntr", "read_df", "write_df"] src/MEDS_transforms/mapreduce/mapper.py:7:from ..dataframe import READ_FN_T, WRITE_FN_T, read_df, write_df src/MEDS_transforms/mapreduce/mapper.py:18: write_fn: WRITE_FN_T = write_df, src/MEDS_transforms/mapreduce/reducer.py:58: >>> def write_dfs(in_fps: list[Path], dfs: list[pl.DataFrame] = dfs, delay_per: float = 0): src/MEDS_transforms/mapreduc
6 more lines
…[1174 chars omitted]… ITE_FN_T, read_and_filter_fntr, read_df, write_df src/MEDS_transforms/mapreduce/stage.py:461: write_fn = write_df src/MEDS_transforms/mapreduce/stage.py:850: write_fn = write_df src/MEDS_transforms/stages/reshard_to_split/reshard_to_split.py:17:from ...dataframe import read_and_filter_fntr, write_df src/MEDS_transforms/stages/reshard_to_split/reshard_to_split.py:328: write_df,
$ cat src/MEDS_transforms/dataframe/read_fn.py; sed -n '1,80p' src/MEDS_transforms/mapreduce/mapper.py; grep -R "reduce_over(" -n src | grep -v 'def'output · 36 lines
"""Functionality for reading input data shards.""" import logging from collections.abc import Callable from pathlib import Path import polars as pl from .types import DF_T logger = logging.getLogger(__name__)
24 more lines
READ_FN_T = Callable[[Path], DF_T]
def read_df(in_fp: Path) -> DF_T:
"""A generic helper to read a dataframe without accounting for globs."""
return pl.scan_parquet(in_fp, glob=False)
def read_and_filter_fntr(filter_expr: pl.Expr, read_fn: READ_FN_T = read_df) -> READ_FN_T:
"""Create a function that reads a DataFrame from a file and filters it based on a given expression.
This is specified as a functor in this way to allow it to modify arbitrary other read functions for use in
different mapreduce pipelines.
Args:
filter_expr: The filter expression to apply to the DataFrame.
read_fn: The read function to use to read the DataFrame
…[4580 chars omitted]…
merge_fn=merge_fn, **kwargs)
src/MEDS_transforms/mapreduce/reducer.py:123: ... reduce_over(in_fps, out_fp, **kwargs)
src/MEDS_transforms/mapreduce/reducer.py:177: ... reduce_over(in_fps, out_fp, **kwargs)
src/MEDS_transforms/mapreduce/reducer.py:186: ... reduce_over(in_fps, out_fp, do_overwrite=True, **kwargs)
src/MEDS_transforms/mapreduce/stage.py:872: reduce_over($ python - <<'PY'
output · 31 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py
index 76cd21e..2417375 100644
--- a/src/MEDS_transforms/dataframe/write_fn.py
+++ b/src/MEDS_transforms/dataframe/write_fn.py
@@ -1,3 +1,4 @@
+import os
from collections.abc import Callable
from pathlib import Path
@@ -13,4 +14,24 @@ def write_df(df: DF_T, out_fp: Path) -> None:
if isinstance(df, pl.LazyFrame):
df = df.collect()19 more lines
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ # Write beside the destination and publish it in one filesystem operation. In
+ # particular, creating the destination first would allow readers to observe a
+ # parquet file without its footer while it is being written.
+ staging_fp = out_fp.with_name(f"{out_fp.name}.tmp")
+ final_existed = out_fp.exists()
+ try:
+ d
…[3333 chars omitted]…
fp) for fp in missing))
+ raise TimeoutError("Timed out waiting for inputs (" + "; ".join(details) + ")")
+ sleep_for = min(polling_time, remaining)
+ else:
+ sleep_for = polling_time
+ logger.info("Waiting to begin reduction for all inputs to be readable...")
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])$ python - <<'PY'
output · 10 lines
shape: (2, 1) ┌─────┐ │ x │ │ --- │ │ i8 │ ╞═════╡ │ 1 │ │ 2 │ └─────┘ TimeoutError Timed out waiting for inputs (missing: /tmp/tmp1ls06u3p/none.parquet)
$ python - <<'PY'
output · 1 lines
False False
$ pytest -q src/MEDS_transforms/mapreduce/reducer.py src/MEDS_transforms/dataframe/write_fn.py --disable-warnings --maxfail=1 2>&1 | tail -80; git diff --check; git status --short
output · 4 lines
[32m.[0m[32m [100%][0m [32m[32m[1m1 passed[0m[32m in 2.28s[0m[0m M src/MEDS_transforms/dataframe/write_fn.py M src/MEDS_transforms/mapreduce/reducer.py
$ python - <<'PY'
output · 36 lines
if out_fp.is_file() and not do_overwrite:
raise FileExistsError(f"Output file already exists: {out_fp.resolve()!s}")
# Existence is not sufficient here: a producer may have created the path and
# still be writing its parquet footer. Read each file eagerly as a readiness
# probe, independently of read_fn (which is commonly lazy).
deadline = None if max_poll_time is None else time.monotonic() + max_poll_time
while True:
missing = [fp for fp in in_fps if not fp.exists()]
unreadable = []
for fp in in_fps:
if fp in missing:24 more lines
continue
try:
if not fp.is_file():
raise ValueError("input is not a file")
pl.read_parquet(fp)
except Exception:
unreadable.append(fp)
if not missing and not unreadable:
bre
…[1072 chars omitted]…
ept BaseException:
# Do not leave either an incomplete staging file or a destination which
# may have been created by an unusual failing replace implementation.
try:
staging_fp.unlink()
except OSError:
pass
if not final_existed:
try:
out_fp.unlink()
except OSError:
pass
raise$ git diff --check && git diff --stat && git diff -- src/MEDS_transforms/dataframe/write_fn.py src/MEDS_transforms/mapreduce/reducer.py | tail -100
output · 29 lines
src/MEDS_transforms/dataframe/write_fn.py | 23 +++++++++++++++-
src/MEDS_transforms/mapreduce/reducer.py | 44 ++++++++++++++++++++++++++++---
2 files changed, 63 insertions(+), 4 deletions(-)
from pathlib import Path
@@ -13,4 +14,24 @@ def write_df(df: DF_T, out_fp: Path) -> None:
if isinstance(df, pl.LazyFrame):
df = df.collect()
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ # Write beside the destination and publish it in one filesystem operation. In
+ # particular, creating the destination first would allow readers to observe a17 more lines
+ # parquet file without its footer while it is being written.
+ staging_fp = out_fp.with_name(f"{out_fp.name}.tmp")
+ final_existed = out_fp.exists()
+ try:
+ df.write_parquet(staging_fp, use_pyarrow=True)
+ os.replace(staging_fp, out_fp)
+ exc
…[3317 chars omitted]…
fp) for fp in missing))
+ raise TimeoutError("Timed out waiting for inputs (" + "; ".join(details) + ")")
+ sleep_for = min(polling_time, remaining)
+ else:
+ sleep_for = polling_time
+ logger.info("Waiting to begin reduction for all inputs to be readable...")
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])$ done: {"notes": "Implemented atomic parquet publication in write_df and bounded, readability-aware polling with timeout diagnostics in reduce_over. Focused reducer/write_fn tests pass; full suite has an unrelated existing doctest failure in rwlock.py."}output · 1 lines
submitted.
GPT-5.6 Luna · attempt 4 · solved · 9 steps · $0.012
Submitted diff · 135 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py
index 76cd21e..2c53be1 100644
--- a/src/MEDS_transforms/dataframe/write_fn.py
+++ b/src/MEDS_transforms/dataframe/write_fn.py
@@ -1,3 +1,6 @@
+"""Helpers for writing dataframe outputs."""
+
+import os
from collections.abc import Callable
from pathlib import Path
@@ -9,8 +12,28 @@ WRITE_FN_T = Callable[[DF_T, Path], None]
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """Write an eager or lazy dataframe to *out_fp* atomically as parquet."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ staging_fp = out_fp.with_name(f"{out_fp.name}.tmp")
+ had_existing_output = out_fp.exists()
+
+ try:
+ df.write_parquet(staging_fp, use_pyarrow=True)
+ os.replace(staging_fp, out_fp)
+ except Exception:
+ # Do not leave a partially written staging file behind. If publishing
+ # failed after replacing the destination, remove only an output that
+ # did not exist before this operation; an existing output is retained.
+ try:
+ staging_fp.unlink(missing_ok=True)
+ except OSError:
+ pass
+ if not had_existing_output:
+ try:
+ out_fp.unlink(missing_ok=True)
+ except OSError:
+ pass
+ raise
diff --git a/src/MEDS_transforms/mapreduce/reducer.py b/src/MEDS_transforms/mapreduce/reducer.py
index 3b29ceb..3316f57 100644
--- a/src/MEDS_transforms/mapreduce/reducer.py
+++ b/src/MEDS_transforms/mapreduce/reducer.py
@@ -28,6 +28,7 @@ def reduce_over(
merge_fn: REDUCE_FN_T | None = None,
do_overwrite: bool = False,
polling_time: float = 0.1,
+ max_poll_time: float | None = None,
):
"""Performs a reduction operation on a list of input file paths, with optional merging to existing data.
@@ -35,6 +36,7 @@ def reduce_over(
in_fps: List of input file paths containing data over which the reduction should be performed.
out_fp: Output file path where the reduced data will be saved.
polling_time: Time in seconds to wait between checks for file readiness.
+ max_poll_time: Maximum time to wait for readable input files. ``None`` waits indefinitely.
read_fn: Function to read data from the input file paths.
write_fn: Function to write the reduced data to the output file path.
reduce_fn: Function to perform the reduction operation on the data. It should take two dataframe
@@ -200,12 +202,70 @@ def reduce_over(
└─────┴─────┘
"""
+ if max_poll_time is not None and max_poll_time <= polling_time:
+ raise ValueError("max_poll_time must be greater than polling_time")
+
if out_fp.is_file() and not do_overwrite:
raise FileExistsError(f"Output file already exists: {out_fp.resolve()!s}")
- while not all(fp.is_file() for fp in in_fps):
- logger.info("Waiting to begin reduction for all files to be written...")
- time.sleep(polling_time)
+ # A mapper may create its destination before it has finished writing it.
+ # Checking only is_file() therefore races with the writer. Read each file
+ # eagerly here: this both checks that it is parquet and that its footer and
+ # contents are complete. Keep track of paths observed on disk so timeout
+ # diagnostics distinguish paths which never appeared from unreadable ones.
+ deadline = None if max_poll_time is None else time.monotonic() + max_poll_time
+ seen_inputs: set[Path] = set()
+ unreadable: list[Path] = []
+ missing: list[Path] = []
+
+ while True:
+ unreadable = []
+ missing = []
+ all_ready = True
+ for fp in in_fps:
+ if not fp.is_file():
+ if fp not in seen_inputs:
+ missing.append(fp)
+ all_ready = False
+ continue
+
+ seen_inputs.add(fp)
+ try:
+ pl.read_parquet(fp)
+ except Exception:
+ unreadable.append(fp)
+ all_ready = False
+
+ if all_ready:
+ break
+
+ if deadline is not None:
+ remaining = deadline - time.monotonic()
+ if remaining <= 0:
+ # Recompute the lists at the deadline so the exception reports
+ # the state observed by the final check, not a prior poll.
+ unreadable = []
+ missing = []
+ for fp in in_fps:
+ if fp.is_file():
+ try:
+ pl.read_parquet(fp)
+ except Exception:
+ unreadable.append(fp)
+ else:
+ missing.append(fp)
+ details = []
+ if unreadable:
+ details.append("present but unreadable: " + ", ".join(map(str, unreadable)))
+ if missing:
+ details.append("missing: " + ", ".join(map(str, missing)))
+ raise TimeoutError("Timed out waiting for input files (" + "; ".join(details) + ")")
+ sleep_for = min(polling_time, remaining)
+ else:
+ sleep_for = polling_time
+
+ logger.info("Waiting to begin reduction for all files to be readable parquet...")
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])Plan:
$ cd /repo && find . -maxdepth 3 -type f | sort | head -200 && grep -R "def reduce_over\|def write_df" -n .
output · 42 lines
./.codecov.yml ./.editorconfig ./.git/COMMIT_EDITMSG ./.git/HEAD ./.git/config ./.git/description ./.git/hooks/applypatch-msg.sample ./.git/hooks/commit-msg.sample ./.git/hooks/fsmonitor-watchman.sample ./.git/hooks/post-update.sample ./.git/hooks/pre-applypatch.sample ./.git/hooks/pre-commit.sample
30 more lines
./.git/hooks/pre-merge-commit.sample ./.git/hooks/pre-push.sample ./.git/hooks/pre-rebase.sample ./.git/hooks/pre-receive.sample ./.git/hooks/prepare-commit-msg.sample ./.git/hooks/push-to-checkout.sample ./.git/hooks/sendemail-validate.sample ./.git/hooks/update.sample ./.git/index ./.git/info/exclude ./.git/logs/HEAD ./.github/workflows/code-quality-main.yaml ./.github/workflows/code-quality-pr.yaml ./.github/workflows/python-build.yaml ./.github/workflows/tests.yaml ./.gitignore ./.pre-commit-config.yaml ./.readthedocs.yaml ./CITATION.cff ./LICENSE ./README.md ./conftest.py ./docs/gen_r …[1028 chars omitted]… _stages.py ./uv.lock ./src/MEDS_transforms/dataframe/write_fn.py:11:def write_df(df: DF_T, out_fp: Path) -> None: ./src/MEDS_transforms/mapreduce/reducer.py:21:def reduce_over( ./src/MEDS_transforms/mapreduce/reducer.py:58: >>> def write_dfs(in_fps: list[Path], dfs: list[pl.DataFrame] = dfs, delay_per: float = 0): ./src/MEDS_transforms/mapreduce/shard_iteration.py:103: >>> def write_dfs(
$ cd /repo && cat pyproject.toml | sed -n '1,220p'
output · 53 lines
[build-system]
requires = ["setuptools>=64", "setuptools-scm>=8.0", "wheel"]
build-backend = "setuptools.build_meta"
[project]
name = "MEDS_transforms"
dynamic = ["version"]
authors = [
{ name="Matthew McDermott", email="mattmcdermott8@gmail.com" },
]
description = "A framework for compiling simple, mapreduce style pipelines over MEDS datasets."
readme = "README.md"41 more lines
requires-python = ">=3.11"
classifiers = [
"Programming Language :: Python :: 3",
"License :: OSI Approved :: MIT License",
"Operating System :: OS Independent",
]
dependencies = [
"polars>=1.35,<2",
"pyarrow",
"hydra-core",
"numpy",
"meds~=0.4.0",
"filelock",
"pytest",
"meds_testing_helpers~=0.3.0",
"pretty-print-directory>=0.1.3",
"yaml_to_disk",
]
[dependency-groups]
dev = ["pre-commit<4", "ruff", "pytest-cov", "pip"]
[tool.setuptools_scm]
[project.optional-dependencies]
local_paral
…[2385 chars omitted]…
"B", # flake8-bugbear
"C4", # flake8-comprehensions
"DTZ", # flake8-datetimez
"RUF", # Ruff-specific rules
"SIM", # flake8-simplify
"TCH", # flake8-type-checking
"N", # pep8-naming
]
[tool.ruff.lint.per-file-ignores]
"__init__.py" = ["F401"] # Unused imports in __init__ files
[tool.ruff.lint.pep8-naming]
extend-ignore-names = ["*MEDS*"] # Ignore MEDS in variable names$ cd /repo && cat src/MEDS_transforms/dataframe/read_fn.py && sed -n '840,900p' src/MEDS_transforms/mapreduce/stage.py && grep -R "write_df" -n tests | head -30
output · 45 lines
"""Functionality for reading input data shards.""" import logging from collections.abc import Callable from pathlib import Path import polars as pl from .types import DF_T logger = logging.getLogger(__name__)
33 more lines
READ_FN_T = Callable[[Path], DF_T]
def read_df(in_fp: Path) -> DF_T:
"""A generic helper to read a dataframe without accounting for globs."""
return pl.scan_parquet(in_fp, glob=False)
def read_and_filter_fntr(filter_expr: pl.Expr, read_fn: READ_FN_T = read_df) -> READ_FN_T:
"""Create a function that reads a DataFrame from a file and filters it based on a given expression.
This is specified as a functor in this way to allow it to modify arbitrary other read functions for use in
different mapreduce pipelines.
Args:
filter_expr: The filter expression to apply to the DataFrame.
read_fn: The read function to use to read the DataFrame
…[2411 chars omitted]…
fn(cfg, cfg.stage_cfg, reduce_fn)
reduce_over(
in_fps=_out_fps,
out_fp=reduce_stage_out_fp,
read_fn=read_fn,
write_fn=write_fn,
reduce_fn=reduce_fn,
merge_fp=merge_fp,
merge_fn=merge_fn,
do_overwrite=cfg.do_overwrite,
polling_time=cfg.polling_time,
)
logger.info(f"Finished reduction in {datetime.now(tz=UTC) - start}")$ cd /repo && python - <<'PY'
output · 39 lines
"""Helpers for writing dataframe outputs.""" import os from collections.abc import Callable from pathlib import Path import polars as pl from .types import DF_T WRITE_FN_T = Callable[[DF_T, Path], None]
27 more lines
def write_df(df: DF_T, out_fp: Path) -> None:
"""Write an eager or lazy dataframe to *out_fp* atomically as parquet."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
out_fp.parent.mkdir(parents=True, exist_ok=True)
staging_fp = out_fp.with_name(f"{out_fp.name}.tmp")
had_existing_output = out_fp.exists()
try:
df.write_parquet(staging_fp, use_pyarrow=True)
os.replace(staging_fp, out_fp)
except Exception:
# Do not leave a partially written staging file behind. If publishing
# failed after replacing the destination, remove only an output that
# did not exist before this operation; an existing output is retained.
try:
staging_fp.unlink(missing_ok=True)
except OSError:
pass
if not had_existing_output:
try:
out_fp.unlink(missing_ok=True)
except OSError:
pass
raise$ cd /repo && python - <<'PY'
output · 35 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py index 76cd21e..2c53be1 100644 --- a/src/MEDS_transforms/dataframe/write_fn.py +++ b/src/MEDS_transforms/dataframe/write_fn.py @@ -1,3 +1,6 @@ +"""Helpers for writing dataframe outputs.""" + +import os from collections.abc import Callable from pathlib import Path @@ -9,8 +12,28 @@ WRITE_FN_T = Callable[[DF_T, Path], None]
23 more lines
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """Write an eager or lazy dataframe to *out_fp* atomically as parquet."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ staging_fp = out_fp.with_name(f"{out_fp.name}.tmp")
+ had_existing_output
…[4403 chars omitted]…
issing)))
+ raise TimeoutError("Timed out waiting for input files (" + "; ".join(details) + ")")
+ sleep_for = min(polling_time, remaining)
+ else:
+ sleep_for = polling_time
+
+ logger.info("Waiting to begin reduction for all files to be readable parquet...")
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])$ cd /repo && python - <<'PY'
output · 29 lines
shape: (2, 1) ┌─────┐ │ x │ │ --- │ │ i8 │ ╞═════╡ │ 1 │ │ 2 │ └─────┘ TimeoutError Timed out waiting for input files (missing: /tmp/tmpfbkzu1cz/missing.parquet) TimeoutError Timed out waiting for input files (present but unreadable: /tmp/tmpwl0_l_n1/x.parquet) bash: line 146: ruff: command not found
17 more lines
/usr/local/lib/python3.12/site-packages/pretty_print_directory/pretty_print_directory.py:204: SyntaxWarning: invalid escape sequence '\.'
>>> config = PrintConfig(ignore_regex=r"\.txt")
.F
=================================== FAILURES ===================================
_________________ test_example_pipeline_parallel_stage_runner __________________
@pytest.mark.parallelized
def test_example_pipeline_parallel_stage_runner():
> log_out = pipeline_tester(PIPELINE_YAML, PARALLEL_STAGE_RUNNER_YAML, STAGE_SCENARIO_SEQUENCE)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
…[8629 chars omitted]…
values/sum_sqd not found in schema
WARNING meds_testing_helpers.dataset:dataset.py:942 Column code/vocab_index not found in schema
=========================== short test summary info ============================
FAILED tests/test_pipeline.py::test_example_pipeline_parallel_stage_runner - ...
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!!
1 failed, 1 passed in 17.99s$ cd /repo && python - <<'PY'
output · 4 lines
replace False False write False False ........................................ [100%] 40 passed in 63.21s (0:01:03)
$ cd /repo && git diff --check && git diff --stat && git diff -- src/MEDS_transforms/mapreduce/reducer.py src/MEDS_transforms/dataframe/write_fn.py | tail -120 && python -m pytest -q tests/test_generic_input.py --disable-warnings --maxfail=1 -o addopts=''
output · 32 lines
src/MEDS_transforms/dataframe/write_fn.py | 27 ++++++++++++-
src/MEDS_transforms/mapreduce/reducer.py | 66 +++++++++++++++++++++++++++++--
2 files changed, 88 insertions(+), 5 deletions(-)
+ """Write an eager or lazy dataframe to *out_fp* atomically as parquet."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ staging_fp = out_fp.with_name(f"{out_fp.name}.tmp")
+ had_existing_output = out_fp.exists()
+20 more lines
+ try:
+ df.write_parquet(staging_fp, use_pyarrow=True)
+ os.replace(staging_fp, out_fp)
+ except Exception:
+ # Do not leave a partially written staging file behind. If publishing
+ # failed after replacing the destination, remove only an output that
+ # did not exist before this operation; an existing output
…[4105 chars omitted]…
ails) + ")")
+ sleep_for = min(polling_time, remaining)
+ else:
+ sleep_for = polling_time
+
+ logger.info("Waiting to begin reduction for all files to be readable parquet...")
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])
.......... [100%]
10 passed in 0.20s$ done: {"notes": "Implemented atomic parquet publication in write_df and readable-parquet polling with bounded timeout diagnostics in reduce_over. Focused tests, reproducers, compile, and diff checks passed."}output · 1 lines
submitted.
GPT-5.6 Luna · attempt 5 · solved · 7 steps · $0.009
Submitted diff · 110 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py
index 76cd21e..41309bc 100644
--- a/src/MEDS_transforms/dataframe/write_fn.py
+++ b/src/MEDS_transforms/dataframe/write_fn.py
@@ -1,3 +1,4 @@
+import os
from collections.abc import Callable
from pathlib import Path
@@ -9,8 +10,27 @@ WRITE_FN_T = Callable[[DF_T, Path], None]
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """Write an eager or lazy dataframe to parquet and publish it atomically."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ staging_fp = out_fp.with_name(out_fp.name + ".tmp")
+ final_existed = out_fp.exists()
+ try:
+ df.write_parquet(staging_fp, use_pyarrow=True)
+ os.replace(staging_fp, out_fp)
+ except Exception:
+ # In particular, don't leave a partial parquet around after either the
+ # write or the publish fails. Only remove the final path if this call
+ # could have created it; preserve an output which predated this call.
+ try:
+ staging_fp.unlink()
+ except FileNotFoundError:
+ pass
+ if not final_existed:
+ try:
+ out_fp.unlink()
+ except FileNotFoundError:
+ pass
+ raise
diff --git a/src/MEDS_transforms/mapreduce/reducer.py b/src/MEDS_transforms/mapreduce/reducer.py
index 3b29ceb..b9746f5 100644
--- a/src/MEDS_transforms/mapreduce/reducer.py
+++ b/src/MEDS_transforms/mapreduce/reducer.py
@@ -28,6 +28,7 @@ def reduce_over(
merge_fn: REDUCE_FN_T | None = None,
do_overwrite: bool = False,
polling_time: float = 0.1,
+ max_poll_time: float | None = None,
):
"""Performs a reduction operation on a list of input file paths, with optional merging to existing data.
@@ -35,6 +36,7 @@ def reduce_over(
in_fps: List of input file paths containing data over which the reduction should be performed.
out_fp: Output file path where the reduced data will be saved.
polling_time: Time in seconds to wait between checks for file readiness.
+ max_poll_time: Maximum time to wait for the input files. ``None`` waits indefinitely.
read_fn: Function to read data from the input file paths.
write_fn: Function to write the reduced data to the output file path.
reduce_fn: Function to perform the reduction operation on the data. It should take two dataframe
@@ -200,12 +202,49 @@ def reduce_over(
└─────┴─────┘
"""
+ if max_poll_time is not None and max_poll_time <= polling_time:
+ raise ValueError("max_poll_time must be greater than polling_time")
+
if out_fp.is_file() and not do_overwrite:
raise FileExistsError(f"Output file already exists: {out_fp.resolve()!s}")
- while not all(fp.is_file() for fp in in_fps):
+ deadline = None if max_poll_time is None else time.monotonic() + max_poll_time
+ while True:
+ missing: list[Path] = []
+ unreadable: list[Path] = []
+ for fp in in_fps:
+ if not fp.is_file():
+ # A path which exists but is not a regular file is present, but
+ # cannot be a readable parquet input.
+ (unreadable if fp.exists() else missing).append(fp)
+ continue
+ try:
+ # Merely seeing the path is not enough: writers can create it
+ # before the parquet footer has been written.
+ pl.read_parquet(fp)
+ except Exception:
+ unreadable.append(fp)
+
+ if not missing and not unreadable:
+ break
+
+ if deadline is not None:
+ remaining = deadline - time.monotonic()
+ if remaining <= 0:
+ details = []
+ if unreadable:
+ details.append(
+ "present but unreadable: " + ", ".join(str(fp) for fp in unreadable)
+ )
+ if missing:
+ details.append("missing: " + ", ".join(str(fp) for fp in missing))
+ raise TimeoutError("; ".join(details))
+ sleep_for = min(polling_time, remaining)
+ else:
+ sleep_for = polling_time
+
logger.info("Waiting to begin reduction for all files to be written...")
- time.sleep(polling_time)
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])Plan:
$ ls && find . -maxdepth 3 -type f | sort | head -200 && grep -R "def reduce_over\|def write_df" -n .
output · 48 lines
CITATION.cff LICENSE README.md conftest.py docs example mkdocs.yml pyproject.toml src static tests uv.lock
36 more lines
./.codecov.yml ./.editorconfig ./.git/COMMIT_EDITMSG ./.git/HEAD ./.git/config ./.git/description ./.git/hooks/applypatch-msg.sample ./.git/hooks/commit-msg.sample ./.git/hooks/fsmonitor-watchman.sample ./.git/hooks/post-update.sample ./.git/hooks/pre-applypatch.sample ./.git/hooks/pre-commit.sample ./.git/hooks/pre-merge-commit.sample ./.git/hooks/pre-push.sample ./.git/hooks/pre-rebase.sample ./.git/hooks/pre-receive.sample ./.git/hooks/prepare-commit-msg.sample ./.git/hooks/push-to-checkout.sample ./.git/hooks/sendemail-validate.sample ./.git/hooks/update.sample ./.git/index ./.git/info/exclude ./.git/logs/HEAD ./.github/workflows/code-quality-main.yaml ./.github/workflows/code-quality-pr.yaml ./.github/workflows/python-build.yaml ./.github/workflows/tests.yaml ./.gitignore ./ …[1135 chars omitted]… _stages.py ./uv.lock ./src/MEDS_transforms/dataframe/write_fn.py:11:def write_df(df: DF_T, out_fp: Path) -> None: ./src/MEDS_transforms/mapreduce/reducer.py:21:def reduce_over( ./src/MEDS_transforms/mapreduce/reducer.py:58: >>> def write_dfs(in_fps: list[Path], dfs: list[pl.DataFrame] = dfs, delay_per: float = 0): ./src/MEDS_transforms/mapreduce/shard_iteration.py:103: >>> def write_dfs(
$ cat pyproject.toml | sed -n '1,180p'; find src/MEDS_transforms/mapreduce src/MEDS_transforms/dataframe -type f -maxdepth 2 -print
output · 49 lines
[build-system]
requires = ["setuptools>=64", "setuptools-scm>=8.0", "wheel"]
build-backend = "setuptools.build_meta"
[project]
name = "MEDS_transforms"
dynamic = ["version"]
authors = [
{ name="Matthew McDermott", email="mattmcdermott8@gmail.com" },
]
description = "A framework for compiling simple, mapreduce style pipelines over MEDS datasets."
readme = "README.md"37 more lines
requires-python = ">=3.11"
classifiers = [
"Programming Language :: Python :: 3",
"License :: OSI Approved :: MIT License",
"Operating System :: OS Independent",
]
dependencies = [
"polars>=1.35,<2",
"pyarrow",
"hydra-core",
"numpy",
"meds~=0.4.0",
"filelock",
"pytest",
"meds_testing_helpers~=0.3.0",
"pretty-print-directory>=0.1.3",
"yaml_to_disk",
]
[dependency-groups]
dev = ["pre-commit<4", "ruff", "pytest-cov", "pip"]
[tool.setuptools_scm]
[project.optional-dependencies]
local_paral
…[2800 chars omitted]…
ransforms/mapreduce/mapper.py
src/MEDS_transforms/mapreduce/reducer.py
src/MEDS_transforms/mapreduce/rwlock.py
src/MEDS_transforms/mapreduce/shard_iteration.py
src/MEDS_transforms/mapreduce/stage.py
src/MEDS_transforms/mapreduce/__init__.py
src/MEDS_transforms/dataframe/read_fn.py
src/MEDS_transforms/dataframe/types.py
src/MEDS_transforms/dataframe/write_fn.py
src/MEDS_transforms/dataframe/__init__.py$ cat > /tmp/reducer.py.new <<'EOF'
output · 36 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py index 76cd21e..41309bc 100644 --- a/src/MEDS_transforms/dataframe/write_fn.py +++ b/src/MEDS_transforms/dataframe/write_fn.py @@ -1,3 +1,4 @@ +import os from collections.abc import Callable from pathlib import Path @@ -9,8 +10,27 @@ WRITE_FN_T = Callable[[DF_T, Path], None]
24 more lines
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """Write an eager or lazy dataframe to parquet and publish it atomically."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ staging_fp = out_fp.with_name(out_fp.name + ".tmp")
+ final_existed = out_fp.exists()
+ try:
+ df.write_parq
…[3378 chars omitted]…
in(str(fp) for fp in missing))
+ raise TimeoutError("; ".join(details))
+ sleep_for = min(polling_time, remaining)
+ else:
+ sleep_for = polling_time
+
logger.info("Waiting to begin reduction for all files to be written...")
- time.sleep(polling_time)
+ time.sleep(sleep_for)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])$ python - <<'PY'
output · 15 lines
TimeoutError present but unreadable: /tmp/tmp986fqjsz/x.parquet
TimeoutError missing: /tmp/tmp986fqjsz/x.parquet, /tmp/tmp986fqjsz/bad
/usr/local/lib/python3.12/site-packages/pretty_print_directory/pretty_print_directory.py:204: SyntaxWarning: invalid escape sequence '\.'
>>> config = PrintConfig(ignore_regex=r"\.txt")
[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[32m.[0m[31mF[0m
=================================== FAILURES ===================================
[31m[1m_________________ test_example_pipeline_parallel_stage_runner __________________[0m
[0m[37m@pytest[39;49;00m.mark.parallelized[90m[39;49;00m
[94mde
…[8630 chars omitted]…
============= short test summary info ============================[0m3 more lines
[31mFAILED[0m tests/test_pipeline.py::[1mtest_example_pipeline_parallel_stage_runner[0m - AssertionError: Pipeline failed with error: Pipeline returned code 1. [31m!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!![0m [31m[31m[1m1 failed[0m, [32m31 passed[0m[31m in 86.83s (0:01:26)[0m[0m
$ git status --short && git diff --check
output · 2 lines
M src/MEDS_transforms/dataframe/write_fn.py M src/MEDS_transforms/mapreduce/reducer.py
$ pytest -q src/MEDS_transforms/mapreduce/reducer.py src/MEDS_transforms/dataframe/write_fn.py --disable-warnings --maxfail=1 && git diff --stat && git diff --check
output · 5 lines
[32m.[0m[32m [100%][0m [32m[32m[1m1 passed[0m[32m in 2.08s[0m[0m src/MEDS_transforms/dataframe/write_fn.py | 24 +++++++++++++++-- src/MEDS_transforms/mapreduce/reducer.py | 43 +++++++++++++++++++++++++++++-- 2 files changed, 63 insertions(+), 4 deletions(-)
$ done: {"notes": "Implemented parquet-readiness polling with bounded deadlines, validation errors, and max_poll_time; made write_df atomically publish via .tmp staging with cleanup and preserved exceptions. Focused doctests and existing tests passed; one parallel pipeline test failed in its unrelated integration execution."}output · 1 lines
submitted.
GPT-5.6 Luna · attempt 6 · solved · 13 steps · $0.013
Submitted diff · 109 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py
index 76cd21e..570db93 100644
--- a/src/MEDS_transforms/dataframe/write_fn.py
+++ b/src/MEDS_transforms/dataframe/write_fn.py
@@ -1,3 +1,4 @@
+import os
from collections.abc import Callable
from pathlib import Path
@@ -9,8 +10,28 @@ WRITE_FN_T = Callable[[DF_T, Path], None]
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """Write a dataframe to parquet, publishing the result atomically."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ staging_fp = out_fp.with_name(out_fp.name + ".tmp")
+ output_existed = out_fp.exists()
+ try:
+ df.write_parquet(staging_fp, use_pyarrow=True)
+ os.replace(staging_fp, out_fp)
+ except Exception:
+ # In particular, do not leave a file which another process could mistake
+ # for a completed output. Keep an output which predates this operation.
+ try:
+ staging_fp.unlink()
+ except OSError:
+ # Preserve the exception raised by writing or publishing. This
+ # also handles a staging file which was already removed.
+ pass
+ if not output_existed:
+ try:
+ out_fp.unlink()
+ except OSError:
+ pass
+ raise
diff --git a/src/MEDS_transforms/mapreduce/reducer.py b/src/MEDS_transforms/mapreduce/reducer.py
index 3b29ceb..26ad802 100644
--- a/src/MEDS_transforms/mapreduce/reducer.py
+++ b/src/MEDS_transforms/mapreduce/reducer.py
@@ -28,6 +28,7 @@ def reduce_over(
merge_fn: REDUCE_FN_T | None = None,
do_overwrite: bool = False,
polling_time: float = 0.1,
+ max_poll_time: float | None = None,
):
"""Performs a reduction operation on a list of input file paths, with optional merging to existing data.
@@ -35,6 +36,7 @@ def reduce_over(
in_fps: List of input file paths containing data over which the reduction should be performed.
out_fp: Output file path where the reduced data will be saved.
polling_time: Time in seconds to wait between checks for file readiness.
+ max_poll_time: Maximum time to wait for inputs. ``None`` waits indefinitely.
read_fn: Function to read data from the input file paths.
write_fn: Function to write the reduced data to the output file path.
reduce_fn: Function to perform the reduction operation on the data. It should take two dataframe
@@ -200,12 +202,47 @@ def reduce_over(
└─────┴─────┘
"""
+ if max_poll_time is not None and max_poll_time <= polling_time:
+ raise ValueError("max_poll_time must be greater than polling_time")
+
if out_fp.is_file() and not do_overwrite:
raise FileExistsError(f"Output file already exists: {out_fp.resolve()!s}")
- while not all(fp.is_file() for fp in in_fps):
+ deadline = None if max_poll_time is None else time.monotonic() + max_poll_time
+
+ def is_readable_parquet(fp: Path) -> bool:
+ if not fp.is_file():
+ return False
+ try:
+ # Reading (rather than only checking the parquet footer) catches a
+ # writer which has published a path before all row data is present.
+ pl.read_parquet(fp)
+ except Exception:
+ return False
+ return True
+
+ while True:
+ unreadable = [fp for fp in in_fps if not is_readable_parquet(fp)]
+ if not unreadable:
+ break
+
logger.info("Waiting to begin reduction for all files to be written...")
- time.sleep(polling_time)
+ if deadline is not None:
+ remaining = deadline - time.monotonic()
+ if remaining <= 0:
+ missing = [fp for fp in unreadable if not fp.exists()]
+ present_unreadable = [fp for fp in unreadable if fp.exists()]
+ details = []
+ if present_unreadable:
+ details.append(
+ "present but unreadable: " + ", ".join(str(fp) for fp in present_unreadable)
+ )
+ if missing:
+ details.append("missing: " + ", ".join(str(fp) for fp in missing))
+ raise TimeoutError("Timed out waiting for input files (" + "; ".join(details) + ")")
+ time.sleep(min(polling_time, remaining))
+ else:
+ time.sleep(polling_time)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])Plan:
$ cd /repo && ls && grep -R "def reduce_over\|def write_df" -n MEDS_transforms tests || true && find . -maxdepth 3 -type f | head -80
output · 55 lines
CITATION.cff LICENSE README.md conftest.py docs example mkdocs.yml pyproject.toml src static tests uv.lock
43 more lines
grep: MEDS_transforms: No such file or directory ./.codecov.yml ./.editorconfig ./.git/COMMIT_EDITMSG ./.git/config ./.git/description ./.git/HEAD ./.git/hooks/applypatch-msg.sample ./.git/hooks/commit-msg.sample ./.git/hooks/fsmonitor-watchman.sample ./.git/hooks/post-update.sample ./.git/hooks/pre-applypatch.sample ./.git/hooks/pre-commit.sample ./.git/hooks/pre-merge-commit.sample ./.git/hooks/pre-push.sample ./.git/hooks/pre-rebase.sample ./.git/hooks/pre-receive.sample ./.git/hooks/prepare-commit-msg.sample ./.git/hooks/push-to-checkout.sample ./.git/hooks/sendemail-validate.sample ./.git/hooks/update.sample ./.git/index ./.git/info/exclude ./.git/logs/HEAD ./.github/workflows/code-quality-main.yaml ./.github/workflows/code-quality-pr.yaml ./.github/workflows/python-build.ya …[800 chars omitted]… quires.txt ./src/MEDS_transforms.egg-info/scm_file_list.json ./src/MEDS_transforms.egg-info/SOURCES.txt ./src/MEDS_transforms.egg-info/top_level.txt ./static/logo_bg.png ./static/logo_dark.svg ./static/logo_light.svg ./tests/conftest.py ./tests/test_example_pkg.py ./tests/test_generic_input.py ./tests/test_main.py ./tests/test_pipeline.py ./tests/test_registered_stages.py ./tests/__init__.py ./uv.lock
$ cd /repo && find src/MEDS_transforms -type f | sort | grep -E 'mapreduce|dataframe|write|reducer' && grep -R "reduce_over\|write_df" -n src tests | head -100
output · 21 lines
src/MEDS_transforms/dataframe/__init__.py src/MEDS_transforms/dataframe/read_fn.py src/MEDS_transforms/dataframe/types.py src/MEDS_transforms/dataframe/write_fn.py src/MEDS_transforms/mapreduce/__init__.py src/MEDS_transforms/mapreduce/mapper.py src/MEDS_transforms/mapreduce/reducer.py src/MEDS_transforms/mapreduce/rwlock.py src/MEDS_transforms/mapreduce/shard_iteration.py src/MEDS_transforms/mapreduce/stage.py src/MEDS_transforms/dataframe/write_fn.py:11:def write_df(df: DF_T, out_fp: Path) -> None: src/MEDS_transforms/dataframe/__init__.py:3:from .write_fn import WRITE_FN_T, write_df
9 more lines
src/MEDS_transforms/dataframe/__init__.py:5:__all__ = ["DF_T", "READ_FN_T", "WRITE_FN_T", "read_and_filter_fntr", "read_df", "write_df"] src/MEDS_transforms/mapreduce/mapper.py:7:from ..dataframe import READ_FN_T, WRITE_FN_T, read_df, write_df src/MEDS_transforms/mapreduce/mapper.py:18: write_fn: WRI …[2125 chars omitted]… ransforms/mapreduce/stage.py:461: write_fn = write_df src/MEDS_transforms/mapreduce/stage.py:850: write_fn = write_df src/MEDS_transforms/mapreduce/stage.py:872: reduce_over( src/MEDS_transforms/stages/reshard_to_split/reshard_to_split.py:17:from ...dataframe import read_and_filter_fntr, write_df src/MEDS_transforms/stages/reshard_to_split/reshard_to_split.py:328: write_df,
$ cd /repo && cat pyproject.toml | sed -n '1,180p' && find tests -type f | sort
output · 55 lines
[build-system]
requires = ["setuptools>=64", "setuptools-scm>=8.0", "wheel"]
build-backend = "setuptools.build_meta"
[project]
name = "MEDS_transforms"
dynamic = ["version"]
authors = [
{ name="Matthew McDermott", email="mattmcdermott8@gmail.com" },
]
description = "A framework for compiling simple, mapreduce style pipelines over MEDS datasets."
readme = "README.md"43 more lines
requires-python = ">=3.11"
classifiers = [
"Programming Language :: Python :: 3",
"License :: OSI Approved :: MIT License",
"Operating System :: OS Independent",
]
dependencies = [
"polars>=1.35,<2",
"pyarrow",
"hydra-core",
"numpy",
"meds~=0.4.0",
"filelock",
"pytest",
"meds_testing_helpers~=0.3.0",
"pretty-print-directory>=0.1.3",
"yaml_to_disk",
]
[dependency-groups]
dev = ["pre-commit<4", "ruff", "pytest-cov", "pip"]
[tool.setuptools_scm]
[project.optional-dependencies]
local_paral
…[2549 chars omitted]…
flake8-type-checking
"N", # pep8-naming
]
[tool.ruff.lint.per-file-ignores]
"__init__.py" = ["F401"] # Unused imports in __init__ files
[tool.ruff.lint.pep8-naming]
extend-ignore-names = ["*MEDS*"] # Ignore MEDS in variable names
tests/__init__.py
tests/conftest.py
tests/test_example_pkg.py
tests/test_generic_input.py
tests/test_main.py
tests/test_pipeline.py
tests/test_registered_stages.py$ cd /repo && python - <<'PY'
output · 34 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py index 76cd21e..e08c64a 100644 --- a/src/MEDS_transforms/dataframe/write_fn.py +++ b/src/MEDS_transforms/dataframe/write_fn.py @@ -1,3 +1,4 @@ +import os from collections.abc import Callable from pathlib import Path @@ -9,8 +10,26 @@ WRITE_FN_T = Callable[[DF_T, Path], None]
22 more lines
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """Write a dataframe to parquet, publishing the result atomically."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ staging_fp = out_fp.with_name(out_fp.name + ".tmp")
+ output_existed = out_fp.exists()
+ try:
+ df.write_parquet(st
…[3288 chars omitted]…
)
+ if missing:
+ details.append("missing: " + ", ".join(str(fp) for fp in missing))
+ raise TimeoutError("Timed out waiting for input files (" + "; ".join(details) + ")")
+ time.sleep(min(polling_time, remaining))
+ else:
+ time.sleep(polling_time)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])$ cd /repo && pytest -q tests --disable-warnings --maxfail=1`, timeout 120000}迷?}
output · 1 lines
(command exceeded 600s)
$ cd /repo && ruff check src/MEDS_transforms/dataframe/write_fn.py src/MEDS_transforms/mapreduce/reducer.py || true
output · 1 lines
bash: line 37: ruff: command not found
$ cd /repo && pytest -q src/MEDS_transforms/mapreduce/reducer.py src/MEDS_transforms/dataframe/write_fn.py --disable-warnings --maxfail=1 --doctest-modules -o doctest_optionflags=NORMALIZE_WHITESPACE --timeout=30 2>/dev/null || true; git diff --check; git status --short
output · 2 lines
M src/MEDS_transforms/dataframe/write_fn.py M src/MEDS_transforms/mapreduce/reducer.py
$ cd /repo && python - <<'PY'
output · 37 lines
import os from collections.abc import Callable from pathlib import Path import polars as pl from .types import DF_T WRITE_FN_T = Callable[[DF_T, Path], None] def write_df(df: DF_T, out_fp: Path) -> None:
25 more lines
"""Write a dataframe to parquet, publishing the result atomically."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
out_fp.parent.mkdir(parents=True, exist_ok=True)
staging_fp = out_fp.with_name(out_fp.name + ".tmp")
output_existed = out_fp.exists()
try:
df.write_parquet(staging_fp, use_pyarrow=True)
os.replace(staging_fp, out_fp)
except Exception:
# In particular, do not leave a file which another process could mistake
# for a completed output. Keep an output which predates this operation.
try:
staging_fp.unlink()
except OSError:
# Preserve the exception raised by writing or publishing. This
# also handles a staging file which was already removed.
pass
if not output_existed:
try:
out_fp.unlink()
except OSError:
pass
raise$ cd /repo && python - <<'PY'
output · 35 lines
shape: (1, 1)
┌─────┐
│ a │
│ --- │
│ i64 │
╞═════╡
│ 1 │
└─────┘ [PosixPath('/tmp/tmpk8q6j105/nested/a.parquet')]
logger.info("Starting reduction process")
start = datetime.now(tz=UTC)
merge_fp = Path(cfg.stage_cfg.metadata_input_dir) / "codes.parquet"23 more lines
reduce_stage_out_fp = Path(cfg.stage_cfg.reducer_output_dir) / "codes.parquet"
if merge_fn is None:
join_cols = [DataSchema.code_name, *cfg.get("code_modifier_cols", [])]
merge_fn = partial(join_and_replace, join_cols=join_cols)
reduce_fn = bind_compute_fn(cfg, cfg.stage_cfg, reduce_fn)
reduce_over(
in_fps=_out_fps,
out_fp=reduce_stage_out_fp,
read_fn=read_fn,
write_fn=write_fn,
reduce_fn=reduce_fn,
merge_fp=merge_fp,
merge_fn=merge_fn,
do_overwrite=cfg.do_overwrite,
polling_time=cfg.polling_time,
)
logger.info(f"Finished reduction in {datetime.now(tz=UTC) - start}")
src/MEDS_transforms/dataframe/write_fn.py | 25 +++++++++++++++++--
src/MEDS_transforms/mapreduce/reducer.py | 41 +++++++++++++++++++++++++++++--
2 files changed, 62 insertions(+), 4 deletions(-)$ cd /repo && python -m doctest src/MEDS_transforms/dataframe/write_fn.py src/MEDS_transforms/mapreduce/reducer.py -o NORMALIZE_WHITESPACE -o ELLIPSIS && git diff -- src/MEDS_transforms/dataframe/write_fn.py src/MEDS_transforms/mapreduce/reducer.py | tail -100
output · 12 lines
Traceback (most recent call last):
File "<frozen runpy>", line 198, in _run_module_as_main
File "<frozen runpy>", line 88, in _run_code
File "/usr/local/lib/python3.12/doctest.py", line 2845, in <module>
sys.exit(_test())
^^^^^^^
File "/usr/local/lib/python3.12/doctest.py", line 2833, in _test
m = __import__(filename[:-3])
^^^^^^^^^^^^^^^^^^^^^^^^^
File "/repo/src/MEDS_transforms/dataframe/write_fn.py", line 7, in <module>
from .types import DF_T
ImportError: attempted relative import with no known parent package$ cd /repo && python - <<'PY'
output · 3 lines
src/MEDS_transforms/dataframe/write_fn.py | 25 +++++++++++++++++-- src/MEDS_transforms/mapreduce/reducer.py | 41 +++++++++++++++++++++++++++++-- 2 files changed, 62 insertions(+), 4 deletions(-)
$ cd /repo && python -m compileall -q src && git diff --check && git diff -- src/MEDS_transforms/mapreduce/reducer.py src/MEDS_transforms/dataframe/write_fn.py
output · 34 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py index 76cd21e..570db93 100644 --- a/src/MEDS_transforms/dataframe/write_fn.py +++ b/src/MEDS_transforms/dataframe/write_fn.py @@ -1,3 +1,4 @@ +import os from collections.abc import Callable from pathlib import Path @@ -9,8 +10,28 @@ WRITE_FN_T = Callable[[DF_T, Path], None]
22 more lines
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """Write a dataframe to parquet, publishing the result atomically."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ staging_fp = out_fp.with_name(out_fp.name + ".tmp")
+ output_existed = out_fp.exists()
+ try:
+ df.write_parquet(st
…[3413 chars omitted]…
)
+ if missing:
+ details.append("missing: " + ", ".join(str(fp) for fp in missing))
+ raise TimeoutError("Timed out waiting for input files (" + "; ".join(details) + ")")
+ time.sleep(min(polling_time, remaining))
+ else:
+ time.sleep(polling_time)
reduced = reduce_fn(*[read_fn(fp) for fp in in_fps])$ done: {"notes": "Implemented atomic parquet publishing with cleanup in write_df, and bounded readable-parquet polling with timeout classification in reduce_over. Ran reproducers, compileall, and diff checks."}output · 1 lines
submitted.
GLM-5.3 Flash · attempt 2 · solved · 10 steps · $0.002
Submitted diff · 100 lines
diff --git a/src/MEDS_transforms/dataframe/write_fn.py b/src/MEDS_transforms/dataframe/write_fn.py
index 76cd21e..81d207b 100644
--- a/src/MEDS_transforms/dataframe/write_fn.py
+++ b/src/MEDS_transforms/dataframe/write_fn.py
@@ -1,3 +1,4 @@
+import os
from collections.abc import Callable
from pathlib import Path
@@ -9,8 +10,25 @@ WRITE_FN_T = Callable[[DF_T, Path], None]
def write_df(df: DF_T, out_fp: Path) -> None:
- """A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
+ """A generic helper to write a dataframe, either lazy or eager, to a parquet file.
+
+ The output is written to a temporary staging path and then atomically published to the final
+ location, so that readers never observe a partially written output file. If writing or publishing
+ fails, the staging file is removed and the original error is re-raised.
+
+ Args:
+ df: The dataframe (eager or lazy) to write.
+ out_fp: The output parquet file path.
+ """
if isinstance(df, pl.LazyFrame):
df = df.collect()
+
out_fp.parent.mkdir(parents=True, exist_ok=True)
- df.write_parquet(out_fp, use_pyarrow=True)
+ tmp_fp = out_fp.parent / (out_fp.name + ".tmp")
+
+ try:
+ df.write_parquet(tmp_fp, use_pyarrow=True)
+ os.replace(tmp_fp, out_fp)
+ except Exception:
+ tmp_fp.unlink(missing_ok=True)
+ raise
diff --git a/src/MEDS_transforms/mapreduce/reducer.py b/src/MEDS_transforms/mapreduce/reducer.py
index 3b29ceb..6af3a14 100644
--- a/src/MEDS_transforms/mapreduce/reducer.py
+++ b/src/MEDS_transforms/mapreduce/reducer.py
@@ -28,6 +28,7 @@ def reduce_over(
merge_fn: REDUCE_FN_T | None = None,
do_overwrite: bool = False,
polling_time: float = 0.1,
+ max_poll_time: float | None = None,
):
"""Performs a reduction operation on a list of input file paths, with optional merging to existing data.
@@ -42,9 +43,13 @@ def reduce_over(
merge_fp: If this file exists, merge the output with the data stored in this file before finalization.
merge_fn: A special reducer to perform the merging, if `merge_fp` is specified and is a file.
do_overwrite: Should this overwrite an existing out file?
+ max_poll_time: Maximum total time (in seconds) to wait for all input files to become readable
+ parquet files. If `None`, the reducer will wait indefinitely.
Raises:
FileExistsError: If the output file already exists.
+ ValueError: If `max_poll_time` is not greater than `polling_time`.
+ TimeoutError: If `max_poll_time` is exceeded while waiting for inputs.
Examples:
>>> def reduce_fn(*dfs: pl.DataFrame) -> pl.DataFrame:
@@ -203,7 +208,36 @@ def reduce_over(
if out_fp.is_file() and not do_overwrite:
raise FileExistsError(f"Output file already exists: {out_fp.resolve()!s}")
- while not all(fp.is_file() for fp in in_fps):
+ if max_poll_time is not None and max_poll_time <= polling_time:
+ raise ValueError(
+ f"max_poll_time ({max_poll_time}) must be greater than polling_time ({polling_time})"
+ )
+
+ deadline = None if max_poll_time is None else time.monotonic() + max_poll_time
+
+ while True:
+ missing = [fp for fp in in_fps if not fp.is_file()]
+ unreadable = []
+ for fp in in_fps:
+ if fp.is_file():
+ try:
+ pl.read_parquet_schema(fp)
+ except Exception:
+ unreadable.append(fp)
+
+ if not missing and not unreadable:
+ break
+
+ if deadline is not None and time.monotonic() >= deadline:
+ reasons = []
+ if unreadable:
+ reasons.append(f"present but unreadable: {unreadable}")
+ if missing:
+ reasons.append(f"missing: {missing}")
+ raise TimeoutError(
+ "Timed out waiting for input files to be readable parquet: " + "; ".join(reasons)
+ )
+
logger.info("Waiting to begin reduction for all files to be written...")
time.sleep(polling_time)$ find /repo -name "*.py" | head -30; echo ---; cat /repo/MEDS_transforms/mapreduce/reducer.py; echo ---; cat /repo/MEDS_transforms/dataframe/write_fn.py
output · 26 lines
/repo/conftest.py /repo/docs/gen_ref_pages.py /repo/docs/gen_stage_docs.py /repo/example/simple_example_pkg/src/simple_example_pkg/csv_to_meds/csv_to_meds.py /repo/example/simple_example_pkg/src/simple_example_pkg/csv_to_meds/__init__.py /repo/example/simple_example_pkg/src/simple_example_pkg/drop_regex/drop_regex.py /repo/example/simple_example_pkg/src/simple_example_pkg/drop_regex/__init__.py /repo/example/simple_example_pkg/src/simple_example_pkg/export_code_summary/export_code_summary.py /repo/example/simple_example_pkg/src/simple_example_pkg/export_code_summary/__init__.py /repo/example/simple_example_pkg/src/simple_example_pkg/__init__.py /repo/example/simple_example_pkg/tests/test_registered_stages.py /repo/src/MEDS_transforms/compute_modes/compute_fn.py
14 more lines
/repo/src/MEDS_transforms/compute_modes/match_revise.py /repo/src/MEDS_transforms/compute_modes/__init__.py /repo/src/MEDS_tr …[471 chars omitted]… duce/reducer.py /repo/src/MEDS_transforms/mapreduce/rwlock.py /repo/src/MEDS_transforms/mapreduce/shard_iteration.py /repo/src/MEDS_transforms/mapreduce/stage.py /repo/src/MEDS_transforms/mapreduce/__init__.py /repo/src/MEDS_transforms/parser.py --- cat: /repo/MEDS_transforms/mapreduce/reducer.py: No such file or directory --- cat: /repo/MEDS_transforms/dataframe/write_fn.py: No such file or directory
$ cat /repo/src/MEDS_transforms/mapreduce/reducer.py; echo ---; cat /repo/src/MEDS_transforms/dataframe/write_fn.py
output · 51 lines
"""Basic utilities for serialized reduce operations on sharded MEDS datasets with caching and locking.""" import logging import time from pathlib import Path from typing import Protocol import polars as pl from ..dataframe import DF_T, READ_FN_T, WRITE_FN_T logger = logging.getLogger(__name__)
39 more lines
class REDUCE_FN_T(Protocol): # noqa: N801
"""Protocol for a function that takes a variable number dataframes and returns one dataframe."""
def __call__(self, *dfs: DF_T) -> DF_T: ...
def reduce_over(
in_fps: list[Path],
out_fp: Path,
read_fn: READ_FN_T,
write_fn: WRITE_FN_T,
reduce_fn: REDUCE_FN_T,
merge_fp: Path | None = None,
merge_fn: REDUCE_FN_T | None = None,
do_overwrite: bool = False,
polling_time: float = 0.1,
):
"""Performs a reduction operation on a list of input file paths, with optional merging to existing data.
Args:
…[7409 chars omitted]…
ib import Path
import polars as pl
from .types import DF_T
WRITE_FN_T = Callable[[DF_T, Path], None]
def write_df(df: DF_T, out_fp: Path) -> None:
"""A generic helper to write a dataframe, either lazy or eager, to a parquet file."""
if isinstance(df, pl.LazyFrame):
df = df.collect()
out_fp.parent.mkdir(parents=True, exist_ok=True)
df.write_parquet(out_fp, use_pyarrow=True)$ ls /repo/tests; grep -rn "reduce_over\|write_df" /repo/tests | head -30
output · 7 lines
__init__.py conftest.py test_example_pkg.py test_generic_input.py test_main.py test_pipeline.py test_registered_stages.py
$ cat /repo/src/MEDS_transforms/dataframe/read_fn.py; sed -n 1,60p /repo/src/MEDS_transforms/dataframe/__init__.py
output · 45 lines
"""Functionality for reading input data shards.""" import logging from collections.abc import Callable from pathlib import Path import polars as pl from .types import DF_T logger = logging.getLogger(__name__)
33 more lines
READ_FN_T = Callable[[Path], DF_T]
def read_df(in_fp: Path) -> DF_T:
"""A generic helper to read a dataframe without accounting for globs."""
return pl.scan_parquet(in_fp, glob=False)
def read_and_filter_fntr(filter_expr: pl.Expr, read_fn: READ_FN_T = read_df) -> READ_FN_T:
"""Create a function that reads a DataFrame from a file and filters it based on a given expression.
This is specified as a functor in this way to allow it to modify arbitrary other read functions for use in
different mapreduce pipelines.
Args:
filter_expr: The filter expression to apply to the DataFrame.
read_fn: The read function to use to read the DataFrame
…[1279 chars omitted]…
═══╡
│ 5 ┆ 8 │
└─────┴─────┘
"""
def read_and_filter(in_fp: Path) -> DF_T:
return read_fn(in_fp).filter(filter_expr)
return read_and_filter
from .read_fn import READ_FN_T, read_and_filter_fntr, read_df
from .types import DF_T
from .write_fn import WRITE_FN_T, write_df
__all__ = ["DF_T", "READ_FN_T", "WRITE_FN_T", "read_and_filter_fntr", "read_df", "write_df"]$ cat > /repo/src/MEDS_transforms/dataframe/write_fn.py <<'EOF'
output · 1 lines
ok
$ cd /repo && python - <<'PY'
output · 3 lines
Traceback (most recent call last): File "<stdin>", line 14, in <module> AttributeError: 'PosixPath' object has no attribute 'write_parquet'
$ cd /repo && python - <<'PY'
output · 6 lines
OK: Timed out waiting for input files to be readable parquet: missing: [PosixPath('/tmp/tmp7gg_t9ok/in_1.parquet'), PosixPath('/tmp/tmp7gg_t9ok/in_2.parquet')]
OK: max_poll_time (0.5) must be greater than polling_time (0.5)
(6, 1)
OK: boom
staging exists: False final: False
lazy ok: (2, 1)$ cd /repo && python -m pytest --doctest-modules src/MEDS_transforms/mapreduce/reducer.py src/MEDS_transforms/dataframe/write_fn.py -q 2>&1 | tail -5; python -m pytest tests -q 2>&1 | tail -3
output · 7 lines
/usr/local/lib/python3.12/site-packages/pretty_print_directory/pretty_print_directory.py:204: SyntaxWarning: invalid escape sequence '\.' >>> config = PrintConfig(ignore_regex=r"\.txt") [32m.[0m[32m [100%][0m [32m[32m[1m1 passed[0m[32m in 2.30s[0m[0m [31mFAILED[0m tests/test_pipeline.py::[1mtest_example_pipeline_parallel_stage_runner[0m - AssertionError: Pipeline failed with error: Pipeline returned code 1. [31mFAILED[0m tests/test_pipeline.py::[1mtest_example_pipeline_parallel_pipeline_cfg[0m - AssertionError: Pipeline failed with error: Pipeline returned code 1. [31m[31m[1m2 failed[0m, [32m52 passed[0m[31m in 117.58s (0:01:57)[0m[0m
$ cd /repo && git stash && python -m pytest tests/test_pipeline.py -q 2>&1 | tail -3; git stash pop
output · 13 lines
Saved working directory and index state WIP on master: 551fbd0 base [31mFAILED[0m tests/test_pipeline.py::[1mtest_example_pipeline_parallel_stage_runner[0m - AssertionError: Pipeline failed with error: Pipeline returned code 1. [31mFAILED[0m tests/test_pipeline.py::[1mtest_example_pipeline_parallel_pipeline_cfg[0m - AssertionError: Pipeline failed with error: Pipeline returned code 1. [31m[31m[1m2 failed[0m, [32m5 passed[0m[31m in 27.92s[0m[0m On branch master Changes not staged for commit: (use "git add <file>..." to update what will be committed) (use "git restore <file>..." to discard changes in working directory) modified: src/MEDS_transforms/dataframe/write_fn.py modified: src/MEDS_transforms/mapreduce/reducer.py no changes added to commit (use "git add" and/or "git commit -a")
1 more lines
Dropped refs/stash@{0} (0307adf9d1b815a7b8a4dbf3cd45a5874e17a77b)Those two pipeline test failures are pre-existing (they fail on the unmodified base too). All checks pass.
$ done: {"notes": "Added max_poll_time polling/deadline logic with readable-parquet checks and TimeoutError/ValueError handling to reduce_over; made write_df atomic via staging file + os.replace with cleanup on failure."}output · 1 lines
submitted.