Streaming: fetch only the parts of a dataset an analysis actually touches (split from #151)
Nessuno ha ancora preso questa issue.
Valutazione
- Difficoltà
- 5/5
- Tempo stimato
- Più di una settimana
- Idoneità per principianti
- 35/100
- Tipo di issue
- Funzionalità
- Chiarezza
- Abbastanza chiara
- Stato di attività
- Attiva
- Stack tecnologico
- huggingface, python
- Ambito
- data-engineering, distributed-systems
Direzione di ricerca
Inizia leggendo #151 e confrontando le quattro alternative elencate qui con lo schema di esecuzione speculativa. Affronta esplicitamente la sicurezza rispetto alle interruzioni e il modo in cui l’approccio scelto segue le regole di #151 relative a dichiarazione, verifica del digest e assenza di fallback sull’intero dataset. Il lavoro sarà considerato completato quando saranno disponibili un confronto scritto e una validazione rispetto a uno store remoto reale usando un dataset più grande della memoria del worker.
Scritto dal modello di indicizzazione a partire dal testo della issue.
Descrizione
Summary
Split out of #151 at the repository owner's request. #151 ("support data sharing") builds the
data-package object: declared data is staged to a private HuggingFace repo (or carried inline when
small) and dereferenced on the worker on demand. That gets the whole declared dataset to the
worker. This issue is about the harder half: getting only the parts of a dataset that a given
analysis actually touches, so an analysis can run on a machine that cannot hold the whole thing.
The owner's framing, quoted in full
From @jeremymanning on #151:
re: streaming: i'm not 100% sure how to leverage this properly. if there is a way we can
systematically figure out which part(s) of the dataset are needed at a given time, then this would
be incredibly powerful because we could (for example) run an analysis on a machine that couldn't
actually fit the full dataset-- along the lines of how HF streaming enables inference on models
when the machine executing the inference can't fit the full dataset. but generalizing this may be
more complicated. perhaps it's possible by doing something like running a mock simulation that
wraps every function recursively to either execute the function (if it takes less than some very
small threshold amount of time) OR simply track which parts of the dataset are accessed inside
that piece. i'm imagining this would use threading and timers in some way, and then interrupt
function calls if they took too long. and then we'd have to do some sort of fancy bookkeeping to
tag different parts of the execution and figure out which parts of the dataset go with which tags,
and then organize the to-be-streamed dataset so that parts could be copied over on demand. unless
i'm missing an obvious solution, this is likely more complex than should be attempted in this
initial "support data sharing" implementation. if so, we should open a new issue to address the
streaming functionality separately, and defer this part of this issue accordingly.
And, on why the deferral:
this is likely more complex than should be attempted in this initial "support data sharing"
implementation.
What this issue must decide before any code
The sketch above has two separable mechanisms, and they have very different risk profiles.
- Access tracking — discover which parts of a dataset a function touches.
- Speculative partial execution — run the function under a time budget, interrupting calls that
run long, and attribute the accesses observed so far to execution "tags".
(2) is the expensive and dangerous part. Interrupting arbitrary user code part-way through, by
threads and timers, is not generally safe: a partially executed function may have already written a
file, posted to an API, or mutated shared state, and there is no way to know from the outside. Any
plan here has to say what class of function it is willing to speculate on, and how a user opts in.
(1) is tractable on its own and may be most of the value. Concretely cheaper alternatives worth
pricing before building the simulator:
- Lazy chunked handles. The data package hands the function an array-like/table-like object whose
__getitem__fetches the covering chunk on demand and caches it. No tracing, no interruption; the
access pattern is the fetch pattern. Costs a wrapper type per supported format. - Ride existing streaming.
datasets.load_dataset(..., streaming=True)andzarr/h5pyover a
remote store already solve this for their own formats. Clustrix may only need to hand the worker
credentials plus a URI, not invent a mechanism. - Declared partitioning. The caller says how the dataset splits and which partition a call needs.
Explicit, unglamorous, and consistent with #151's "declaration, never inference" rule. - Recorded first run. Run once with full data and record accesses; use the recording to prefetch
on subsequent runs. Sidesteps interruption entirely, at the cost of needing one full-size run.
Acceptance criteria (to be refined once an approach is chosen)
- A written comparison of the four alternatives above against the simulator sketch, with the
interruption-safety question answered explicitly. - Whatever is built must obey #151's rules: declaration over inference, digest verification of every
fetched part, and no silent fallback to fetching the whole dataset when streaming fails. - Verified against a real remote store with a dataset larger than the worker's memory. Per this
repository's standard, a mocked demonstration proves nothing here.
Explicitly out of scope
Same list as #151: no sync/mirror/watch, no DAG or data-derived ordering, no cluster-to-cluster
transfer, no provenance database.
Related
- #151 (the data-package implementation this was split out of)
- Lingua principale
- Python
- Stelle
- 10
- Fork
- 4
- Merge medio
- 15h 10m
- PR unite (30g)
- 2
Preparare l'ambiente
Come iniziare
- Leggi tutta la issue e poi la guida ai contributi del progetto.
- Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
- Fai un fork del repository e lavora su un branch.
- Apri una pull request che faccia riferimento al numero della issue.
Altre issue di ContextLab/clustrix
-
Difficoltà 4/5 3-5 giorni Idoneità per principianti 52/100
ContextLab/clustrix#175 ·
-
Difficoltà 5/5 Più di una settimana Idoneità per principianti 25/100
ContextLab/clustrix#170 ·
-
enhancement epic
Difficoltà 5/5 Più di una settimana Idoneità per principianti 20/100
ContextLab/clustrix#160 ·
-
enhancement
Difficoltà 5/5 Più di una settimana Idoneità per principianti 42/100
ContextLab/clustrix#151 · 5 commenti ·
-
Difficoltà 5/5 Più di una settimana Idoneità per principianti 25/100
ContextLab/clustrix#146 · 2 commenti ·
Tutte le issue di ContextLab/clustrix
Issue simili
-
Broken links found in docsApertadocs pydanty:is-working
Difficoltà 2/5 1-3 ore Idoneità per principianti 75/100
pydantic/pydantic-ai#8863 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 68/100
run-llama/llama_index#23278 ·
I maintainer di solito rispondono entro 2 giorni
-
documentation from-review-extraction github-actions priority: low severity:nit
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 92/100
LearningCircuit/local-deep-research#6946 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 82/100
oracle/langchain-oracle#323 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 88/100
tenstorrent/tt-metal#58057 · 1 commento ·
I maintainer di solito rispondono entro 1 giorno