Skip to content

perf(read): serve concurrent reads of one file in parallel - #246

Open
XciD wants to merge 20 commits into
mainfrom
perf/mmap-concurrent-reads
Open

XciD wants to merge 20 commits into
mainfrom
perf/mmap-concurrent-reads

Conversation

@XciD

@XciD XciD commented Sep 29, 2026 •

Copy link
Copy Markdown
Member

Fixes #234. Supersedes #241.

Summary

Parallel page faults on one memory-mapped file ran one at a time, at about 3 MB/s. transformers, torch and safetensors load checkpoints this way. The per-handle PrefetchState mutex was held across network I/O, and interleaved streams replaced each other in its single buffer, so almost every fault became a new 256 KiB range download of about 50 ms. This is the stall that made transformers turn mmap off on hf-mount (huggingface/transformers#45547).

Many cold readers at once also wedged the mount (#234). Each read stream held xet download buffer for data that only a FUSE worker could consume, and the workers were blocked on streams that could not get buffer for their first bytes. Reads timed out at offset 0 and ended in EIO.

This PR replaces PrefetchState with a RemoteReader per lazy handle (src/virtual_fs/remote_reader.rs):

  • Blocks, fetched concurrently. The file is cached in 256 KiB blocks. Missing blocks are fetched by background tasks, each streaming one bounded range. No lock is held across network I/O, and a read of a block in flight waits for it instead of fetching it again.
  • Read-ahead per stream. Up to 32 sequential streams are tracked per handle. Each keeps a window of blocks requested ahead of its reads: 2 MiB at first (64 MiB for a read from the start of the file), growing fourfold then twofold up to 512 MiB, and topped up in fetches that download in parallel. A fetch close to the reader is short (1 MiB, then as long as its distance to the reader) so that the first blocks arrive early; fetches far ahead are 32 MiB, and no fetch crosses a multiple of 32 MiB, so that every pass over a file asks for the same ranges and a re-read finds them in the chunk cache. A stream's read-ahead stops where another active stream ahead of it started, also inside its last block, so parallel copies of one tensor do not fetch each other's ranges. A stream that starts up to 16 MiB past the furthest block read (the next tensor, or the next record header of a torch zip) reads ahead at once; further reads (the slices of a tensor-parallel rank) wait for a second read, as with the old prefetch code.
  • Fetching goes on past a slow fetch, within what the link carries. A xet range delivers its data only when the whole range has arrived, so fetches finish out of order. At most half of the read-ahead budget is in flight, and each completed fetch tops its stream's window up again while the stream reads. While a read waits for a slow fetch, the next ones keep downloading instead of the network idling. xet downloads the terms of the fetches in the order they were asked for, so on a slow link a deep pipeline would make the last fetches wait past the fetch timeout: a reader caps the read-ahead in flight that no read waits for yet (64 MiB at first, plus the bytes of each read-ahead fetch that arrives within a quarter of the fetch timeout, halved by each slower or failed one, between 8 MiB and half its limit). On a fast link the cap soon reaches half the limit.
  • Bounded memory. A handle holds at most 512 MiB of pending or unread blocks. Read-ahead is reserved in a budget per mount before its fetches start (1 GiB, --read-ahead-mb), and each handle may hold an even share of it. A handle that holds read-ahead or that the budget turns away counts in the split; one whose only fetches are blocks its reads wait for does not. A handle above its share gives fetched read-ahead back at its next read (first what no stream heads to, then the blocks furthest ahead). Blocks already read (for re-reads and for streams that meet in a block; the page cache keeps the rest) are capped at 8 MiB per handle and 256 MiB across the handles of a mount, keeping the last block of each live stream. A handle that has no read waiting and read no new block for 10 s gives back its read-ahead beyond its share; after 120 s it gives back all of it and the blocks it read, and its next read starts small windows again. So a reader that reads in bursts keeps what its next burst needs, and a file kept open does not keep a share. Every 256 reads, a handle also drops the unread blocks of streams that stopped. Blocks are copied out of xet terms, so that a cached block does not keep a whole 64 MiB term alive.
  • Fetches do not wait for readers. A fetch copies its data into blocks as it arrives, whether a reader waits for it or not, so no xet download buffer stays held until a FUSE worker consumes it. This breaks the cycle of prefetch hangs / deadlocks in some way during mass fetch #234.
  • FUSE reads answer from a task when they wait. A read of data already fetched is still served on the worker thread. A read that waits for the network finishes in a task, so it does not hold one of the --max-threads workers while other faults queue behind it.
  • xet settings from fix(xet): give each read stream its own download buffer (#234) #241. The download buffer limit goes from 256 MiB to 1 GiB, as one fast reader alone keeps up to 256 MiB in flight, and the cap of the adaptive download concurrency goes from 64 to 124, unless HF_XET_FIXED_DOWNLOAD_CONCURRENCY is set. The per-stream download buffers of fix(xet): give each read stream its own download buffer (#234) #241 are not needed: fetches are bounded and drained as they arrive.

Retries (3 attempts in a row that deliver nothing, resuming at the first missing block), the per-chunk stall timeout (--read-fetch-timeout-ms), EOF handling, and --direct-io (a re-read fetches the bytes again, even inside a block read in part) keep their behavior. A read that waits for several fetches fails as soon as one of them gives up; a read whose failed fetch it only joined tries once more with fetches of its own.

New unit tests cover parallel reads, sharing a block in flight, interleaved streams fetching each byte once, read-ahead up to the end of the file, forward scans, refill while a read waits, the budgets of a handle and of a mount (with a handle the budget turns away), the read-ahead of an idle handle (and of one whose read waits for a slow fetch), blocks a read waits for staying at hand until it ends, the block of a cancelled read, a read that ends after its failed block was fetched again, a fetch task that ends early, reads of 8 and 40 blocks whose fetches fail and stall at once, fetches spawned on a runtime that shuts down, the eviction queue and partial re-reads with --direct-io, cancellation on close, streams that meet inside a block, sparse forward reads, bursts with pauses and a read after a long pause, refill once a stream stops, readers with only demand fetches, the read-ahead of a stopped stream, blocks already read split across a mount and kept at the end of a live stream, attempts that deliver blocks, a read that joined a failed fetch, fetch runs within cells, blocks reads wait for and the cap, and a slow link (a mock link that downloads terms in the order they were asked for, as xet does). test_fuse_parallel_mmap_read has 16 threads copy one memory-mapped 48 MiB file through a mount with 4 worker threads and checks the content. test_fuse_parallel_cold_reads, from #241, has 24 cold readers of 64 MiB files go through a mount with 4 worker threads and a 5 s fetch timeout, and fails on any read error or fetch timeout.

Benchmarks

m6i.2xlarge in us-east-1, --no-disk-cache, page cache dropped and a fresh mount for every run. Qwen/Qwen2.5-1.5B-Instruct/model.safetensors (3.09 GB) unless noted. Throughput to CAS changed a lot during the runs (a plain sequential read of the same file went from about 260 to 830 MB/s from one run to the next), so the sequential rows interleave the binaries and give medians. Only the last table includes the xet settings commit.

Workload main This PR
safetensors, 4 threads, natural key order (transformers-like) 2.7 MB/s (77 of 338 tensors in 411 s) 195 to 224 MB/s
safetensors, 1 thread, file order (torch copies each tensor with 8 threads) 5.9 MB/s (161 of 338 tensors in 302 s) 301 to 364 MB/s
torch.load(mmap=True) then copy, 4 threads, EleutherAI/pythia-1.4b (2.93 GB) 2.0 MB/s (47 of 364 tensors in 358 s) 68 to 82 MB/s (36 to 43 s)
8 threads copying random 4 MiB ranges of one mapping 1.8 MB/s, p50 17.9 s per copy 157 to 210 MB/s, p50 152 to 195 ms
Sequential read(), 54 MB file (median of 3) 110 MB/s 166 MB/s
Sequential read(), 91 MB file (median of 3) 149 MB/s 275 MB/s
Sequential read(), 268 MB file (median of 3) 285 MB/s 290 MB/s
Sequential read(), 3.09 GB file (median of 4) 590 MB/s 946 MB/s
Sequential mmap scan, 3.09 GB file (median of 4) 575 MB/s 893 MB/s
Random 4 KiB page faults, p50 / p99 (3 rounds) 62 to 63 ms / 183 to 227 ms 63 to 67 ms / 195 to 281 ms

The runs on main stopped at their deadline, hence the partial tensor counts. The small files are google/electra-small-discriminator/pytorch_model.bin, sentence-transformers/all-MiniLM-L6-v2/model.safetensors and distilbert/distilbert-base-uncased/model.safetensors.

Peak RSS of the daemon during a sequential read of the 3.09 GB file: 875 to 959 MiB with the first commits, about 720 to 830 MiB after the review fixes (below), against 698 to 747 MiB on main (read-ahead budget of 512 MiB per handle).

Bytes read through a mount were compared with a copy downloaded directly from the Hub, for 16 threads reading contiguous ranges of one mapping, 16 threads reading random ranges, and a safetensors load tensor by tensor: all identical.

Many cold readers (#234)

Shards of Qwen/Qwen2.5-72B-Instruct read at once from a cold mount, each from its start, with --max-threads 4 --read-fetch-timeout-ms 5000 like test_fuse_parallel_cold_reads:

Load main This PR
24 files, first 64 MiB of each 177 s, EIO on 19 files 2.6 s (619 MB/s), no errors
37 files, first 256 MiB of each 301 s, EIO on 32 files 15 s (661 MB/s), no errors

In the second run with this PR, one fetch got no data for 5 s and its retry succeeded.

xet settings

3 rounds with the default thread count and fetch timeout, peak RSS of the daemon in parentheses:

Workload Without the settings With the settings
37 cold readers, first 256 MiB of each 978 to 1106 MB/s (1156 to 1185 MiB) 1097 to 1216 MB/s (1172 to 1258 MiB)
Sequential read() of the 3.09 GB file 577 to 845 MB/s (771 to 907 MiB) 669 to 976 MB/s (832 to 895 MiB)

For one reader of a small file, the settings make no measurable difference. Whole-file read() as in the CI benchmark, medians of 5 rounds:

File main This PR without the settings This PR
google/electra-small-discriminator/pytorch_model.bin (54 MB) 112 MB/s 192 MB/s 182 MB/s
distilbert/distilbert-base-uncased/model.safetensors (268 MB) 308 MB/s 350 MB/s 353 MB/s
openai-community/gpt2/model.safetensors (548 MB) 368 MB/s 548 MB/s 515 MB/s

The CI benchmark varies more than that between runs of the same code: three runs of its 500 MB FUSE read gave 893, 524 and 1000 MB/s.

Review fixes

The mount budget reserved before fetches start, idle handles, and the bookkeeping cleanups (4552bf4 to 32f3b26), against the code before them with the same xet settings. 3 rounds, peak RSS of the daemon in parentheses:

Workload Before After
37 cold readers, first 256 MiB of each 829 to 1173 MB/s (1190 to 1246 MiB) 1148 to 1236 MB/s (1095 to 1118 MiB)
24 cold readers of 64 MiB, --max-threads 4 --read-fetch-timeout-ms 5000 761 to 953 MB/s 675 to 845 MB/s
Sequential read() of the 3.09 GB file 794 to 889 MB/s (756 to 784 MiB) 702 to 864 MB/s (741 to 831 MiB)
safetensors, 4 threads, natural key order (1 round) 210 MB/s (822 MiB) 234 MB/s (905 MiB)

The cleanup that followed (eb85607) reads at the same speed as 32f3b26, within the noise: over 3 rounds, 628 to 779 against 707 to 789 MB/s with 37 readers, 574 to 698 against 648 to 754 MB/s for one reader, and 189 to 225 against 189 to 219 MB/s for safetensors with 4 threads.

No fetch timed out. The same binary varies by up to 30% between rounds here (the code before gave 599 to 965 MB/s in the 24-reader case over two runs), so only the 37-reader case moved beyond the noise. When many readers start at once, the first ones give back the end of their window as the others arrive and fetch it again later.

Second review round

PR before the fixes of the second review (506487e) against this PR, m6i.2xlarge, 5 rounds interleaved, medians:

Workload Before After
safetensors, 4 threads, natural key order 254 MB/s 230 MB/s
Sequential read(), 548 MB / 3.09 GB 546 / 966 MB/s 572 / 959 MB/s
fio sequential, 548 MB, FUSE / NFS 555 / 527 MB/s 591 / 518 MB/s
37 cold readers, first 256 MiB of each 1283 MB/s 1250 MB/s
Headers of 37 shards 3.3 s 3.2 s

The safetensors load is about 9% slower in every comparison of the day; variants that restore the 64 MiB forward-scan gap, drop the cap of read-ahead in flight or drop the cell alignment did not bring it back, and CDN stalls during the runs (30 s without data on some fetches, for both binaries) blurred the smaller differences.

With all traffic into the instance shaped to 3 MB/s (tc), reading 400 MiB of the 3.09 GB file: main read about 3 MiB in 15 minutes (its 128 MiB prefetch windows got no data within 30 s, EIO every 3 minutes); the PR before these fixes had 143 fetch timeouts and an EIO at 165 MiB; this PR reads it in 136 s (3.08 MB/s, the link rate), with no timeout, no error and 275 MiB of peak RSS.

Notes

  • A sequential read makes about three CAS reconstruction queries per 32 MiB fetch (304 for the 3.09 GB file): xet splits each fetch in 8 and 16 MiB reconstruction blocks (HF_XET_RECONSTRUCTION_MIN_RECONSTRUCTION_FETCH_SIZE, unchanged) that it downloads in parallel and delivers in turn. The single stream of main made a few queries for the whole file. With one 32 MiB block per fetch, torch.load(mmap=True) opened faster (12 to 14 s instead of 23 to 30 s), but small files read 40% slower. Fetches now fall on the same 32 MiB cells on every pass, so the plan cache serves a re-read within the hour; a first read can only make fewer queries once CachedXetClient can derive tight range plans from a full plan, which needs chunk sizes in the terms of xet-core (derive_range_response).
  • fix(fuse): flush deadlock in advanced-writes #228 moves every FUSE operation to a task. This PR only moves the reads that wait for the network: a read of data already fetched is answered on the worker thread, without a task handoff for each 128 KiB request. Both change the same read handler.
  • transformers still reads safetensors fully into memory on hf-mount (_is_on_hf_mount). Its first forward pass over mmap-backed weights touches them one parameter at a time, in module order, not file order, so each parameter costs at least one round trip: about 55 s for this model, against about 8 s for the full read. Dropping that workaround needs more than this PR (for example, filling the page cache ahead of the reads).

Page faults of a memory-mapped file ran one at a time. The per-handle
PrefetchState mutex was held across network I/O, and interleaved streams
replaced each other in its single buffer, so parallel loaders
(transformers, torch, safetensors) turned almost every fault into a new
256 KiB range download. Loading a checkpoint this way ran at about 3 MB/s.

Replace PrefetchState with a RemoteReader per lazy handle: a 256 KiB
block cache filled by concurrent bounded fetches, with blocks in flight
shared between reads, read-ahead per sequential stream (topped up as
fetches complete), and memory budgets per handle and per mount. FUSE
reads that wait for the network finish in a task instead of holding a
worker thread.
@github-actions

Copy link
Copy Markdown
Contributor

POSIX Compliance (pjdfstest)

============================================================
  pjdfstest POSIX Compliance Results
------------------------------------------------------------
  Files: 130/130 passed    Tests: 832 total (0 subtests failed)
  Result: PASS
------------------------------------------------------------
  Category               Passed    Total   Status
  -------------------- -------- -------- --------
  chflags                     5        5       OK
  chmod                       8        8       OK
  chown                       6        6       OK
  ftruncate                  13       13       OK
  granular                    5        5       OK
  mkdir                       9        9       OK
  open                       19       19       OK
  posix_fallocate             1        1       OK
  rename                     10       10       OK
  rmdir                      11       11       OK
  symlink                    10       10       OK
  truncate                   13       13       OK
  unlink                     11       11       OK
  utimensat                   9        9       OK
============================================================

@github-actions

github-actions Bot commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Benchmark Results

============================================================
  Benchmark — 50MB
------------------------------------------------------------
  Metric                                 FUSE          NFS
  ------------------------------ ------------ ------------
  Sequential read                    274.1 MB/s     287.3 MB/s
  Sequential re-read                2228.4 MB/s    2205.1 MB/s
  Range read (1MB@25MB)                0.2 ms         0.2 ms
  Random reads (100x4KB avg)           0.0 ms         0.0 ms
  Sequential write (FUSE)           1244.4 MB/s
  Close latency (CAS+Hub)            0.111 s
  Write end-to-end                   329.9 MB/s
  Dedup write                       1567.1 MB/s
  Dedup close latency                0.101 s
  Dedup end-to-end                   377.6 MB/s
============================================================
============================================================
  Benchmark — 200MB
------------------------------------------------------------
  Metric                                 FUSE          NFS
  ------------------------------ ------------ ------------
  Sequential read                    630.3 MB/s     836.0 MB/s
  Sequential re-read                2232.6 MB/s    2227.9 MB/s
  Range read (1MB@25MB)                0.2 ms         0.2 ms
  Random reads (100x4KB avg)           0.0 ms         0.0 ms
  Sequential write (FUSE)           1364.5 MB/s
  Close latency (CAS+Hub)            0.091 s
  Write end-to-end                   840.9 MB/s
  Dedup write                       1341.8 MB/s
  Dedup close latency                0.108 s
  Dedup end-to-end                   777.9 MB/s
============================================================
============================================================
  Benchmark — 500MB
------------------------------------------------------------
  Metric                                 FUSE          NFS
  ------------------------------ ------------ ------------
  Sequential read                    977.8 MB/s    1106.4 MB/s
  Sequential re-read                2129.3 MB/s    2322.9 MB/s
  Range read (1MB@25MB)                0.3 ms         0.2 ms
  Random reads (100x4KB avg)           0.0 ms         0.0 ms
  Sequential write (FUSE)           1274.9 MB/s
  Close latency (CAS+Hub)            0.099 s
  Write end-to-end                  1018.5 MB/s
  Dedup write                       1362.7 MB/s
  Dedup close latency                0.130 s
  Dedup end-to-end                  1005.7 MB/s
============================================================
============================================================
  fio Benchmark Results
------------------------------------------------------------
  Job                        FUSE MB/s   NFS MB/s  FUSE IOPS   NFS IOPS
  ------------------------- ---------- ---------- ---------- ----------
  seq-read-100M                  207.5      350.9                      
  seq-reread-100M               1470.6      186.6                      
  rand-read-4k-100M                0.1        0.1         16         18
  seq-read-5x10M                1000.0      609.8                      
  rand-read-10x1M                160.7       24.8      41130       6341
  Random Read Latency           FUSE avg      NFS avg
  ------------------------- ------------ ------------
  rand-read-4k-100M           62551.1 us   55998.9 us
  rand-read-10x1M                23.5 us     156.1 us
============================================================

Files of a few dozen MiB read sequentially got slower (CI benchmark:
94 MB/s for 50 MB, against 159 to 263 MB/s on main). Read-ahead started
at 2 MiB and doubled, and a fetch delivers nothing until its whole range
has arrived, so each step of the ramp cost a full round trip.

A read from the start of a file now asks for 64 MiB at once, windows grow
fourfold while below 64 MiB, and fetches close to the reader are short
(1 MiB, then as long as their distance to the reader) so that the first
blocks arrive early. Fetches far ahead stay at 32 MiB.

Keep HF_XET_RECONSTRUCTION_MIN_RECONSTRUCTION_FETCH_SIZE at 8 MiB: xet
then splits each fetch in reconstruction blocks of 8 and 16 MiB that it
downloads in parallel and delivers in turn. Small files read 50 to 85%
faster than with one 32 MiB block per fetch.
Fold in the parts of #241 that still apply with RemoteReader: a download
buffer limit of 1 GiB (one fast reader alone keeps up to 256 MiB in
flight, the whole former limit), a download concurrency cap of 124, and
its regression test for #234 (many cold readers, 4 FUSE workers, 5 s
fetch timeout), which now looks for the RemoteReader timeout message.

#241's per-stream download buffers are not needed: fetches are bounded
and drained as they arrive, so no stream holds buffer while its reader
waits for a FUSE worker.
@XciD
XciD marked this pull request as ready for review September 30, 2026 07:23
XciD added 16 commits September 30, 2026 09:27
- A reader kept the blocks it fetched ahead until its handle closed. The
  mount budget is split between readers when they ask for read-ahead, so
  files read a little and then kept open held their earlier, larger
  shares: 64 files with one byte read from each held about 1.5 GiB. A
  reader that has not read for 10 s now drops its unread blocks, and
  fetches that complete after that no longer top its window up.
- With --direct-io, a block dropped once read left its entry in the
  eviction queue. Behind a block read only in part, which stays at the
  front, the queue grew by one entry for each block read.
- stream_calls_record_read_ahead_and_seek_ranges relied on the order in
  which fetch tasks on a multi-threaded runtime record their calls. It
  now finds each call by its range.
- A block a read waits for is flagged when the read plans, and arrives
  as read: the accounting of a first read lives in State::mark_read, and
  a refill between the fetch and the read can no longer evict the block
  and fetch it again.
- The eviction queue holds only blocks that reads touched (remove_read
  keeps it exact), so evict_behind no longer skips stale entries.
- trim_read_ahead counts the blocks it wants only when the range may not
  fit, and make_room finds the live streams once per call.
- A fetch drops the sender of each block it delivered, so a block
  dropped after its read no longer stays in memory until its run ends.
- Smaller cleanups: now_or_never in the FUSE read handler, block_len,
  next_chunk, Stream::is_live and min_run, FetchRun.read_ahead derived
  from its stream, test helpers (reader_with_limits, upload_files), Bytes
  in the mock, and comments that still described the prefetch buffer.
…by progress

- Read-ahead is reserved in the mount budget before its fetches start,
  so the readers of a mount hold at most 1 GiB together (16 MiB each
  beyond 64 readers). The split between readers only capped what each
  one could ask for: 16 readers of one file that start at once held
  292 MiB against a 256 MiB budget in the new unit test. A reader whose
  share shrank also gives back, at its next read, the read-ahead that
  none of its streams heads to.
- A reader is idle once no read waits for a fetch and no read reached a
  new block for 10 s. A read that waits for a slow fetch no longer makes
  its reader look idle (which dropped the read-ahead fetched past it),
  and re-reads of blocks at hand no longer keep a reader active forever.
- A refill that starts read-ahead also starts the idle watch, which may
  have ended just before.
- The mount budget is a fixed 1 GiB. The floor of 16 MiB per reader let
  128 readers hold 2 GiB, and a reader that left could bring the capacity
  below what the others held.
- A reader the budget cuts counts in the split until it goes idle, so
  that the shares of the others shrink, and a reader above its share
  gives fetched read-ahead back at its next read: first the blocks no
  stream heads to, then those furthest ahead of their stream. Before, a
  reader opened while two others held the whole budget got no read-ahead
  until they had read theirs. Cut by the budget, a reader also fetches
  runs of any useful size (from 1 MiB) instead of waiting for a whole run.
- A block a read waits for joins the eviction queue once that read ends,
  so that concurrent reads still find it meanwhile instead of fetching
  it again.
- report_ahead skips the atomic update when nothing changed.
- evict_unread replaces make_room, shrink_to and drop_read_ahead: one
  order (blocks no live stream heads to, oldest first, then the blocks
  furthest ahead of their stream) and one stop condition.
- Locking the state returns a guard that reports the read-ahead to the
  mount budget when it unlocks, so that no path can change ahead_bytes
  and forget to. The budget owns its accounting (report, release).
- A block counts the reads waiting for it instead of a flag: it joins
  the eviction queue once the last of them ends, and the block of a read
  cancelled before it arrived is plain read-ahead again instead of
  staying in memory until the handle closes.
- Waiting is built where the count it releases is taken, the idle task
  uses State::idle, a top-up lists its missing blocks once, complete
  looks its block up once, and track reads the tick itself.
- Test helpers: reader_with, ahead.
…e blocks of a fetch that ends early

- A read that waited for a block could end after that block was removed
  and fetched again (its fetch gave up, or a forward-only read dropped
  it): leaving then took the count of the new fetch below zero, a panic
  in debug builds and, in release, a block marked read with a huge count
  that nothing evicted until the handle closed. The count saturates: a
  stray leave at most makes a read block count as read-ahead.
- Waiting is built last in plan: dropped by a panic there, it would have
  locked the state the call still held.
- FetchRun fails the blocks it did not deliver when it is dropped, so a
  fetch task that ends early (aborted, or it panicked) no longer leaves
  blocks in flight that every later read waits on in vain.
…first failed block

- Fetch tasks were spawned under the state lock. A runtime that shuts
  down drops a spawned future at once, on the spawning thread, and since
  b92cbb1 dropping a FetchRun locks the state to fail its blocks: a read
  or a refill during shutdown deadlocked its worker, and the daemon
  never exited. Fetches are now prepared under the lock and spawned
  after it is released.
- A read that waits for several fetches waited for them one by one, so
  a block whose fetch gave up ended the read only after the slower
  fetches before it (forever, with the fetch timeout disabled). The read
  now waits for all of them at once and fails at the first failure.
… fetch, refetch partial re-reads with --direct-io

- A read waits for its blocks with FuturesUnordered. try_join_all keeps
  results in order past 30 futures, so for a read of more than 30 blocks
  (possible when fs.fuse.max_pages_limit is raised) a block whose fetch
  gave up stayed queued behind a stalled earlier fetch.
- Blocks carry the id of the fetch that filled them, and a read that
  ends leaves only the fetch it waited for. After a fetch gave up, a read
  that ended late could count itself out of a later fetch of the same
  block: that block then joined the eviction queue while another read
  still waited for it, or counted as unread read-ahead.
- With --direct-io, a block read in part records how far reads consumed
  it, and a read that starts before that point fetches the block again,
  as the forward buffer of the old prefetch code did. Only whole blocks
  were dropped once read, so re-reads inside a block read in part were
  served from memory.
- The NFS handle pool comment said the remote reader caps fetched blocks
  across handles. It caps read-ahead; blocks already read are capped per
  handle only (8 MiB), so the pool caps those.
…e fetches drop on panic

- HF_XET_CLIENT_AC_MAX_DOWNLOAD_CONCURRENCY now has a default (124), and
  xet-runtime reads HF_XET_FIXED_DOWNLOAD_CONCURRENCY only when the
  canonical variables are unset: a user who pinned a fixed download
  concurrency (for example under a proxy connection cap) got an adaptive
  one that grows to 124. The download defaults now stay unset when the
  fixed alias is set, as the upload defaults already did.
- In plan() and refill(), the prepared fetch runs were declared after the
  state lock, so a panic between prepare_fetches and the unlock dropped
  them with the lock held, and FetchRun::drop locks the state to fail its
  blocks: the thread deadlocked instead of unwinding. They are now
  declared before the lock.
…rs, count only readers that hold read-ahead

- Parallel copies of one tensor split it in ranges that most often meet
  inside a block. next_stream_start ignored a stream that started in the
  last block of another, so the earlier stream's read-ahead ran into the
  next range and fetched it again (1.73x the file for 8 threads after a
  100 KiB header). It now stops there too.
- A new stream up to 64 MiB past the furthest block read continued the
  forward scan, so sparse forward reads (the slices of a tensor-parallel
  rank) fetched the gaps between them, with windows doubling from one
  read to the next (322 blocks fetched for 8 read 17 MiB apart). The gap
  is now 16 MiB, the forward-skip distance of the old prefetch code.
- After 10 s without a read, a reader dropped all its read-ahead but kept
  its large windows, so a reader that reads in bursts (training steps)
  fetched up to a whole window again at every burst (2x the file for
  bursts of 8 MiB with 30 s pauses). An idle reader now gives back only
  what exceeds its share of the mount budget at 10 s, and everything at
  120 s, after which its next read starts the windows small again.
- A finished read-ahead fetch topped the window up even once the stream
  had stopped reading. It now does only while a read waits or the stream
  read since the fetch started.
- Blocks that reads wait for count in ahead_bytes while in flight, so a
  reader with only such fetches counted among the readers that split the
  mount budget, and other readers dropped fetched read-ahead down to the
  smaller share. Only read-ahead counts now (pending_read_bytes).
- A reader that keeps reading elsewhere never goes idle, so the
  read-ahead of a stream that stopped (64 MiB after a header read at
  offset 0) and a stale starved flag stayed until close. Every 256 reads,
  the reader now drops the unread blocks of stopped streams and clears
  the flag.
…ks, add --read-ahead-mb

- Each handle kept up to 8 MiB of blocks already read until it closed,
  outside any mount budget: hundreds of memory-mapped dataset shards
  kept 8 MiB each. The readers open on a mount now split 256 MiB of them
  (8 MiB at most each), and an idle reader gives them back with its
  read-ahead after 120 s.
- All streams of a handle share one queue of blocks already read, so
  during a parallel mmap load a stream that paused found its boundary
  block evicted by the others and fetched it again, block after block.
  The eviction now keeps the last block of each live stream.
- The read-ahead budget of a mount (1 GiB) is now an option,
  --read-ahead-mb, so that a daemon with a tight memory limit (a CSI
  sidecar) can lower it. HF_XET_RECONSTRUCTION_DOWNLOAD_BUFFER_LIMIT
  still bounds the download buffers of xet-core.
…pts after progress, retry a read whose joined fetch failed

- On a slow link, read-ahead kept up to half the reader limit in flight
  (256 MiB), in fetches whose terms xet downloads in the order they were
  asked for, so the first term of the last fetch arrived after all the
  others: past the 30 s fetch timeout, with retries queued behind the
  same load, and reads failed with EIO (6 errors over 400 MiB read at
  3 MB/s in a model of that link). Read-ahead now waits while one of its
  fetches is older than a quarter of the fetch timeout, and a reader caps
  its read-ahead in flight: each fetch that arrives in time adds its bytes
  to the cap (from 64 MiB, up to half the reader limit), each slow or
  failed one halves it (down to 8 MiB). The same model then reads at the
  link rate, with no error and no retry.
- A fetch gave up after 3 attempts in total, even when every attempt
  delivered blocks. Only attempts in a row that deliver nothing count
  now, as the old prefetch code allowed.
- A read that joined a fetch in flight got its EIO with no attempt of its
  own, even after the network came back. A read whose failed fetch it
  only joined now tries once more with fetches of its own.
… the chunk cache

Fetch runs started wherever the reader and the budget left them, so a
second pass over a file (a remount, a pod restart, another process)
asked for other ranges than the first. xet splits each run in
reconstruction queries from its start, and the chunk cache serves a
range only from one item that holds all of it: most of a warm re-read
came from CAS again (22-33% hits in a model). A run now never crosses a
multiple of MAX_FETCH_BLOCKS blocks, and a cut by the budget backs up to
the last such boundary when that still leaves a run worth a fetch: past
the first cell, a scan fetches whole cells, the same on every pass. Their
reconstruction queries are the same too, so the plan cache also serves
them within the hour.
trim_read_ahead reserved all the read-ahead the budgets allowed, then
released what the range did not need. In between, a concurrent reader
could find the mount budget full and get its read-ahead cut. It now
reserves at most the bytes of the missing blocks.
Each read that steps a stream collected the missing blocks of its
read-ahead range (up to a whole window of map lookups) and went through
the budget, even when the bytes in flight were at their bound and no
fetch could start. It now returns at once in that case, as the trim
would have fetched nothing.
…p the fetch-age gate

- The cap of read-ahead in flight counted every pending block, also the
  blocks reads wait for. Many concurrent reads (mmap faults, NFS read
  RPCs) used it up and left no room for read-ahead until the cap grew.
  It now counts the bytes in flight that no read waits for yet.
- No read-ahead started while one read-ahead fetch of the reader was
  older than a quarter of the fetch timeout. On a slow link the adaptive
  cap alone keeps fetches within the timeout (400 MiB read at the link
  rate, no error, in the model of a 3 MB/s link), and a single stalled
  connection (30 s without data, retried) stopped all the read-ahead of
  its reader for up to 90 s, which showed in benchmarks during CDN
  stalls. The gate is gone.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

prefetch hangs / deadlocks in some way during mass fetch

1 participant