Input pipelines
helia_edge.data reads indexed records with Grain and hands the batches to either backend: to_tf_dataset for TensorFlow, to_torch_loader for Torch. The generator helpers further down adapt a Python reader to a TensorFlow dataset while keeping control of which records are read and how often; they need the TensorFlow backend.
Read records with Grain
Section titled “Read records with Grain”Install helia-edge[grain]. Any object with len() and integer indexing is a source, including a list. The transform runs per record, in the worker processes when workers is above zero, and receives a random generator derived from the seed:
import numpy as npfrom helia_edge.data import to_grain
records = [{"x": np.full((8,), index, dtype=np.float32), "y": np.int32(index % 2)} for index in range(10)]
def add_noise(record, rng): return {"x": record["x"] + rng.normal(size=8).astype(np.float32), "y": record["y"]}
dataset = to_grain(records, seed=0, shuffle=True, transform=add_noise, batch_size=4, num_epochs=2)print([batch["x"].shape for batch in dataset])The shapes are five batches of (4, 8): batches run across epoch boundaries, so the 20 records of two epochs fill five batches. With num_epochs=1 the last batch would be (2, 8); set drop_remainder=True to drop a partial last batch. Each of the num_epochs passes over the records is a new permutation with new transform draws.
- Determinism. The record order and every transform draw depend only on
seed, not onworkers. A seed is required whenever you shuffle or transform. - Training epochs. Iterating the dataset again replays the same elements, so a Keras
fitthat reads it once per epoch sees the same order and augmentation every epoch. Givefitone stream instead:num_epochs=Nonewithsteps_per_epoch=len(records) // batch_size, so each Keras epoch continues where the last one stopped. - Torch.
to_torch_loader(dataset)yields the same batches as tensors, keeping dict and tuple structure. It adds no batching or workers of its own, and it copies arrays out of worker shared memory. - TensorFlow.
to_tf_dataset(dataset, output_signature)yields the same batches as atf.data.Dataset. UseNonefor the batch dimension when the last batch can be smaller. - Parallelism. Put heavy per-record work in the transform and raise
workers: it then runs in separate processes, outside the training process’s GIL, and only finished batches reach the training process. Workers add inter-process overhead, so for light per-record workworkers=0is faster; measure with your own transform.
Build a finite dataset
Section titled “Build a finite dataset”The ID generator chooses records. The data generator reads them. spec describes one unbatched sample, including its shape and dtype. This small example uses synthetic records so you can run it without a dataset:
import numpy as npimport tensorflow as tffrom helia_edge.utils import create_interleaved_dataset_from_generator
records = {index: np.full((8, 1), index, dtype=np.float32) for index in range(3)}
def read_records(record_ids): for record_id in record_ids: yield records[record_id]
dataset = create_interleaved_dataset_from_generator( data_generator=read_records, id_generator=iter, ids=list(records), spec=tf.TensorSpec(shape=(8, 1), dtype=tf.float32), stream_mode="global",)batches = list(dataset.batch(2).as_numpy_iterator())print([batch.shape for batch in batches])The shapes are [(2, 8, 1), (1, 8, 1)]. Replace read_records with your reader and describe its actual output in spec. For supervised training, yield (inputs, targets) and supply a matching pair of tensor specifications. Batching, shuffling and prefetching remain explicit dataset operations.
Choose the sampling mode
Section titled “Choose the sampling mode”| Mode | Choose it when | Ordering |
|---|---|---|
global |
You own a single finite, repeated or weighted schedule. | Preserves that schedule; num_workers does not partition it. |
finite |
Each reader terminates and partitioning IDs does not change what it yields. | Deterministic mode concatenates partitions; otherwise they can interleave. |
Go further
Section titled “Go further”Read independent finite partitions
Reuse the reader and tensor specification above. Choose finite only when every reader terminates and splitting IDs does not affect the samples it yields:
partitioned = create_interleaved_dataset_from_generator( read_records, iter, list(records), tf.TensorSpec(shape=(8, 1), dtype=tf.float32), stream_mode="finite", num_workers=2, deterministic=True,)Deterministic mode concatenates contiguous partitions in order. num_workers counts generators in the same process, not worker processes. Set deterministic=False only when interleaving is acceptable.
Use a weighted, repeating schedule
Keep repeating schedules in global mode. Take a bounded sample while inspecting the reader, and set an explicit epoch length when training:
from functools import partialfrom helia_edge.utils import random_id_generator
weighted = create_interleaved_dataset_from_generator( read_records, partial(random_id_generator, weights=[1, 1, 2]), list(records), tf.TensorSpec(shape=(8, 1), dtype=tf.float32), stream_mode="global",)preview = list(weighted.take(6).as_numpy_iterator())Sampling is with replacement. Weights are relative, not exact per-epoch quotas; their count must match the IDs and their finite total must be positive. Use .batch(...) and steps_per_epoch for an unbounded training stream.
Scheduling contracts and compatibility with earlier experiments
create_interleaved_dataset_from_generator defaults to stream_mode="global":
one caller-owned ID/sample schedule, including repeated or weighted schedules.
num_workers does not partition this mode. This intentionally changes the old
implicit partitioning, which dropped remainder IDs and could oversample subjects
in smaller repeated partitions. Existing frozen experiments must retain their
pinned version or declare a new loader/order policy when migrating.
Use StreamMode.GLOBAL or StreamMode.FINITE from helia_edge.utils for typed
configuration. Their string values remain accepted at the API boundary.
Use stream_mode="finite" only when each partition terminates and splitting IDs
does not change the caller’s generator semantics. Every ID is assigned once;
deterministic=True concatenates contiguous partitions in order. Set it to False
to interleave finite partitions when order is unimportant. Full coverage applies
to complete enumerations; arbitrary early stopping does not preserve the same
sampling distribution. This is in-process Python generation, not multiprocessing.
Neither mode changes batching, shuffling, labels, timestamps or subject/window weights. Empty inputs produce a typed empty dataset without invoking readers. Omitted preprocessing is an identity operation. Invalid worker counts fail early. Generators receive fresh ID lists for each epoch, so in-place shuffling cannot mutate the caller’s schedule.
random_id_generator(ids, weights=...) samples IDs with replacement using
nonnegative relative weights; omitted weights preserve uniform sampling. Weights
must match the IDs and have a finite positive total. Earlier versions ignored
supplied weights, so adopting this fix changes those weighted experiments.