Datasets

Many files, or a Hugging Face Hub dataset, as one lazy table

marrow.datasets.load_dataset is shaped like the datasets function of the same name. It resolves one split of a dataset to the files it is made of and returns a lazy table scanning all of them in order. Building the table reads the first file’s schema and nothing else; the data is read when the plan runs, so filters and column pruning apply to every file.

Files of one format

Name the format — "parquet", "json" (newline-delimited) or "arrow" (the IPC file format) — and the files with data_files. A pattern or a list of them is the train split, as in datasets; a mapping names each split.

import os, tempfile
import marrow.parquet as pq
from marrow.datasets import load_dataset

data = tempfile.mkdtemp()
for shard, start in enumerate([0, 3, 6]):
    pq.write_table(
        ma.table({"id": ma.array([start, start + 1, start + 2])}),
        os.path.join(data, f"train-{shard}.parquet"),
    )
pq.write_table(ma.table({"id": ma.array([100])}), os.path.join(data, "test-0.parquet"))

train = load_dataset("parquet", data_files=os.path.join(data, "train-*.parquet"))
print(train.collect().column("id").to_pylist())
[0, 1, 2, 3, 4, 5, 6, 7, 8]
files = {
    "train": os.path.join(data, "train-*.parquet"),
    "test": os.path.join(data, "test-*.parquet"),
}
print(load_dataset("parquet", data_files=files, split="test").collect().num_rows)
1

The result is an ordinary plan over one scan of every file, so the verbs, the optimizer and the engine treat it like any other:

query = train.filter(col("id") >= 5).optimize()
print(query.explain())
print(query.collect().column("id").to_pylist())
Filter(ParquetScan(/tmp/tmpu4kaxumv/train-0.parquet, /tmp/tmpu4kaxumv/train-1.parquet, /tmp/tmpu4kaxumv/train-2.parquet) pruned by 1, ge(id, 5))
[5, 6, 7, 8]

Patterns

A pattern is a path or a URI whose last segments may hold wildcards, matched the way datasets and fsspec match data_files:

Wildcard Matches
* any run of characters within one path segment
** any run of segments, none included — **/*.parquet reaches the top level too
? one character other than /, in a local path
\ makes the next character literal

A local path is listed from the filesystem; any other URI (hf://, s3://, gs://, …) through its OpenDAL service, which must be able to list — http(s):// cannot, so name each file there. In a URI ? starts the query, as it does in any URI, and every file found carries the query on. A pattern without wildcards names one file and needs no listing.

Every file must have the first file’s columns; a file that does not is named in the error when the scan reaches it.

Hugging Face Hub datasets

Pass a dataset’s owner/name instead of a format. Its configs and their splits come from the dataset card’s configs: block, or, when the card has none, from the file names the way datasets infers them (train, validation/valid/dev/val, test/eval):

from marrow.datasets import (
    get_dataset_config_names,
    get_dataset_split_names,
    load_dataset,
)

get_dataset_config_names("openai/gsm8k")          # ['main', 'socratic']
get_dataset_split_names("openai/gsm8k", "main")   # ['train', 'test']

gsm = load_dataset("openai/gsm8k", "main", split="test")
gsm.filter(col("question").char_length() > 300).select("question").collect()

A split stored as Parquet, JSON Lines or Arrow is read from its own files, through hf://, fetching only the column chunks the query needs. One stored otherwise — CSV, a .json document, compressed files, images — or with no data files at all is read from the Parquet copy the Hub converts every public dataset into, which exists for the default branch:

iris = load_dataset("scikit-learn/iris")           # Iris.csv, read as Parquet
iris.collect().num_rows                            # 150

data_files works for a Hub dataset too, its patterns relative to the repository. revision= picks a branch, a tag or a commit; a revision holding a /, such as refs/pr/1, cannot be named, so use its commit. HF_TOKEN is sent when it is set, which private and gated datasets need.

Note

Reading from the Hub, or any object store, needs libopendal_c; a local path never does. See the opendal environment in pixi.toml.

In Mojo

The same entry point returns a DynRelation. HubDataset is what the Hub says about a dataset, DataFiles the patterns of each split, and the patterns themselves are marrow.io.Glob, usable on their own to match or expand paths.

docs/snippets/load_dataset.mojo
"""A Hub dataset's split and a glob of local files, each as a plan.

Compiled by `pixi run docs_check`; included by docs/guide/datasets.qmd.
"""

from marrow.datasets import DataFiles, HubDataset, load_dataset
from marrow.dtypes import int64
from marrow.expr import col, lit


def main() raises:
    # What the Hub says about a dataset: its configs, and each one's splits.
    var hub = HubDataset.fetch("openai/gsm8k")
    print(hub)

    # One split of one config, as a scan over every file in it.
    var gsm = load_dataset("openai/gsm8k", "main", split="test")
    print(gsm.limit(3).execute())

    # Files of one format, named by a glob.
    var orders = load_dataset(
        "parquet", DataFiles(["data/orders-*.parquet"])
    ).filter(col("amount", int64) > lit(100, int64))
    print(orders.execute())

Limits

  • A remote Parquet scan fetches one row group per request, one after another, so a file written with many small row groups reads slowly over the network.
  • A config is resolved when the plan is built; it is not a parameter of the plan.
  • Hive-style key=value/ directories do not become columns.
Back to top