Data packages¶
Clustrix ships your function and its arguments. It does not ship your dataset.
A function that opens "data/subjects.h5" finds that file on your laptop and
does not find it on the worker, and no amount of decorating changes that. A
data package is how a dataset travels: you name the files, clustrix moves them,
and the function reads them back through the package on whichever machine it
happens to be running on.
Nothing is inferred. A file moves because you named it, never because a string
in your source code looked like a path. That restraint is deliberate: an upload
triggered by the literal "s3://bucket/notes.log" is the worst failure mode
available here, so declaration is the only route.
The shape of it¶
Three steps, and the middle one is the ordinary one:
# cluster-required: needs a configured cluster and a real data/ directory
import clustrix
# 1. Declare
subjects = clustrix.data_package("data/subjects.h5")
# 2. Pass it like any other argument
@clustrix.cluster(cores=8)
def fit(pkg):
# 3. Dereference inside the function
with open(pkg.path("subjects.h5"), "rb") as handle:
return len(handle.read())
fit(subjects)
clustrix.data_package() accepts a path, a list of paths, a directory, or
raw bytes. A directory expands to the files beneath it, and a list may mix
files and directories freely. Whatever you name, the paths inside the package
are relative to the common ancestor of everything named, so a function written
against "data/subjects.h5" keeps working against "data/subjects.h5" on
the worker. Pass base= when you want a different root.
Dereference with path(), which returns a filesystem
path, or with read_bytes(), which returns the
contents. Both check the file against the digest recorded when the package was
built, so a file that arrived wrong is an error rather than a wrong answer.
Both take no argument at all when the package holds exactly one file. A package
dereferenced on the machine that built it reads your original files in place –
no copy, no fetch – provided their digests still match what was packaged.
Here is the same round trip against the local backend, which needs no
cluster and no network, and which this page executes for real every time it is
checked:
import os
import clustrix
from clustrix import configure
configure(cluster_type="local")
os.makedirs("trials", exist_ok=True)
with open("trials/run1.csv", "w") as handle:
handle.write("trial,rt\n1,0.42\n")
trials = clustrix.data_package("trials", force_local=True)
print(trials.filenames(), trials.total_bytes) # ['run1.csv'] 16
@clustrix.cluster(cores=2)
def count_rows(pkg):
with open(pkg.path("run1.csv")) as handle:
return len(handle.readlines())
print(count_rows(trials)) # 2
force_local=True says “carry the contents inside the object regardless of
size”. Nothing is uploaded, and there is nothing to clean up afterwards.
Several packages travel as easily as one, and a list of them is walked the same way a single one is:
# cluster-required: needs a configured cluster to execute
import os
import clustrix
packages = [
clustrix.data_package("data/subjects.h5"),
clustrix.data_package("data/stimuli/"),
]
@clustrix.cluster(cores=4)
def summarize(pkgs):
roots = clustrix.materialize_packages(pkgs)
return [len(os.listdir(root)) for root in roots]
summarize(packages)
clustrix.materialize_packages() walks lists, tuples and dicts, writes
every package it finds to local disk, and hands back the directory holding
each. Anything that is not a package comes back unchanged.
Where the bytes actually live¶
Two places, and which one you get depends on size.
Small packages ride inside the object. Below stage_inline_max_bytes,
the file contents are carried in the package itself, pickled alongside your
function’s other arguments, and shipped over the transport that already moves
the payload – SFTP for ssh and slurm, the payload channel for
huggingface. No second transport, no remote store, nothing to delete.
Larger packages go to a private HuggingFace dataset repo. Above that threshold, the contents are uploaded and the package carries only the coordinates plus a digest per file. Three things follow, and you want to know all of them before it happens rather than after:
It needs HuggingFace credentials. The token comes from
hf_tokenin your config, thenHF_TOKENin the environment, then the cache thathf auth loginwrites. Without one, staging refuses and says so. The worker needs one too, from its own environment: clustrix deliberately does not pickle your token into the job payload, so a package staged remotely is unreadable on a worker with no HuggingFace credentials of its own.Clustrix creates a repo in your account. The first package that does not fit inline calls
create_repo(private=True, exist_ok=True)for<namespace>/clustrix-data. The namespace ishf_namespaceif you set one, otherwisehf_username, otherwise whatever the token’swhoami()reports. Sethf_data_repoto name a different repo outright.This applies on every backend. A
slurmjob with a package too big to inline still stages that package through HuggingFace, because that is the only remote store clustrix has.
Each package gets its own folder in that repo, keyed by a fresh identifier, so two packages never share a stored blob and deleting one cannot pull data out from under another. The cost of that is duplication: identical content packaged twice is stored twice.
Digests are computed locally, from your own files, and travel to the worker inside the function payload – which is an upload-only, local-origin artifact. Bytes fetched back out of the store are checked against those digests. In this way a tampered store is caught, because the expected digest never went through it.
The three size thresholds¶
Field |
Default |
Effect |
|---|---|---|
|
1 MB |
Under this, the package is inline. At or above it, the package is uploaded. |
|
100 MB |
At or above this, staging logs a warning before it starts. A transfer that takes minutes with no output is indistinguishable from a hang. |
|
5 GB |
At or above this, staging raises |
The inline threshold is measured on the serialized package, not on the raw data. Those two numbers are nowhere near each other once a package holds many small files, because every file also carries a relative path and a 64-character digest: ten thousand four-byte files are forty kilobytes of data and a 1.09 MB pickle. Measured on the data alone, that package would be a megabyte over the limit and still call itself inline, which is why the limit is applied to the pickle instead.
Deleting a package¶
Nothing is ever cleaned up for you. There is no TTL, no reaper, no eviction
policy, and no deletion when a job finishes. cleanup_on_success governs the
job directory on the cluster and does not touch staged data. Whether a dataset
is still needed is your judgement rather than clustrix’s, so a package you
staged stays staged, and stays billable, until you say otherwise.
The handle that created a package can end it:
# cluster-required: needs HuggingFace credentials
import clustrix
pkg = clustrix.data_package("data/subjects.h5")
pkg.delete()
delete() returns whether anything was actually
removed. Calling it twice is not an error, and neither is calling it on a
package somebody already cleaned up elsewhere. A remote copy that is present
and refuses to delete does raise, because a warning there would leave you
paying for storage you believe you released.
Four things it does not touch, each for its own reason:
The files you packaged.
local_rootpoints at your own data.A directory you named yourself. If you called
materialize(dest=...), that directory is yours; it may hold anything, and clustrix cannot tell what it put there from what was already there.The cache directory above the package. Exactly one directory is removed,
<local_cache_dir>/data-packages/<package id>, which clustrix created and which is keyed by an identifier nothing else uses.The repo itself, only the package’s folder inside it. An account whose last package is deleted keeps an empty
clustrix-datadataset, which is yours to remove by hand. A user who pointedhf_data_repoat a repo they own and care about would not thank clustrix for deleting it because the last package went away.
There is deliberately no context-manager form. A with block that quietly
deleted the upload on the way out would be exactly the automatic cleanup this
design rejects.
When the object is gone¶
Losing the handle does not mean losing the ability to clean up, because the object is not the only key to its own deletion:
# cluster-required: needs HuggingFace credentials
import clustrix
for record in clustrix.list_data_packages():
print(record["package_id"], record["name"], record["total_bytes"])
clustrix.delete_data_package("0123456789abcdef0123456789abcdef")
clustrix.list_data_packages() returns one record per package in the
store; an upload that was interrupted before its manifest went up appears with
"complete": False. clustrix.delete_data_package() takes an id and
validates it before anything reaches the Hub, since the id becomes a path in
the store and an id that is not one is a deletion aimed somewhere else.
Keeping a package across sessions¶
The package object is the durable handle, and it is plain data: strings, ints, bytes. No client, no socket, no credential. So the way to keep a staged dataset across sessions is to pickle the object and load it later, and a copy loaded in a fresh interpreter still reaches the same remote data:
# cluster-required: needs HuggingFace credentials
import pickle
import clustrix
pkg = clustrix.data_package("data/subjects.h5")
with open("subjects.pkl", "wb") as handle:
pickle.dump(pkg, handle)
# ... a week later, a different interpreter ...
with open("subjects.pkl", "rb") as handle:
pkg = pickle.load(handle)
pkg.path("subjects.h5") # still resolves
pkg.delete() # still deletes
Credentials are deliberately not among the attributes that survive that round trip. A saved package must not be a token sitting on disk, so it re-authenticates from the ordinary config path every time it is used.
What it refuses¶
Every refusal below raises clustrix.StagingError before anything
moves.
Paths that look like credentials. *.pem, *.key, .env,
.netrc, id_rsa*, kubeconfig, anything under .ssh/, .aws/,
.gnupg/, .kube/ or .config/gcloud/, a .git/config (which is one
of the commonest places a personal access token ends up on disk), and so on.
Both ends of a symlink are tested, since data/notes pointing at
~/.ssh/id_rsa is credential-shaped at the far end and innocuous at the
near one. Pass allow_sensitive=True if you really do mean to move a keypair
to the worker.
Anything that is not a regular file. Reading a fifo or /dev/zero does
not fail, it blocks, so packaging one would hang with no output and no
timeout. A refusal costs you one message.
A data repo that already exists and is public. create_repo(private=True,
exist_ok=True) creates a private repo but does not make an existing public
one private; it returns the repo as it is. Clustrix refuses rather than
flipping the setting, because a repo can be public on purpose and silently
changing someone’s visibility is its own incident. Make it private yourself, or
point hf_data_repo somewhere else.
A package at or above stage_max_bytes – with the largest file named.
Two different files that would land on the same name. One would silently overwrite the other on the worker and the run would produce a wrong answer rather than an error. Naming the same file twice is harmless and collapses.
Files that change while they are being staged. Files go up by path, so the bytes on the wire are whatever the file held at upload time rather than what was hashed a moment earlier. Clustrix re-hashes after the commit and, if anything moved, removes the folder it just created and tells you which files changed.
When not to use this¶
If the data is already reachable from the worker, staging it is pure waste. On
an HPC cluster with a shared filesystem your /scratch directory is visible
from every compute node, so the file your function opens is already there and
moving a copy of it through HuggingFace buys nothing but a transfer and a
storage bill. The same goes for data in object storage the worker can
authenticate to.
Use the read-only filesystem utilities to check.
cluster_exists answers the question directly:
from clustrix import cluster_exists
from clustrix.config import ClusterConfig
config = ClusterConfig(cluster_type="local", local_work_dir=".")
if cluster_exists("setup.py", config):
print("already there; nothing to stage")
Data packages are for the case where that answer is no.
API¶
- clustrix.data_package(source, *, name=None, base=None, config=None, force_local=False, allow_sensitive=False, filename='data.bin')[source]¶
Package data or files into an object a
@clusterfunction can read.- Parameters:
source (
Union[str,Path,bytes,Sequence[Union[str,Path]]]) – A path, a list of paths, or rawbytes. Directories expand to the files beneath them.name (
Optional[str]) – Label for messages; defaults to the source’s basename.base (
Union[str,Path,None]) – What the package’s relative paths are relative to. Defaults to the common ancestor of everything named, which is what lets a function keep using the paths it already uses.config – A
ClusterConfig; the global one by default.force_local (
bool) – Carry the contents inside the object regardless of size. Nothing is uploaded and there is nothing to clean up.allow_sensitive (
bool) – Permit paths that look like credentials.filename (
str) – The name rawbytesget inside the package.
- Return type:
- Returns:
A
DataPackage. Pass it – or a list of them – to a@cluster-decorated function as an ordinary argument.- Raises:
StagingError – on a missing file, a credential-shaped path, a path that escapes the package root, or a package at or above
stage_max_bytes.
- class clustrix.DataPackage(name, package_id, files, local_root=None, repo_id=None, path_in_repo=None, inline=None, _materialised=None)[source]¶
A named bundle of files that can travel to a worker and be read there.
Build one with
data_package(); the constructor here is the low-level form and does no staging of its own.- name¶
A label, used in messages and in the materialisation directory.
- package_id¶
Unique per package. Also the folder name in the remote store, which is what keeps two packages from sharing a blob.
- files¶
One
PackagedFileper file, in declaration order.
- local_root¶
Where the files are on the machine that built the package.
Nonefor a package built from in-memory data. Used as a zero-cost source when the same machine dereferences the package.
- repo_id / path_in_repo
The private remote location, or
Nonefor an inline package.
- inline¶
relpath -> byteswhen the payload rides inside the object.
Every attribute is plain data – strings, ints, bytes. No HF client, no socket, no file handle, so the object pickles cleanly and a copy loaded in a fresh interpreter still reaches the same remote data. Credentials are deliberately not among those attributes: a saved package must not be a token sitting on disk, so it re-authenticates from the ordinary config path whenever it is used.
- path(relpath=None, config=None)[source]¶
Local filesystem path to one file, fetching it if it is not here yet.
Called with no argument on a single-file package. This is the “on demand” half: nothing is fetched until something asks for it, and a package dereferenced on the machine that built it reads the original files without copying them.
- Return type:
- read_bytes(relpath=None, config=None)[source]¶
The contents of one file, verified against its recorded digest.
- Return type:
- materialize(dest=None, config=None)[source]¶
Put every file on local disk and return the directory holding them.
Files keep the relative paths they had when the package was built, so a function that opened
"data/x.csv"opens"data/x.csv"under the returned root on either machine.- Return type:
- exists(config=None)[source]¶
Whether the remote copy is still there. Inline packages: always.
- Return type:
- __getstate__()[source]¶
Pickle without the materialisation path.
Where this package was last unpacked is true of one machine at one moment. Carrying it into a pickle means a worker – or this machine a week later – can find a stale directory at that path and reuse it. Everything else here is durable; this one field is not.
- delete(config=None)[source]¶
Remove the remote copy and any clustrix-owned local cache of it.
Returns whether anything was actually removed remotely. Calling this twice, or on a package that was already cleaned up elsewhere, is not an error – but a remote copy that is present and cannot be deleted raises, because a warning here would leave the caller paying for storage they believe they released.
Nothing calls this for you. There is no TTL and no reaper: whether a dataset is still needed is the user’s judgement, not clustrix’s.
The files you packaged are never touched.
local_rootpoints at the user’s own data; deleting that because a transfer was cleaned up would be indefensible. Only the remote folder, and the copy clustrix wrote into its own cache, are removed.A directory you named yourself is never touched either. If you called
materialize(dest=...), that directory is yours – it may hold anything, and clustrix has no way to know what it created there versus what was already in it. It used to be removed recursively, which deleted whatever else the caller kept alongside the data. Clear it yourself if you want it gone.The repo itself is never touched either, only the package’s folder inside it.
hf_data_repomay well point at a repo the user owns and cares about. An account whose last package is deleted keeps an emptyclustrix-datadataset, which is theirs to remove.- Return type:
- __init__(name, package_id, files, local_root=None, repo_id=None, path_in_repo=None, inline=None, _materialised=None)¶
- clustrix.materialize_packages(value, config=None)[source]¶
Walk a structure and materialise every
DataPackagein it.Convenience for a worker that would rather have paths than handles:
roots = clustrix.materialize_packages(packages)
Lists, tuples, and dicts are walked; anything else is returned unchanged.
- Return type:
- clustrix.list_data_packages(config=None)[source]¶
Every package clustrix has staged in the remote store.
A user who lost the handle still needs a way to see what is costing them storage, so the object is not the only key to its own deletion. Returns a list of manifests; incomplete uploads appear with
"complete": False.
- clustrix.delete_data_package(package_id, config=None)[source]¶
Delete a staged package by id, without needing the object.
The counterpart to
list_data_packages(). Returns whether anything was removed; a package that is already gone is not an error, a package that is there and will not delete raises.The id is validated first, before anything reaches the Hub. It becomes a path in the store, so an id that is not one is a deletion aimed somewhere else:
""addressed the wholepackages/prefix and removed every package in the account, and"../README.md"climbed out of the prefix and removed a file that was never a package.- Return type: