nemo_automodel.components.speculative.streaming.loader
nemo_automodel.components.speculative.streaming.loader
Streaming consumer for speculative-decoding draft training.
FeatureDataLoader is a Python iterator over
Eagle3TargetBatch. It pulls one SampleRef lease at a
time from a SampleRefQueue, materializes the tensors through a
FeatureStore, hands the trainer a fresh
Eagle3TargetBatch, and releases the previous lease on every
__next__ call — so the trainer can hold one batch across one full
forward without it being freed mid-forward.
The loader is per-rank and the leases stay on-rank; the queue + store live in the same Python process as the consumer. FSDP / CP / EP parallelism lives inside the trainer’s forward / backward and is unaffected by the loader’s lifecycle.
Module Contents
Classes
Functions
Data
API
Iterator over Eagle3TargetBatch materialized from a streaming queue.
Lifecycle
Each iterator pull yields an Eagle3TargetBatch whose
tensors come from a fresh store.get(). The previous
batch’s lease is ack’d and its store handle released on the
NEXT pull — so the trainer can hold one batch across one
forward pass without it being freed mid-forward. A consumer
that wants eager reclamation (e.g. to free memory before
pulling the next batch) calls consume_now after
computing its loss. Iteration ends on close or when
the queue drains.
Parameters:
The metadata-only queue the producer puts SampleRef
onto. queue.close() at any time cuts the iterator short.
The FeatureStore each lease will be materialized
through. Must match ref.store_uri for the leased refs.
FeatureAlgorithm the loader runs the
per-algorithm schema check for; defaults to EAGLE-3.
Ack the most recent lease / release its store handle and shut the queue.
After close the iterator raises StopIteration on
the next pull. Idempotent so a trainer can use with safely.
Release the most recently yielded batch eagerly.
Useful for trainer hooks that want to free memory before
pulling the next batch (e.g. immediately after backward()).
Idempotent.
Build an Eagle3TargetBatch from the store’s tensors.
Pulled into a module-level helper so the per-algorithm batch-building logic lives in one place and adding DFlash / DSpark only requires a new branch here.