Almost every Hugging Face training script starts with load_dataset, and most people treat it as a download function. It is closer to a small database engine. It resolves files from the Hub or disk, converts them into Apache Arrow tables, memory-maps those tables so that a dataset larger than RAM opens in seconds, records every transformation under a content fingerprint so repeated runs skip work, and, in streaming mode, turns remote files into shards that PyTorch workers and training ranks can divide between them.

Knowing that architecture separates a pipeline that tokenizes a corpus once and reuses it for a month from one that silently re-tokenizes on every launch, fills a disk with orphaned cache files, or feeds two GPUs the same examples. This article explains each layer, works through tokenizing and packing a corpus, covers streaming and distributed loading, and ends with failure modes and a checklist. It describes the 4.x line of the library; how token blocks should be packed and stored in general, independent of any library, is covered in the dataloader and tokenization pipeline.

Advertisement

The architecture in one picture

The library has two data paths that share one API. The map-style path pays an up-front cost to write Arrow files into a local cache and in exchange gets cheap random access, which is what shuffling, indexing and multi-epoch training want. The streaming path writes nothing, reads remote or local files lazily and only supports sequential iteration, which is what a corpus bigger than your disk requires.

How the datasets library turns files into training batchesHub repo or local filesParquet, JSONL, CSV, imagesload_datasetresolve, download, prepareArrow files in cachememory-mapped, not loadedDataset.mapnew fingerprint, new fileindices mappingshuffle, select, splitformat or transformtorch tensors on readDataLoader workerseach re-maps the filesstreaming=TrueIterableDataset, no cacheshards per workerand per nodefilesArrowtransformrowsskip preparesplitMap-style path (top): pay once to write Arrow, then random access is almost free.Streaming path (bottom): nothing written to disk, but only sequential access.
Map-style datasets are Arrow files on disk plus a chain of fingerprinted transforms; streaming datasets are a lazy iterator over file shards.

Dataset and IterableDataset share method names such as map and shuffle, but not their semantics, and most bugs come from assuming they match.

Arrow and memory mapping

Apache Arrow is a columnar in-memory format: each column is stored as contiguous buffers, with an offsets buffer for variable-length values such as strings or token lists. The crucial property is that the on-disk Arrow IPC file has the same layout as the in-memory table. The library therefore does not read a dataset into RAM; it asks the operating system to memory-map the file, and pages are loaded only when a row is touched. A 100 GB dataset opens instantly and resident memory grows with what you actually read, not with dataset size.

Two consequences follow. DataLoader workers are cheap, because each re-maps the same files and the operating system shares the page cache between them. And contiguous reads are fast while random-order reads touch pages all over the file, and on network storage each page miss is a round trip.

That is why the indices mapping matters. Operations such as shuffle, select, train_test_split and sort do not rewrite data. They create a small indices table that says which physical row to read for each logical row. That is quick to create, but every subsequent read is a random read. If you shuffle once and then iterate many times, or run map after a shuffle, call flatten_indices() to rewrite the data in the new order, or shuffle at the loader level instead.

Advertisement

What load_dataset actually does

A call such as load_dataset("org/name", split="train") runs a resolution pipeline. It finds the repository on the Hub, reads the dataset card metadata for configurations, splits and file patterns, and picks a packaged builder by file type: Parquet, JSON Lines, CSV, text, Arrow, or folder-based loaders for images and audio. Since version 4.0, arbitrary Python loading scripts are gone and trust_remote_code is no longer supported, so a dataset is data files plus metadata, which is both safer and more predictable.

The builder downloads files into the Hugging Face cache and then prepares them: it parses each file, infers or applies Features, which is the schema, and writes Arrow shards. The default location is under ~/.cache/huggingface/datasets; set HF_HOME or HF_DATASETS_CACHE to move it, which you should do on any machine whose home directory is small. num_proc parallelises the preparation step across files, which only helps when there are several files.

For reproducibility, revision pins the repository to a commit hash so an upstream push cannot change your data between runs. For Parquet sources, columns and filters keep unneeded columns and row groups from being read at all.

The map fingerprint cache

Every Dataset carries a content fingerprint. When you call map, the library computes a new fingerprint from the previous one plus a hash of a pickled form of the function and its arguments, and writes the output to a file named after that fingerprint, cache-<fingerprint>.arrow, next to the source files. Run the script again and the fingerprint matches, so the result loads from disk instead of being recomputed.

It follows that caching fails whenever the hash is unstable. If the function captures an object that cannot be pickled deterministically, the library logs a warning that a random hash was used, and every run recomputes and writes a fresh copy. Datasets built in memory with from_dict or from_list have no cache directory, so their transforms are not cached unless you save them. And a change the hash cannot see, such as editing a file your function reads, will not invalidate the cache; pass load_from_cache_file=False or bump a version argument when that happens.

  • batched=True hands the function up to batch_size rows (default 1,000) as column lists, which is what fast tokenizers want; it also lets the function return more or fewer rows than it received, which is how packing works.
  • num_proc splits the dataset into contiguous shards, runs one process per shard and writes one cache file each; the combined result is a concatenation.
  • remove_columns drops input columns from the output, which avoids carrying raw text through every later step and shrinks the cache.

Worked example: tokenize and pack a corpus

Suppose you have 40 GB of Parquet text, about 10 million documents, and want 2,048-token training blocks. The code below tokenizes in batches across 16 processes, appends an end-of-sequence token to each document, then concatenates and cuts fixed blocks. The pack step is a batched map that returns fewer rows than it receives, so remove_columns in the first step is required: otherwise the row count mismatch raises an error.

from datasets import load_dataset
from transformers import AutoTokenizer

tok = AutoTokenizer.from_pretrained("your-org/your-tokenizer")
BLOCK = 2048

ds = load_dataset("parquet", data_files={"train": "corpus/train-*.parquet"},
                  split="train", num_proc=8)

def tokenize(batch):
    # batched=True: batch["text"] is a list of up to batch_size strings
    out = tok(batch["text"], add_special_tokens=False)
    return {"input_ids": [ids + [tok.eos_token_id] for ids in out["input_ids"]]}

def pack(batch):
    # concatenate a batch of documents and cut fixed-length blocks;
    # the tail shorter than BLOCK is dropped (a few tokens per 1,000 docs)
    flat = [t for ids in batch["input_ids"] for t in ids]
    n = len(flat) // BLOCK * BLOCK
    return {"input_ids": [flat[i:i + BLOCK] for i in range(0, n, BLOCK)]}

tokenized = ds.map(tokenize, batched=True, num_proc=16,
                   remove_columns=ds.column_names, desc="tokenize")
packed = tokenized.map(pack, batched=True, batch_size=1000, num_proc=16, desc="pack")
packed = packed.with_format("torch")
print(packed, packed.cache_files[:1])

The first run spends most of its time in the tokenizer. The second run finds both fingerprints on disk and returns in seconds. Check packed.cache_files to see which files back the dataset. The disk bill is three copies: prepared source, tokenized table and packed table, so save_to_disk the final table and delete intermediates with cleanup_cache_files(). If you need document-aware attention masks rather than plain concatenation, pack in the collator instead.

Formats and transforms

with_format("torch") returns numeric columns as tensors at read time, zero-copy where types allow, without changing stored data. with_transform registers a function that runs on every read and is never cached, which suits random augmentation but not tokenization. The rule: deterministic, expensive work goes in map and is cached; random or cheap work goes in a transform or the collator.

Streaming and shards

With streaming=True you get an IterableDataset. Nothing is prepared; transforms are recorded and applied lazily as you iterate. Its unit of parallelism is the shard, usually one data file or, after reshard(), a Parquet row group. A PyTorch DataLoader gives each worker a subset of shards, and if there are fewer shards than workers, the extra workers sit idle. You can also convert a prepared dataset with to_iterable_dataset(num_shards=64), which is faster than streaming the same data from the Hub because it reads local Arrow files.

Shuffling is approximate. shuffle(seed, buffer_size) shuffles the order of shards and then samples from a rolling buffer of buffer_size examples, 1,000 by default. If each file holds one source, a small buffer still feeds long runs of correlated data, so raise the buffer or pre-shuffle when writing files. Call set_epoch(epoch) each epoch so the effective seed becomes the initial seed plus the epoch number.

For resumable training, IterableDataset has state_dict() and load_state_dict(), which record the current shard and the example offset within it; torchdata's StatefulDataLoader uses them automatically. Resuming skips completed shards and re-reads the current one up to the offset. One caveat is documented: examples sitting in the shuffle buffer at checkpoint time are lost on resume, and the buffer refills with new data.

Splitting across ranks

datasets.distributed.split_dataset_by_node gives each training rank its own slice. For map-style datasets each rank gets a contiguous chunk. For iterable datasets the behaviour depends on arithmetic: if the number of shards is a multiple of the world size, shards are divided evenly between ranks, which is efficient; otherwise every rank reads every shard and keeps one example in world_size, which multiplies network and parsing work by the number of ranks. Shuffle with the same fixed seed on every rank before splitting, or ranks will disagree about the shard order and overlap.

import os
from datasets import load_dataset
from datasets.distributed import split_dataset_by_node
from torchdata.stateful_dataloader import StatefulDataLoader

rank, world = int(os.environ["RANK"]), int(os.environ["WORLD_SIZE"])

ds = load_dataset("your-org/web-corpus", split="train", streaming=True,
                  revision="3f2c1e0")            # pin a commit, not a branch
ds = ds.shuffle(seed=1234, buffer_size=10_000)   # same seed on every rank
ds = split_dataset_by_node(ds, rank=rank, world_size=world)
ds = ds.map(tokenize, batched=True, remove_columns=["text"])

print("shards:", ds.num_shards)                  # per rank now: want a multiple of num_workers
loader = StatefulDataLoader(ds, batch_size=8, num_workers=4)

for epoch in range(2):
    ds.set_epoch(epoch)                          # effective seed = 1234 + epoch
    for step, batch in enumerate(loader):
        train_step(batch)
        if step % 1000 == 0:
            save_checkpoint(model, opt, data_state=loader.state_dict())

For map-style datasets, a different trap applies: if eight ranks each call map on a fresh cache, all eight tokenize at once and race on the same files. Run preprocessing on the main process first and let the others load from the cache; the Transformers Trainer exposes this as the main_process_first context manager on its training arguments, and the cache must live on storage that every rank can see. The Trainer architecture article shows where that fits in a full training job.

Failure modes seen in practice

SymptomCauseFix
Every launch re-tokenizesUnhashable function gives a random fingerprintRead the hashing warning; pass objects by name, or tokenize once and save_to_disk
Disk fills up over weeksEvery changed map writes a new cache filecleanup_cache_files, a dedicated cache volume, periodic pruning
Training slower after shuffleIndices mapping turns reads into random readsflatten_indices, or shuffle in the loader
Workers idle while GPU starvesnum_shards below num_workersreshard, to_iterable_dataset with more shards, smaller files
Ranks see duplicate dataDifferent shuffle seeds per rankOne fixed seed, then split_dataset_by_node
Loss curve has periodic stepsCorrelated files and a small shuffle bufferBigger buffer, pre-shuffled files
Old script fails on 4.xDataset relied on a loading scriptUse a Parquet revision of the dataset or convert the files yourself

The last row is common in 2026: version 4.0 also replaced the Sequence feature type with List and moved audio and video decoding to TorchCodec, so pin the library version in your environment and upgrade deliberately.

Operational guidance and trade-offs

Treat the prepared dataset as a build artifact. Tokenize once with a pinned revision, a pinned tokenizer and a pinned library version, save the result with save_to_disk or push_to_hub, and record the fingerprint or commit in your experiment tracker. Training jobs then load that artifact and never re-tokenize. For cluster jobs without internet access, set HF_HUB_OFFLINE=1 so a missing file fails fast instead of hanging on the network.

Map-style gives exact shuffling, random access and cached transforms at the cost of a full local copy; streaming needs no disk but shuffles approximately and re-reads and re-transforms remote data every epoch. A common middle path is to stream a huge raw corpus once through filtering and tokenization, write the result as many Parquet shards, then train from those. Tokenizer choices that affect this pipeline are covered in the tokenizers architecture article, and how to filter, deduplicate and mix sources before any of this is covered in pretraining data filtering and mixture.

What to do next

  1. Move the cache to a large volume with HF_HOME or HF_DATASETS_CACHE and check free space before a big map.
  2. Pin revision on every load_dataset call and record it with the run.
  3. Run your tokenize map twice and confirm the second run hits the cache; if it does not, fix the hashing warning.
  4. Use batched=True, num_proc and remove_columns on every heavy map, and save the final table with save_to_disk.
  5. For streaming, print num_shards and make it a multiple of ranks times workers; reshard if needed.
  6. Shuffle with one fixed seed on every rank before split_dataset_by_node, and call set_epoch each epoch.
  7. Checkpoint the loader with StatefulDataLoader and test a kill-and-resume before a long run.
  8. Pin the datasets version and read the 4.0 changes before upgrading older scripts.
Key takeaway: The datasets library is an Arrow-backed storage engine with a content-addressed cache and a sharded streaming reader. Map-style datasets are memory-mapped files plus fingerprinted transforms, so deterministic preprocessing should run once in map and be reused, while shuffles and selections create an indices mapping that slows reads until flattened. Streaming datasets trade exact shuffling and caching for zero disk, and their parallelism is set entirely by the shard count, which must divide evenly across workers and ranks. Pin revisions, own your cache directory, make shard counts work for your cluster and test resume before you trust a long run.