Skip to content

Add endpoint-keyed flush for batching targets - #652

Closed
royischoss wants to merge 7 commits into
mlrun:developmentfrom
royischoss:ml-12996
Closed

royischoss wants to merge 7 commits into
mlrun:developmentfrom
royischoss:ml-12996

Conversation

@royischoss

@royischoss royischoss commented Sep 10, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

Add an opt-in, endpoint-keyed flush fence to ParquetTarget. This lets MLRun wait until all buffered and already-in-flight Parquet writes accepted for one model endpoint are durable, without flushing or blocking unrelated endpoints.

Changes Made

  • Add serialized, keyword-only flush_key_field configuration and public async ParquetTarget.flush(flush_key).
  • Keep the logical endpoint key separate from physical time-partition paths, including composite keys, mixed-key batches, and endpoints spanning multiple partitions.
  • Serialize concurrent same-key flushes, shield writes from caller cancellation, retain terminal write failures, and clean up successful per-key state.
  • Reject flushes before startup and during/after termination, and aggregate keyed write/cleanup failures deterministically while preserving legacy termination behavior for other batching targets.
  • Always forward the termination signal downstream, on both the failure and the cancellation path, so steps after the target are torn down and cannot leak file handles or connection pools.
  • Re-check for newly buffered batches after every await during tracked termination, matching the existing _emit_all() loop.
  • Distinguish cancelled writes with BatchWriteCancelledError, report the real batch key and sequence in that message, and safely retain exceptions that cannot be copied.
  • Wait across repeated cancellation during termination so cleanup is never left detached.
  • Name flush_key_field in the error when an event has no flush key, a wrong-typed key, or is missing the configured body field.
  • Bug fix, applies to all ParquetTarget users: move the _last_written_event update out of _blocking_emit (executor thread) into _emit (event loop). It previously advanced even when the write failed, advanced before the file was closed, and raced between concurrent executor threads. Now it advances only after a successful write and stays monotonic.

Out-of-scope change worth reviewing separately

storey/sources.py gains a _copy_exception_or_original helper and applies it to six pre-existing copy.copy call sites in AwaitableResult and AsyncAwaitableResult. Those sites had the same latent bug fixed in the batching path: copying an exception whose __init__ signature does not match its args raises TypeError and replaces the real error. This is unrelated to batching, so reviewers of the source-side result paths should look at it on its own.

Testing

  • make lint — passed
  • pytest tests/test_batching_races.py — 34 passed, stable across multiple PYTHONHASHSEED values
  • pytest tests/test_flow.py — 289 passed
  • Full local unit suite — 499 passed, 1 skipped

Related Issues

Additional Notes

  • Keyed tracking is opt-in and Parquet-only. Other batching targets expose no flush() and reject flush_key_field, because tracked mode turns serial awaited writes into concurrent tasks.
  • The fence covers only events the target has already accepted; it cannot see events still queued in an upstream source or step. MLRun must therefore drive it from inside the graph rather than from a caller that just pushed to the stream.
  • A failed write permanently poisons that logical key for the target lifetime, because a later successful write cannot prove that the lost batch became durable.
  • Once termination starts on a keyed target it runs to completion even under cancellation, so an external timeout around shutdown is advisory rather than a bound.
  • MLRun controller event semantics and retry policy remain outside Storey and will be implemented separately.

Made with Cursor

…ng the logic needed to track in_flight and done events
@royischoss
royischoss requested a review from gtopper September 10, 2026 10:45
@royischoss royischoss added the safe-to-test Green light for integration tests to run label Sep 10, 2026
@github-actions github-actions Bot removed the safe-to-test Green light for integration tests to run label Sep 14, 2026
@royischoss royischoss added the safe-to-test Green light for integration tests to run label Sep 14, 2026
Co-authored-by: Cursor <cursoragent@cursor.com>
@github-actions github-actions Bot removed the safe-to-test Green light for integration tests to run label Sep 14, 2026
royischoss and others added 3 commits September 14, 2026 16:11
Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
@royischoss royischoss added the safe-to-test Green light for integration tests to run label Sep 14, 2026
Co-authored-by: Cursor <cursoragent@cursor.com>
@github-actions github-actions Bot removed the safe-to-test Green light for integration tests to run label Sep 15, 2026
@davesh0812 davesh0812 added the safe-to-test Green light for integration tests to run label Sep 15, 2026
@royischoss royischoss closed this Sep 16, 2026
@royischoss

Copy link
Copy Markdown
Collaborator Author

differed to next version of mlrun

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

Labels

safe-to-test Green light for integration tests to run

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants