Read the tile WCS header log lock-free in make_post_process - #938
Merged
Merged
Conversation
make_post_process looked up the header log once per (exposure, CCD) through SqliteDict: two sqlite transactions per lookup, each taking a POSIX lock on NFS, and each lookup unpickling all of the exposure's WCS again. Load the store once with a new lock-free reader (sqlite immutable=1) and index the resulting dict; the membership check moves out of the CCD loop. shapepipe.pipeline.sqlite_store provides ImmutableSqliteDict, a read-only Mapping with SqliteDict's keys, order and decoded values, and read_sqlitedict, which loads a whole store into a dict. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
immutable=1 skips hot-journal recovery, so a store whose writer died mid-transaction would read back its uncommitted rows. ImmutableSqliteDict now raises when a -journal or -wal file sits next to the store. The module docstring scopes its no-concurrent-writer precondition to the sites that use the reader. Tests cover the guard, a real killed writer, and a path with URI special characters. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
tile_detect's post-processing (make_post_processinsextractor_script.py) builds the per-epochEXP_NAME/CCD_NHDUs andN_EPOCHby looking up the tile's header log (log_exp_headers<tile>.sqlite, written bymerge_headers) once per exposure × CCD:That is ~80 transactions per exposure, and each takes and releases a POSIX lock on the file. On NFS
/scratchthose locks are slow:?mode=ro(locking) took 18.4 s, 80–250 ms per query on a busy store. With?immutable=1(no locking) the same reads took 0.25 s. Some freshly written stores took 27–30 s per locked query, or never returned.sqlitedict.select←__getitem__/__contains__←make_post_process, and the SqliteDict worker thread in D state onrpc_wait_bit_killable.Reading each exposure once would still take one lock per exposure, which costs ~30 s on a bad file, so the read has to avoid locks entirely. The loop also re-unpickles every exposure's WCS list once per CCD. That cost is separate from locking and is not small: see the numbers below.
Change
shapepipe.pipeline.sqlite_store:ImmutableSqliteDictis a read-onlyMappingover a SqliteDict file, opened withfile:<path>?immutable=1. Its keys,rowidorder and values (decoded with sqlitedict's owndecode) match whatSqliteDictreturns for stores written with SqliteDict's default encoding, which is how every ShapePipe store is written.read_sqlitedict(path)loads a whole store into a dict.FileNotFoundError.SqliteDictwould silently create an empty store instead.make_post_processloads the header log once withread_sqlitedict, then indexes the dict. The "exposure missing from header file" check now runs once per exposure, outside the CCD loop. Output is unchanged.immutable=1is only correct when no process writes the file during the read. sqlite neither sees concurrent changes nor replays a hot journal. This holds here becausemerge_headerswrites and closes the header log in an earlier rule. As a guard against a writer that is still running or was killed mid-transaction,ImmutableSqliteDictrefuses to open a store with a-journalor-walfile next to it and raisesRuntimeError. Without the guard, a reviewer reproducedimmutable=1returning 49 of 50 rows with uncommitted values, whereSqliteDictrolled the journal back.Verification
tests/module/test_sqlite_store.py:merge_headersfrom split_exp-style WCS.npyfiles,read_sqlitedictreturns the same keys, in the same order (TILE_IDfirst), with values whose pickles are byte-identical toSqliteDict's.Mappingsemantics againstSqliteDict: lookup,in,len,KeyError, androwidorder after an overwrite.-journalor-walfile next to the store makes the read raise. One case uses a realSqliteDictwriter killed mid-transaction: the read raises, andSqliteDictthen rolls back to the committed row.?,#and%reads correctly through thefile:URI.tests/module/test_sextractor_post_process.pygains a test that an exposure missing from the header log still raisesKeyError. The existing tests pass.log_exp_headers-202-301.sqlite(9.7 MB, 6 exposures × 40 CCDs, NFS/scratch, login node):read_sqlitedictmatchesSqliteDict(flag="r").items()exactly: same key order, byte-identical pickles.Timings:
SqliteDict(482 queries)read_sqlitedict/scratch): develop'smake_post_processtook 78.6 s and this branch's 7.1 s. The output is identical across all 9 HDUs (PRIMARY,LDAC_IMHEAD,LDAC_OBJECTS,EPOCH_0–EPOCH_5), andEPOCH_*/N_EPOCHequal the campaign's own sexcat. Full suite on the branch: 835 passed, 6 skipped.Scope
This PR only touches
make_post_process, where the time is going. Other code reads SqliteDict stores in the same per-key way, and each could switch toImmutableSqliteDictas a drop-in in a follow-up. None of these are converted here, to keep the diff away from files that open PRs are reworking:ngmix.Vignet(6–7 vignet stores read per object) andngmix_runner's pre-check, which scans every key of the PSF and galaxy vignet stores.make_cat.save_psf_data, which does one read ofgalaxy_psf_cat[str(id)]per object (~35K per tile).make_cat.pyis touched by Wire the UNIONS external healsparse masks into the workflow (MASK_<flag value>_<name> columns) #886/Exposure-level HealSparse maps: instrument defects and exposure count #887.f_wcsreads invignetmaker,psfex_interpandmccd_interpolation_script.git merge-treeagainst #933, #925 and #886 is clean. Against #887 it conflicts only inworkflow/README.md, and #887 conflicts there with develop as well. On #933,make_post_processreads the header log exactly as on develop, so this change still applies after #933 lands.Claude Opus 5.5 on behalf of Cail
🤖 Generated with Claude Code