DreamDB

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.

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.

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.