# Training pipelines

A dataset is already a stream of time-ordered batches, so feeding a trainer does not need an export step. Three surfaces: PyTorch, Arrow, and the raw batch iterator underneath both.

## PyTorch

`DreamDBIterableDataset` wraps a query as a `torch.utils.data.IterableDataset`:

```python
from dreamdb.torch import DreamDBIterableDataset
from torch.utils.data import DataLoader

torch_ds = DreamDBIterableDataset(
    ds,
    field="embedding",
    query=query_vector,
    top_k=10_000,
    batch_size=64,
    where_eq={"split": "train"},
    as_numpy=True,
    transform=None,
)

loader = DataLoader(torch_ds, batch_size=None, num_workers=4)

for batch in loader:
    embeddings = batch["embedding"]      # np.ndarray (B, D)
    labels = batch["label"]              # list[str]
```

`batch_size=None` on the `DataLoader` is deliberate: DreamDB is already batching, and letting torch re-batch would shuffle work between two layers that both think they own it. `transform` runs per batch, which is the place to decode images or move tensors.

> **Note:** This section is written from the adapter's signature and source rather than from a training run — it is the one part of the Python SDK's documentation not exercised end to end while writing. The `Dataset` methods underneath it are.

## Arrow batches

For anything that speaks Arrow — Polars, DuckDB, a Parquet export, another framework:

```python
for batch in ds.iter_arrow_batches(batch_size=256, shuffle_seed=42,
                                   window_batches=256):
    print(batch.num_rows, batch.schema.names)
```

`pyarrow` is an optional dependency, and the method says so when it is missing:

```
iter_arrow_batches requires `pyarrow`. Install with `pip install pyarrow`
or use `iter_vector` for the plain-dict path.
```

`shuffle_seed` shuffles within each emitted batch, which is much cheaper than a global shuffle over object storage. General multi-field reads are windowed; `window_batches` controls how many output batches' anchors are joined at a time, while `None` restores the fully eager path.

Embeddings become Arrow `FixedSizeList<float32>` columns. Typed arrays become `FixedShapeTensor` columns carrying the declared shape and dtype instead of opaque binary values.

## The raw iterator

Both of the above sit on `iter_stream`, which walks the whole dataset in time order:

```python
for batch in ds.iter_stream(batch_size=256, fields=["image", "label"]):
    ...
```

Naming `fields` is the important optimization: an unqualified read fetches every track, and image or video tracks dominate that cost. A training loop that only needs embeddings and labels should say so.

> **Warning:** The raw `iter_stream` fast path still requires an embedding field. `fields=["label"]` — or any list without an embedding — raises `streaming iter with no embedding fields — use iter() for scalar-only`. Use `iter_scalar` for filtered scalar reads, or `iter_arrow_batches` for a general scalar, media, or typed-array projection; it falls back to the windowed join path.

`channel_buffer` controls how many batches are prefetched ahead of the consumer. Raise it when the network is the bottleneck; lower it when memory is.

## Reproducibility

Pin the data, not just the seed:

```python
version = ds.snapshot("run-2026-07-30")
train_ds = db.Dataset.open_at(version, backend=BACKEND)
```

The manifest is immutable and content-addressed, so a run that records its snapshot label can be reproduced exactly, even after months of further ingest. See [Versioning](/python-sdk-versioning.md).
