zephon.io

Public IO API: datasets, the in-memory shard primitive, and store options.

Configure store behavior via Pipeline.options(io_options=StoreOptions(...)).

class CacheOptions(enabled=_NotSupplied.TOKEN, root=_NotSupplied.TOKEN, limit_bytes=_NotSupplied.TOKEN, rg_cache_bytes=_NotSupplied.TOKEN, keep_zip=_NotSupplied.TOKEN, validate_hash=_NotSupplied.TOKEN, download_retry=_NotSupplied.TOKEN, download_timeout=_NotSupplied.TOKEN, open_retry_attempts=_NotSupplied.TOKEN, open_retry_initial_backoff=_NotSupplied.TOKEN, open_retry_max_backoff=_NotSupplied.TOKEN, min_slack_bytes=_NotSupplied.TOKEN, max_slack_bytes=_NotSupplied.TOKEN)[source]

Bases: object

User-configurable knobs that control cache behaviour.

enabled: bool = 0
root: str | Path = 0
limit_bytes: int | None = 0
rg_cache_bytes: int | None = 0
keep_zip: bool = 0
validate_hash: str | None = 0
download_retry: int = 0
download_timeout: float = 0
open_retry_attempts: int = 0
open_retry_initial_backoff: float = 0
open_retry_max_backoff: float = 0
min_slack_bytes: int = 0
max_slack_bytes: int = 0
classmethod from_any(obj)[source]
Return type:

CacheOptions

merge(other)[source]

Overlay fields supplied by other onto self.

Return type:

CacheOptions

class Dataset(name, backend, path=None, _ids=None, _counts=None)[source]

Bases: object

Descriptor for a dataset with index and backend info.

Instances are created via from_path (file-backed) or from_dict (in-memory/testing). They always provide:

  • name: a user-facing identifier used in mixtures

  • backend: opaque metadata that lets FetchOp build a reader later

  • path: original filesystem path if file-backed, otherwise None

File-backed datasets also carry a few-KB handle to the node-local shard catalog (set by from_path); this is what travels in ctx, not shard_meta.

Shard counts are not stored as a mapping; read them via ids() / counts() / total() / max_count() / shard_count().

Backend kinds used by the internal store builder:

  • litdata/mds/jsonl/parquet/vortex: kind and path (no per-shard shards graph — that lives in the catalog now)

  • inmem: kind and shards (dict[int, InMemoryShard])

Note: this class does not expose any method to fetch rows; IO is delegated to an internal shard store owned by the FetchOp.

name: str
backend: Mapping[str, object]
path: str | None = None
ids()[source]

Return the sorted shard ids as an int64 array.

Return type:

ndarray

counts()[source]

Return per-shard sample counts (int64), aligned with ids().

Return type:

ndarray

raw_bytes()[source]

Return per-shard byte sizes (int64), aligned with ids().

File-backed datasets read the catalog’s on-disk shard sizes; in-memory shards size their resident payloads lazily (InMemoryShard.raw_bytes).

Return type:

ndarray

total()[source]

Total sample count across all shards.

Return type:

int

max_count()[source]

Largest single-shard sample count (0 when empty).

Return type:

int

shard_count()[source]

Number of shards.

Return type:

int

classmethod from_path(name, path, *, fmt=None)[source]

Construct a file-backed dataset descriptor.

Performs a count-only discovery: it obtains shard_id + num_rows as small numpy arrays (the work source’s input) without materializing the per-shard shard_meta graph. The full columnar catalog is built once per node later, by the Engine’s finalize() (or lazily by the store builder for Engine-less use).

Parameters:
  • name (str) – Logical dataset name used in mixtures and debugging.

  • path (str) – Filesystem directory containing the dataset.

  • fmt (str | None) – Optional explicit format. When None, auto-detects.

Return type:

Dataset

Supported formats:

  • litdata directories containing index.json structured with config and chunks

  • mds directories containing index.json structured with shards

  • jsonl directories where *.jsonl files act as shards

Special URI schemes:

  • hf://org/name[@rev]/[config/]split is served through the HuggingFace backend: the split’s uploaded files when their format is readable, else HuggingFace’s Parquet conversion when it is complete and built from the requested commit. fmt picks the source in that format; appending ~original or ~parquet to the revision (hf://org/name@~parquet/split) forces one. A partial conversion is used only with ZEPHON_HF_ALLOW_PARTIAL=1. Uploaded files are read as stored: the dataset card’s reader options and features casting are not applied. The dataset’s path becomes a URI pinned to the resolved commit, source and config.

Returns:

a descriptor populated with shard counts and a catalog handle for later IO.

Return type:

Dataset

Raises:
classmethod from_dict(name, shards)[source]

Construct an in-memory dataset descriptor.

Return type:

Dataset

class InMemoryShard(rows)[source]

Bases: object

List-backed shard useful for tests and quickstarts.

property raw_bytes: int

Estimated decoded payload bytes for this shard, sampled and cached.

close()[source]
getsamples(indices)[source]
Return type:

list[dict[str, object]]

class ParquetRGCacheOptions(enabled=_NotSupplied.TOKEN, root=_NotSupplied.TOKEN, limit_bytes=_NotSupplied.TOKEN, min_free_bytes=_NotSupplied.TOKEN)[source]

Bases: object

Node-shared decoded Parquet row-group cache configuration.

enabled=None enables the cache with the shard cache or an explicit root. An omitted root uses the shard cache’s reserved .parquet-rg-cache child. An omitted limit uses the deprecated shard option when present, otherwise the larger of 4 GiB and 10% of the shard cache limit. These defaults resolve when the store is built so option merges retain the distinction between omitted and explicit values. Explicit None resets a previously supplied field to its automatic default.

enabled: bool | None = 0
root: str | Path | None = 0
limit_bytes: int | None = 0
min_free_bytes: int = 0
classmethod from_any(obj)[source]

Normalize a decoded row-group cache configuration.

Return type:

ParquetRGCacheOptions

merge(other)[source]

Overlay fields supplied by other onto self.

Return type:

ParquetRGCacheOptions

class StoreOptions(cache=<factory>, parquet_rg_cache=<factory>)[source]

Bases: object

Top-level IO store options passed to FetchOp.

cache: CacheOptions
parquet_rg_cache: ParquetRGCacheOptions
resolved_parquet_rg_cache()[source]

Resolve decoded-cache defaults against the shard-cache options.

Return type:

ParquetRGCacheOptions

classmethod from_any(obj)[source]
Return type:

StoreOptions

merge(other)[source]

Return a new options object with other overriding self.

Return type:

StoreOptions