nemo_rl.data.energon.sft_worker#

Colocated Energon loader extension for Megatron policy workers.

Module Contents#

Classes#

SFTMegatronPolicyWorker

Megatron policy worker with an Energon loader on each DP owner.

Data#

API#

class nemo_rl.data.energon.sft_worker.SFTMegatronPolicyWorker(
*args: Any,
processor: Any = None,
**kwargs: Any,
)#

Bases: nemo_rl.models.policy.workers.megatron_policy_worker.MegatronPolicyWorkerImpl

Megatron policy worker with an Energon loader on each DP owner.

Initialization

Initialize the MegatronPolicyWorker.

setup_sft_dataloader(
*,
data_config: Mapping[str, Any],
batch_size: int,
max_sequence_length: int,
placement_fingerprint: str,
packing_algorithm: str | None,
max_sequences_per_bin: int | None,
sequence_length_pad_multiple: int,
only_unmask_final: bool,
restored_state: Optional[dict[str, Any]] = None,
) bool#

Build the train loader on the TP0/PP0/CP0 rank of this DP replica.

load_next_sft_batch(
*,
only_unmask_final: bool,
make_sequence_length_divisible_by: int,
) nemo_rl.data.energon.sft_types.StepEnvelope#

Load, prepare, and publish one batch into this process’s local store.

commit_sft_batch() None#

Release the active process-local batch after a successful step.

abort_sft_batch() None#

Release the active batch after a failed policy step.

sft_dataloader_state_dict() dict[str, Any]#

Capture this logical loader state after its batch is committed.

close_sft_dataloader() None#

Clear local batch state and release the loader reference.

_require_active_envelope() nemo_rl.data.energon.sft_types.StepEnvelope#
static _source_ids(
batch: Mapping[str, Any],
*,
batch_size: int,
) tuple[str, ...]#
static _source_tags(
batch: Mapping[str, Any],
*,
batch_size: int,
) list[dict[str, Any]]#
nemo_rl.data.energon.sft_worker.__all__#

[‘SFTMegatronPolicyWorker’, ‘StepEnvelope’]