Feat/train buffer spill - #329
Open
viktorbeck98 wants to merge 5 commits into
Open
viktorbeck98 wants to merge 5 commits into
viktorbeck98 wants to merge 5 commits into
Conversation
With use_config_data_as_training=True, CoreComponent kept every configure record in memory until training started, and replayed them with an O(N^2) list.pop(0). On BGL and HDFS this alone passed an 8 GB cap. TrainBuffer (now in _core_op/_train_buffer.py, re-exported from core.py) keeps up to train_buffer_max_records records in memory (default 100000), then writes each full batch as one closed Parquet part file through fsspec and streams the parts back in order when training starts. Order and semantics are unchanged: configure on the prefix, train on the same records, then detect. - New CoreConfig fields train_buffer_max_records and train_buffer_dir. - Default dir is a private per-user folder in the system temp dir (/tmp/detectmatelibrary-<uid>/train_buffer, 0700, owner-checked). - On local disk each run holds an flock; a later component start removes run dirs of killed processes (only the part files the buffer wrote). - Warnings on the first spill and on keep_configuring; docs note and regenerated config tables. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
NewValueComboDetector kept every configure record in a plain list for its second configuration pass, so memory still grew with data_use_configure after the train buffer learned to spill (OOM at ~1.38M BGL rows under an 8 GB cap). The list is now a TrainBuffer with the same train_buffer_max_records / train_buffer_dir settings; set_configuration replays it and deletes its files. TrainBuffer takes an optional reason for its first-spill warning, so the combo buffer does not advise setting use_config_data_as_training=False. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| import fsspec | ||
| import pytest | ||
|
|
||
| import detectmatelibrary.common._core_op._train_buffer as tb |
Addresses the CodeQL review on #329: if locking or creating the run directory failed, the open lock file leaked and a stray .lock/.lock.tmp stayed on disk. Also comment the empty except and import the buffer module one way in its tests. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
ipmach
reviewed
Oct 1, 2026
| try: | ||
| fs.rm(path, recursive=True) | ||
| except FileNotFoundError: | ||
| pass # nothing was spilled, or it is gone already |
Contributor
There was a problem hiding this comment.
missing warning or something
ipmach
requested changes
Oct 1, 2026
ipmach
left a comment
Contributor
There was a problem hiding this comment.
I found the Trainbuffer not that readable, it feels it has to many variables in init and hard to follow.
Question/Suggestion, why can we not use a Polars DataFrame with one column of binaries (each binary is a schema)? The logical will be the same as this, Polars use arrow in the backend and save in the same format but probably the syntax will be more cleaner.
Otherwise, if this work for HDFS and BGL then sounds good to me
|
|
||
| def _spill(self) -> None: | ||
| import pyarrow as pa | ||
| import pyarrow.parquet as pq |
Contributor
There was a problem hiding this comment.
imports should be in top of the document
Three follow-ups to the train buffer spill (review on #47): - Sliding windows: a spilled window that moves on by one record from the window before it stores only its new last record; replay rebuilds it from its predecessor. window=10 spill+replay drops from 55.8 to 8.3 us per record (200k records); single records are unchanged. - No disk access before the first spill: the stale-run cleanup moves from TrainBuffer construction to the first spill, after the default directory's ownership check. Components that never spill never scan. - NewValueComboDetector keeps one copy of its configure records. With use_config_data_as_training=True, process() now keeps the record for training before calling configure(), which checks _kept_for_training() and skips its own copy; the second configuration pass reads buffer_train with the new non-destructive peek() and leaves it for training. Direct configure() calls and use_config_data_as_training=False still use the detector's own buffer.
Review feedback on #329 found TrainBuffer hard to follow: too much state in __init__ and too many helpers. Most of that came from the directory safety and flock-based stale-run cleanup, so drop them and spill to a tempfile.TemporaryDirectory. It is private (0700), unique, and removed after replay, on garbage collection or at interpreter exit. A killed process leaves its detectmate-train-* directory behind; the docs say so. - TrainBuffer(schema_class, max_records, spill_dir, window, name): add(), iterate to replay in order (then empties), clear(). 402 -> 122 lines. - Local paths only: no fsspec, flock, ownership checks, back-off, peek/last/len/__add__ or `why`. pyarrow imports at module top. - Window mode still writes each log once; the window size now comes from the component's DataBuffer instead of being detected per record. - An empty `except: pass` on cleanup is now a warning naming the directory to delete by hand. - NewValueComboDetector keeps its own copy of the configure records again, instead of reading buffer_train through _kept_for_training/peek. This removes hidden coupling; it only costs disk, not RAM. - process() configures before buffering again, as on development.
This branch has not been deployed
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.
Task
Description
With
use_config_data_as_training=True(the default), every configure record stays in memory until training starts.NewValueComboDetectoralso keeps a second copy for its second configuration pass. On BGL (~1.4M rows), every detector was OOM-killed at 8 GB.TrainBufferkeeps up totrain_buffer_max_recordsrecords in memory (default 100 000). Beyond that it writes Parquet parts via fsspec and streams them back into training. Replay is now O(N) instead ofpop(0).CoreConfigfields:train_buffer_max_recordsandtrain_buffer_dir. The default dir is private per user, in the system temp dir.keep_configuring.NewValueComboDetectoruses the same buffer for its configure inputs.Behaviour is unchanged: configure on the prefix, train on the same records in the same order, then detect.
How Has This Been Tested?
Checklist
Need help on this PR? Tag
@codesmith-botwith what you need. Autofix is disabled.