Skip to content
heliaEDGE
User guide
HELIA

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.

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 np
from 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 on workers. A seed is required whenever you shuffle or transform.
  • Training epochs. Iterating the dataset again replays the same elements, so a Keras fit that reads it once per epoch sees the same order and augmentation every epoch. Give fit one stream instead: num_epochs=None with steps_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 a tf.data.Dataset. Use None for 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 work workers=0 is faster; measure with your own transform.

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 np
import tensorflow as tf
from 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.

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.
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 partial
from 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.