Conversation
… remove unecessary dependencies
| not pathlib.Path(record["output_file"]).is_file() | ||
| or pathlib.Path(record["output_file"]).stat().st_size == 0 | ||
| ): | ||
| return True |
There was a problem hiding this comment.
sorry for the long comment - but this re-analyzes the overall resume functionality in depth and has some broader suggestions.
firstly: the resume path here can't ever fire, and since the plan is to build parallelization on top of this eventually I'd rather we fix the foundation in this PR.
_is_new returns True when the output file is missing or empty here. but that's the normal state of an interrupted run, so it fights _merge_needed, and the pipeline checks _is_new first (line 494) and rmtree's the directory. I set up a crashed-run state - matching .config.json, three .done files, no output yet - and ran the pipeline:
_is_done -> False
_merge_needed -> True <- "resume: just merge"
_is_new -> True <- "throw it away and restart"
[INFO] New folder, creating directory and chunking
pre-existing .done files that survived: NONE
so it's either "already done, skip" or "start from scratch" - the resume branch is dead. output is still correct, it just always redoes the whole file. and test_is_not_new_when_config_matches is the hidden test that actually fails once you make the class collectible: it was telling us this.
four things I'd do, roughly in order of value:
1. give _is_new one job. it's answering two questions - "are this directory's parameters still valid" and "is the final output already there". only the first belongs to it; the second is _is_done's. then order the checks so resume wins:
if _is_done(output_dir, record): # output exists and params match -> nothing to do
return
if not _params_match(output_dir, record): # config/batch_size/output path/input changed
reset_run_dir(output_dir, record)
# otherwise: pick up whatever units are already .done and finish the rest2. make "done" atomic. transform_chunk (line 417) opens <chunk>.done and writes at the very end, so a crash between the open and the write leaves a zero-length .done that later looks complete. write <chunk>.done.tmp and os.replace() it into place - atomic within a filesystem, so a .done that exists is always whole. right now that's an unlikely race; with W workers it becomes routine, and every resume decision is built on trusting these files.
3. don't hold a unit in memory if we don't need to. transform_chunk accumulates the whole chunk in a bytearray before one fout.write. writing per record should let the buffered file object do the batching itself
4. consider dropping the chunk files while keeping the chunk abstraction. a work unit doesn't have to be a file - it can be a (start_offset, end_offset) byte range over the input, split on newline boundaries, with each worker reading its own slice. that kills the serial split_file pre-pass (today nothing gets tagged until the entire input has been copied to scratch), takes peak scratch from ~2x input down to 1x, and makes unit boundaries a recomputable function of the input instead of on-disk state that can be half-written. this is the biggest of the four and the least urgent - 1 through 3 matter more.
also: docstrings say "Stream-transform" but this isn't truly streaming so I'd remove that.
Co-authored-by: Alexander VanTol <Avantol13@users.noreply.github.com>
Co-authored-by: Alexander VanTol <Avantol13@users.noreply.github.com>
Co-authored-by: Alexander VanTol <Avantol13@users.noreply.github.com>
Co-authored-by: Alexander VanTol <Avantol13@users.noreply.github.com>
… by including chunk and done files in fhir_input folder and give each test its out subfolder in fhir_outputs
New Features
Dependency updates
Breaking Changes