bridge.models.bagel.data.order#

Module Contents#

Classes#

_RestorableDataset

BagelPlannedLoader

Read an independently planned repeating source stream by Energon restore key.

Functions#

_repeat_to_length

Match BAGEL’s repeated file selection up to num_used_data.

shuffle_and_shard

Apply BAGEL’s sorted epoch shuffle and rank/worker floor sharding.

plan_t2i_sources

Plan stable BAGEL T2I source IDs in file, row-group, and row order.

plan_editing_sources

Plan stable BAGEL Editing source IDs in shuffled row-group order.

plan_vlm_sources

Plan stable BAGEL VLM source IDs through line and epoch shuffles.

plan_manifest_indices

Map BAGEL’s independently planned source stream to physical WDS sample indices.

Data#

T

API#

bridge.models.bagel.data.order.T#

‘TypeVar(…)’

class bridge.models.bagel.data.order._RestorableDataset#

Bases: typing.Protocol

restore_sample(key: tuple[str | int, ...]) → object#
bridge.models.bagel.data.order._repeat_to_length(
items: collections.abc.Sequence[bridge.models.bagel.data.order.T],
length: int,
) → list[bridge.models.bagel.data.order.T]#

Match BAGEL’s repeated file selection up to num_used_data.

bridge.models.bagel.data.order.shuffle_and_shard(
items: collections.abc.Sequence[bridge.models.bagel.data.order.T],
*,
seed: int,
rank: int,
world_size: int,
worker_id: int = 0,
num_workers: int = 0,
sort_key: collections.abc.Callable[[bridge.models.bagel.data.order.T], object] | None = None,
) → list[bridge.models.bagel.data.order.T]#

Apply BAGEL’s sorted epoch shuffle and rank/worker floor sharding.

bridge.models.bagel.data.order.plan_t2i_sources(
parquet_paths: collections.abc.Sequence[str],
row_counts: collections.abc.Mapping[str, collections.abc.Sequence[int]],
**shard_args: int,
) → list[dict[str, object]]#

Plan stable BAGEL T2I source IDs in file, row-group, and row order.

bridge.models.bagel.data.order.plan_editing_sources(
row_groups: collections.abc.Sequence[tuple[str, int]],
row_counts: collections.abc.Mapping[tuple[str, int], int],
**shard_args: int,
) → list[dict[str, object]]#

Plan stable BAGEL Editing source IDs in shuffled row-group order.

bridge.models.bagel.data.order.plan_vlm_sources(
lines: collections.abc.Sequence[str],
*,
jsonl_name: str,
num_used_data: int,
shuffle_seed: int,
**shard_args: int,
) → list[dict[str, object]]#

Plan stable BAGEL VLM source IDs through line and epoch shuffles.

bridge.models.bagel.data.order.plan_manifest_indices(
manifest_path: pathlib.Path,
*,
seed: int,
rank: int,
world_size: int,
worker_id: int,
num_workers: int,
num_used_data: int,
shuffle_seed: int = 0,
) → list[int]#

Map BAGEL’s independently planned source stream to physical WDS sample indices.

class bridge.models.bagel.data.order.BagelPlannedLoader(
dataset: bridge.models.bagel.data.order._RestorableDataset,
sample_indices: collections.abc.Sequence[int],
worker_config: megatron.energon.WorkerConfig,
)#

Bases: collections.abc.Iterator[object]

Read an independently planned repeating source stream by Energon restore key.

Initialization

Store the random-access Energon dataset and canonical source plan.

__iter__() → Self#

Return this stateful source reader.

__next__() → object#

Cook the next planned WDS sample through Energon’s restore API.

save_state_rank() → dict[str, int]#

Save this worker’s source-stream position.

restore_state_rank(
state: collections.abc.Mapping[str, object],
) → None#

Restore this worker’s source-stream position.