Skip to content

Feat/train buffer spill - #329

Open
viktorbeck98 wants to merge 5 commits into
developmentfrom
feat/train-buffer-spill
Open

viktorbeck98 wants to merge 5 commits into
developmentfrom
feat/train-buffer-spill

Conversation

@viktorbeck98

@viktorbeck98 viktorbeck98 commented Sep 30, 2026 •

Copy link
Copy Markdown
Collaborator

Task

Description

With use_config_data_as_training=True (the default), every configure record stays in memory until training starts. NewValueComboDetector also keeps a second copy for its second configuration pass. On BGL (~1.4M rows), every detector was OOM-killed at 8 GB.

  • TrainBuffer keeps up to train_buffer_max_records records 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 of pop(0).
  • New CoreConfig fields: train_buffer_max_records and train_buffer_dir. The default dir is private per user, in the system temp dir.
  • Warnings on the first spill and on keep_configuring.
  • On local disk, files left by killed processes are removed at the next start.
  • NewValueComboDetector uses the same buffer for its configure inputs.
  • Docs note added; config tables regenerated.

Behaviour is unchanged: configure on the prefix, train on the same records in the same order, then detect.

How Has This Been Tested?

  • New unit tests for the buffer and the combo detector. The full suite passes.
  • Metrics are identical on 13 Audit benchmark cells, with and without forced spilling.
  • NewValue and NewValueCombo on BGL now finish under an 8 GB cap with no swap.

Checklist

  • This Pull-Request goes to the development branch.
  • I have successfully run prek locally.
  • I have added tests to cover my changes.
  • I have linked the issue-id to the task-description.
  • I have performed a self-review of my own code.

View with [code]smith Autofix with [code]smith
Need help on this PR? Tag @codesmith-bot with what you need. Autofix is disabled.

viktorbeck98 and others added 2 commits September 30, 2026 14:01
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>
@viktorbeck98
viktorbeck98 changed the base branch from main to development September 30, 2026 12:25
Comment thread src/detectmatelibrary/common/_core_op/_train_buffer.py Dismissed
Comment thread tests/test_common/test_train_buffer.py Outdated
import fsspec
import pytest

import detectmatelibrary.common._core_op._train_buffer as tb
Comment thread src/detectmatelibrary/common/_core_op/_train_buffer.py Dismissed
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>
@viktorbeck98
viktorbeck98 requested a review from ipmach September 30, 2026 12:59
try:
fs.rm(path, recursive=True)
except FileNotFoundError:
pass # nothing was spilled, or it is gone already

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

missing warning or something

@ipmach ipmach left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants