# 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

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:

```python
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.

## 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:

```python
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.

## 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

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:

```python
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:

```python
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.

## Continue with a complete workflow

[Data API](https://ambiqai.github.io/helia-edge/reference/api/helia_edge/data/)[Choose transforms](https://ambiqai.github.io/helia-edge/guide/preprocessing/)[CIFAR-10 training notebook](https://ambiqai.github.io/helia-edge/examples/train-cifar-model/)
