Files
Jungfraujoch/reader/ReadAhead.cpp
T
leonarski_fandClaude Opus 5.5 aa3fa6b9c7 Reader: stream the frame data into the page cache ahead of the image loops
A cold run on a spinning disk waited on the disk twice over. The CBF header
scan read 256 kB from every frame on eight threads that each strode through
their own share of the sweep, so they drifted apart and the scan became a
seek storm (34 s for 2400 frames here); and after it, the pre-scan and the
first-pass indexing touch a few hundred frames and leave the disk idle until
the first image loop reads everything at seek-bound rates.

- ReadAhead (reader/): once the dataset is open, rugnux starts eight threads
  that read the data files - HDF5 data files (legacy, VDS or the integrated
  master) or the per-frame CBF/marCCD/SMV files - in 4 MB pieces taken
  strictly in order, into a throwaway buffer. One stream reads this disk at
  125 MB/s, eight in-order streams at 190 MB/s, 32 at 157 MB/s. It never gets
  more than a quarter of MemAvailable (GlobalMemoryStatusEx on Windows, 4 GiB
  where there is no figure) ahead of what ReadRawImage has handed out, so a
  dataset bigger than the cache does not evict its own start, and it stops
  with the reader. Plain ifstream reads: portable, no POSIX calls.
- Header scans (CBF, marCCD, SMV) hand the files out in order from an atomic
  counter (sweep::ForEachInOrder) instead of striding: 18 s -> 12 s for 2400
  cold CBF headers. The CBF header is first read with a 16 kB probe and again
  with the old 256 kB one only when the separator is not in it, so the parsed
  header is exactly what it was: 12 s -> 6 s.

Output unchanged: p.hkl, p.mtz and p_unmerged.mtz md5-identical to the
rc173 baseline on 6toc (CBF, 2400 frames, 6.0 GB) and 9q41 (HDF5 VDS, 900
frames, 5.1 GB), and on 6z9g (HDF5, 12.8 GB) to the unmodified branch; myob
(p.hkl p.mtz p_P1.mtz p_unmerged.mtz) md5-identical to the reference.

Measured cold (files evicted with POSIX_FADV_DONTNEED before every run),
same code without this commit vs with it, on a shared box (load 20-70, other
agents reading the same disk, so single runs scatter by +-20 s):
  6toc  wall 61.7/62.7 -> 49.3/49.8 s (clean pairs); all data resident
        after 62/51/50 -> 45/41/42 s
  9q41  wall 67.3 -> 57.8 s (clean pair); resident after 58/46/43 -> 48/37/35 s
  6z9g  resident after 81 -> 69 s
The first image loop can look slower with this in CBF runs: the old 256 kB
header probes pulled ~70% of the data in as kernel readahead, so the old
loop started warm - after a 34 s header scan instead of 10 s.
Warm (myob, NVMe, cached): 19.35/20.00 s without, 19.76-20.16 s with; the
read-ahead then only copies 9.3 GB out of the page cache, 0.44 s wall and
3.4 CPU-s measured standalone.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01D1G8gJVAy6gp1K5Dz3NE5C
2026-09-27 10:32:58 +02:00

82 lines
2.4 KiB
C++

// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include "ReadAhead.h"
#include <chrono>
#include <filesystem>
#include <fstream>
#include <limits>
#ifdef _WIN32
#include <windows.h>
#endif
namespace {
constexpr uint64_t PIECE_BYTES = 4 << 20;
constexpr size_t THREADS = 8;
// Memory the system can hand out without swapping. Where there is no such figure (macOS), 4 GiB.
uint64_t AvailableMemory() {
#ifdef _WIN32
MEMORYSTATUSEX status{};
status.dwLength = sizeof(status);
if (GlobalMemoryStatusEx(&status))
return status.ullAvailPhys;
#else
std::ifstream meminfo("/proc/meminfo");
std::string key;
uint64_t kb = 0;
while (meminfo >> key >> kb) {
if (key == "MemAvailable:")
return kb * 1024;
meminfo.ignore(std::numeric_limits<std::streamsize>::max(), '\n');
}
#endif
return 4ULL << 30;
}
} // namespace
ReadAhead::ReadAhead(std::vector<std::string> files, uint64_t n_images) : files_(std::move(files)) {
uint64_t total = 0;
for (size_t i = 0; i < files_.size(); i++) {
std::error_code ec;
const uint64_t size = std::filesystem::file_size(files_[i], ec); // a gap in the sweep is ""
if (ec)
continue;
for (uint64_t offset = 0; offset < size; offset += PIECE_BYTES)
pieces_.push_back({i, offset, total + offset});
total += size;
}
bytes_per_image_ = static_cast<double>(total) / static_cast<double>(n_images);
window_ = static_cast<double>(AvailableMemory() / 4);
for (size_t t = 0; t < THREADS; t++)
threads_.emplace_back(&ReadAhead::Run, this);
}
ReadAhead::~ReadAhead() {
stop_ = true;
for (auto &t: threads_)
t.join();
}
void ReadAhead::Run() {
std::vector<char> buffer(PIECE_BYTES);
for (size_t i = next_piece_++; i < pieces_.size(); i = next_piece_++) {
const auto &piece = pieces_[i];
while (static_cast<double>(piece.position) > window_ + static_cast<double>(images_read_) * bytes_per_image_) {
if (stop_)
return;
std::this_thread::sleep_for(std::chrono::milliseconds(20));
}
if (stop_)
return;
std::ifstream in(files_[piece.file], std::ios::binary);
in.seekg(static_cast<std::streamoff>(piece.offset));
in.read(buffer.data(), static_cast<std::streamsize>(buffer.size()));
}
}