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
This commit is contained in:
@@ -1,5 +1,7 @@
|
||||
ADD_LIBRARY(JFJochReader STATIC
|
||||
JFJochReader.cpp JFJochReader.h
|
||||
ReadAhead.cpp
|
||||
ReadAhead.h
|
||||
JFJochHDF5Reader.cpp
|
||||
JFJochHDF5Reader.h
|
||||
SweepLayout.cpp
|
||||
|
||||
@@ -4,6 +4,8 @@
|
||||
#include "HDF5ImageLocator.h"
|
||||
#include "../common/JFJochException.h"
|
||||
|
||||
#include <algorithm>
|
||||
|
||||
namespace {
|
||||
// Coalesce consecutive single-image mappings into one contiguous range when the source and
|
||||
// virtual images stay contiguous in the same file/dataset.
|
||||
@@ -101,6 +103,22 @@ HDF5ImageLocator::Location HDF5ImageLocator::Resolve(int64_t global_image) const
|
||||
return {layout_.master_file, static_cast<uint32_t>(global_image), layout_.master_filename};
|
||||
}
|
||||
|
||||
std::vector<std::string> HDF5ImageLocator::DataFiles() const {
|
||||
std::vector<std::string> ret;
|
||||
if (layout_.format == FileWriterFormat::NXmxLegacy) {
|
||||
for (const auto &f: layout_.legacy_files)
|
||||
ret.push_back(f.path);
|
||||
} else if (layout_.format == FileWriterFormat::NXmxVDS
|
||||
&& layout_.data_layout == HDF5DataSetLayout::VIRTUAL) {
|
||||
for (const auto &mapping: layout_.vds_mappings)
|
||||
if (CoversFirstChannel(mapping)
|
||||
&& std::find(ret.begin(), ret.end(), mapping.filename) == ret.end())
|
||||
ret.push_back(mapping.filename);
|
||||
} else if (!layout_.master_filename.empty())
|
||||
ret.push_back(layout_.master_filename);
|
||||
return ret;
|
||||
}
|
||||
|
||||
std::vector<HDF5DataSourceMessage> HDF5ImageLocator::GetSourceMapping(uint64_t first_image,
|
||||
std::optional<uint64_t> image_count,
|
||||
uint64_t total_images,
|
||||
|
||||
@@ -64,6 +64,10 @@ public:
|
||||
// by the layout. Does not bounds-check against the total image count - the caller does that.
|
||||
Location Resolve(int64_t global_image) const;
|
||||
|
||||
// The files that hold the pixels, in image order: the data files of a legacy or VDS dataset, the
|
||||
// master itself when the images are in it.
|
||||
std::vector<std::string> DataFiles() const;
|
||||
|
||||
// Source mapping for re-writing a derived file (e.g. _process.h5) so it links back to the
|
||||
// original pixel sources rather than to a master. total_images is supplied by the caller.
|
||||
// stride is the step between consecutive images of the derived file in the SOURCE: image i of the
|
||||
|
||||
@@ -70,6 +70,9 @@ public:
|
||||
// file that holds a legacy/VDS image's per-image metadata.
|
||||
HDF5ImageLocator::Location Resolve(int64_t global) const;
|
||||
|
||||
// The files holding the pixels, in image order.
|
||||
std::vector<std::string> DataFiles() const { return locator_.DataFiles(); }
|
||||
|
||||
// Read the pixels at a resolved location into a CompressedImage backed by `buffer`. Templated on
|
||||
// the allocator so a caller can hand over a buffer that does not zero what it is about to
|
||||
// overwrite (RawByteBuffer).
|
||||
|
||||
@@ -8,10 +8,8 @@
|
||||
#include <cstring>
|
||||
#include <filesystem>
|
||||
#include <array>
|
||||
#include <future>
|
||||
#include <map>
|
||||
#include <optional>
|
||||
#include <thread>
|
||||
#include <tuple>
|
||||
|
||||
#include "../common/JFJochException.h"
|
||||
@@ -299,20 +297,11 @@ void JFJochCBFReader::ReadFiles(const std::string &path) {
|
||||
// sweep of several thousand frames this is seconds of startup before a single image is read, and
|
||||
// the files are independent.
|
||||
std::vector<sweep::Frame> frames(files_.size());
|
||||
{
|
||||
const size_t nthreads = std::min<size_t>(std::max(1u, std::thread::hardware_concurrency()), 8);
|
||||
std::vector<std::future<void>> futures;
|
||||
for (size_t t = 0; t < nthreads; t++)
|
||||
futures.push_back(std::async(std::launch::async, [&, t] {
|
||||
for (size_t i = t; i < files_.size(); i += nthreads) {
|
||||
const auto h = minicbf::ReadHeader(files_[i]);
|
||||
frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m,
|
||||
h.beam_x_px, h.beam_y_px, h.wavelength_A};
|
||||
}
|
||||
}));
|
||||
for (auto &f : futures)
|
||||
f.get();
|
||||
}
|
||||
sweep::ForEachInOrder(files_.size(), [&](size_t i) {
|
||||
const auto h = minicbf::ReadHeader(files_[i]);
|
||||
frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m,
|
||||
h.beam_x_px, h.beam_y_px, h.wavelength_A};
|
||||
});
|
||||
|
||||
// The sweep the headers describe, which is not always the files laid out end to end: a deposited
|
||||
// series can be missing frames, and those are gaps in the rotation rather than images to close up.
|
||||
@@ -406,6 +395,7 @@ bool JFJochCBFReader::ReadRawImage(int64_t image_number, JFJochReaderRawImage &i
|
||||
if (!HasImage(image_number))
|
||||
return false;
|
||||
image.image = DecodeInto(image_number, image.image_buffer, image.read_buffer);
|
||||
NoteImageRead();
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
@@ -37,6 +37,8 @@ class JFJochCBFReader : public JFJochReader {
|
||||
template <class Buffer>
|
||||
CompressedImage DecodeInto(int64_t image_number, Buffer &buffer, std::vector<uint8_t> &scratch) const;
|
||||
|
||||
std::vector<std::string> DataFiles() const override { return files_; }
|
||||
|
||||
public:
|
||||
~JFJochCBFReader() override = default;
|
||||
|
||||
|
||||
@@ -72,14 +72,21 @@ bool JFJochHDF5Reader::ReadRawImage(int64_t image_number, JFJochReaderRawImage &
|
||||
chunk = image_source_.PrepareDirectRead(loc);
|
||||
if (!chunk) {
|
||||
ret.image = image_source_.ReadImageAt(ret.image_buffer, loc);
|
||||
NoteImageRead();
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
ret.image = HDF5ImageSource::ReadDirect(ret.image_buffer, *chunk);
|
||||
NoteImageRead();
|
||||
return true;
|
||||
}
|
||||
|
||||
std::vector<std::string> JFJochHDF5Reader::DataFiles() const {
|
||||
std::unique_lock ul(hdf5_mutex);
|
||||
return image_source_.DataFiles();
|
||||
}
|
||||
|
||||
bool JFJochHDF5Reader::LoadImage_i(std::shared_ptr<JFJochReaderDataset> &dataset,
|
||||
DataMessage &message,
|
||||
std::vector<uint8_t> &buffer,
|
||||
|
||||
@@ -40,6 +40,8 @@ class JFJochHDF5Reader : public JFJochReader {
|
||||
|
||||
HDF5ImageLocator::Location GetImageLocation(int64_t image_number) const;
|
||||
|
||||
std::vector<std::string> DataFiles() const override;
|
||||
|
||||
public:
|
||||
~JFJochHDF5Reader() override = default;
|
||||
|
||||
|
||||
@@ -4,8 +4,6 @@
|
||||
#include "JFJochMarCCDReader.h"
|
||||
|
||||
#include <cmath>
|
||||
#include <future>
|
||||
#include <thread>
|
||||
|
||||
#include "../common/JFJochException.h"
|
||||
#include "../common/JFJochMath.h"
|
||||
@@ -116,20 +114,11 @@ void JFJochMarCCDReader::ReadFiles(const std::string &path) {
|
||||
// Reading one costs a 4 kB read, so on a sweep of several thousand frames this is worth
|
||||
// spreading over the cores, as the CBF path does for the same reason.
|
||||
std::vector<sweep::Frame> frames(files_.size());
|
||||
{
|
||||
const size_t nthreads = std::min<size_t>(std::max(1u, std::thread::hardware_concurrency()), 8);
|
||||
std::vector<std::future<void>> futures;
|
||||
for (size_t t = 0; t < nthreads; t++)
|
||||
futures.push_back(std::async(std::launch::async, [&, t] {
|
||||
for (size_t i = t; i < files_.size(); i += nthreads) {
|
||||
const auto h = marccd::ReadHeader(files_[i]);
|
||||
frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m,
|
||||
h.beam_x_px, h.beam_y_px, h.wavelength_A};
|
||||
}
|
||||
}));
|
||||
for (auto &f : futures)
|
||||
f.get();
|
||||
}
|
||||
sweep::ForEachInOrder(files_.size(), [&](size_t i) {
|
||||
const auto h = marccd::ReadHeader(files_[i]);
|
||||
frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m,
|
||||
h.beam_x_px, h.beam_y_px, h.wavelength_A};
|
||||
});
|
||||
|
||||
// The sweep the headers describe, which is not always the files laid out end to end: a deposited
|
||||
// series can be missing frames, and those are gaps in the rotation rather than images to close
|
||||
@@ -215,6 +204,7 @@ bool JFJochMarCCDReader::ReadRawImage(int64_t image_number, JFJochReaderRawImage
|
||||
if (!HasImage(image_number))
|
||||
return false;
|
||||
image.image = DecodeInto(image_number, image.image_buffer, image.read_buffer);
|
||||
NoteImageRead();
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
@@ -38,6 +38,8 @@ class JFJochMarCCDReader : public JFJochReader {
|
||||
template <class Buffer>
|
||||
CompressedImage DecodeInto(int64_t image_number, Buffer &buffer, std::vector<uint8_t> &scratch) const;
|
||||
|
||||
std::vector<std::string> DataFiles() const override { return files_; }
|
||||
|
||||
public:
|
||||
~JFJochMarCCDReader() override = default;
|
||||
|
||||
|
||||
@@ -64,6 +64,13 @@ std::shared_ptr<JFJochReaderRawImage> JFJochReader::GetRawImage(int64_t image_nu
|
||||
return ret;
|
||||
}
|
||||
|
||||
void JFJochReader::StartReadAhead() {
|
||||
read_ahead_.reset();
|
||||
const auto files = DataFiles();
|
||||
if (!files.empty() && GetNumberOfImages() > 0)
|
||||
read_ahead_ = std::make_unique<ReadAhead>(files, GetNumberOfImages());
|
||||
}
|
||||
|
||||
void JFJochReader::SetStartMessage(const std::shared_ptr<JFJochReaderDataset> &val) {
|
||||
std::unique_lock ul(m);
|
||||
dataset = val;
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
|
||||
#include <unordered_set>
|
||||
#include <map>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
|
||||
#include "../common/JFJochMessages.h"
|
||||
@@ -15,6 +16,7 @@
|
||||
#include "JFJochReaderDataset.h"
|
||||
#include "JFJochReaderImage.h"
|
||||
#include "JFJochReaderSpots.h"
|
||||
#include "ReadAhead.h"
|
||||
|
||||
class JFJochReader {
|
||||
mutable std::mutex m;
|
||||
@@ -31,8 +33,17 @@ class JFJochReader {
|
||||
int64_t n_image,
|
||||
int64_t image_jump,
|
||||
JFJochReaderImage &image);
|
||||
|
||||
std::unique_ptr<ReadAhead> read_ahead_;
|
||||
protected:
|
||||
void SetStartMessage(const std::shared_ptr<JFJochReaderDataset> &val);
|
||||
|
||||
// The files that hold the frame data, in the order the images sit in them: what StartReadAhead()
|
||||
// streams. Empty for a reader with no files of its own.
|
||||
virtual std::vector<std::string> DataFiles() const { return {}; }
|
||||
// Every ReadRawImage() that read an image reports it here, which is what moves the read-ahead on.
|
||||
void NoteImageRead() { if (read_ahead_) read_ahead_->ImageRead(); }
|
||||
|
||||
DiffractionExperiment default_experiment;
|
||||
public:
|
||||
virtual ~JFJochReader() = default;
|
||||
@@ -44,6 +55,11 @@ public:
|
||||
|
||||
virtual void Close() = 0;
|
||||
|
||||
// Start streaming the frame data into the page cache ahead of the image reads (see ReadAhead).
|
||||
// For a batch run over the whole dataset; call it once the dataset is open, before any image is
|
||||
// read. It changes only when the bytes arrive, never what is read.
|
||||
void StartReadAhead();
|
||||
|
||||
std::shared_ptr<JFJochReaderImage> LoadImage(int64_t image_number, int64_t summation_factor = 1);
|
||||
|
||||
void UpdateGeomMetadata(const DiffractionExperiment& experiment);
|
||||
|
||||
@@ -4,8 +4,6 @@
|
||||
#include "JFJochSMVReader.h"
|
||||
|
||||
#include <cmath>
|
||||
#include <future>
|
||||
#include <thread>
|
||||
|
||||
#include "../common/JFJochException.h"
|
||||
#include "../common/JFJochMath.h"
|
||||
@@ -116,20 +114,11 @@ void JFJochSMVReader::ReadFiles(const std::string &path) {
|
||||
// Reading one costs a 4 kB read, so on a sweep of several thousand frames this is worth
|
||||
// spreading over the cores, as the CBF path does for the same reason.
|
||||
std::vector<sweep::Frame> frames(files_.size());
|
||||
{
|
||||
const size_t nthreads = std::min<size_t>(std::max(1u, std::thread::hardware_concurrency()), 8);
|
||||
std::vector<std::future<void>> futures;
|
||||
for (size_t t = 0; t < nthreads; t++)
|
||||
futures.push_back(std::async(std::launch::async, [&, t] {
|
||||
for (size_t i = t; i < files_.size(); i += nthreads) {
|
||||
const auto h = smv::ReadHeader(files_[i]);
|
||||
frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m,
|
||||
h.beam_x_px, h.beam_y_px, h.wavelength_A};
|
||||
}
|
||||
}));
|
||||
for (auto &f : futures)
|
||||
f.get();
|
||||
}
|
||||
sweep::ForEachInOrder(files_.size(), [&](size_t i) {
|
||||
const auto h = smv::ReadHeader(files_[i]);
|
||||
frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m,
|
||||
h.beam_x_px, h.beam_y_px, h.wavelength_A};
|
||||
});
|
||||
|
||||
// The sweep the headers describe, which is not always the files laid out end to end: a deposited
|
||||
// series can be missing frames, and those are gaps in the rotation rather than images to close
|
||||
@@ -215,6 +204,7 @@ bool JFJochSMVReader::ReadRawImage(int64_t image_number, JFJochReaderRawImage &i
|
||||
if (!HasImage(image_number))
|
||||
return false;
|
||||
image.image = DecodeInto(image_number, image.image_buffer, image.read_buffer);
|
||||
NoteImageRead();
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
@@ -38,6 +38,8 @@ class JFJochSMVReader : public JFJochReader {
|
||||
template <class Buffer>
|
||||
CompressedImage DecodeInto(int64_t image_number, Buffer &buffer, std::vector<uint8_t> &scratch) const;
|
||||
|
||||
std::vector<std::string> DataFiles() const override { return files_; }
|
||||
|
||||
public:
|
||||
~JFJochSMVReader() override = default;
|
||||
|
||||
|
||||
+10
-2
@@ -437,14 +437,22 @@ void Slurp(const std::string &path, size_t max_bytes, std::vector<uint8_t> &buf)
|
||||
}
|
||||
|
||||
// Enough to reach the separator on any header seen in the wild (the longest measured is ~6.3 kB).
|
||||
// The first read is kept small because a sweep's startup reads every frame's header: on a spinning
|
||||
// disk 256 kB from each of 2400 files took twice as long as 16 kB (12 s against 6 s). A header
|
||||
// that does not end within it is read again with the large probe, so the result never depends on it.
|
||||
constexpr size_t HEADER_FIRST_PROBE_BYTES = 16 * 1024;
|
||||
constexpr size_t HEADER_PROBE_BYTES = 256 * 1024;
|
||||
|
||||
} // namespace
|
||||
|
||||
Header ReadHeader(const std::string &path) {
|
||||
std::vector<uint8_t> buf;
|
||||
Slurp(path, HEADER_PROBE_BYTES, buf);
|
||||
const auto start = FindBinarySection(buf.data(), buf.size());
|
||||
Slurp(path, HEADER_FIRST_PROBE_BYTES, buf);
|
||||
auto start = FindBinarySection(buf.data(), buf.size());
|
||||
if (!start.has_value()) {
|
||||
Slurp(path, HEADER_PROBE_BYTES, buf);
|
||||
start = FindBinarySection(buf.data(), buf.size());
|
||||
}
|
||||
if (!start.has_value())
|
||||
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
|
||||
"No CBF binary section in " + path);
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
// 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()));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
|
||||
// SPDX-License-Identifier: GPL-3.0-only
|
||||
|
||||
#pragma once
|
||||
|
||||
#include <atomic>
|
||||
#include <cstdint>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
#include <vector>
|
||||
|
||||
// Streams the files that hold a dataset's frames into the page cache, front to back, ahead of the
|
||||
// image loops that read them.
|
||||
//
|
||||
// On a spinning disk a cold run otherwise waits on the disk: the header scan, the pre-scan and the
|
||||
// first-pass indexing touch a few hundred frames and leave the disk idle for tens of seconds, and
|
||||
// then the first image loop reads everything at whatever rate its scattered requests get. Streaming
|
||||
// from the moment the dataset is open keeps the disk busy through all of that, at its sequential
|
||||
// rate, so the loops find their frames already in memory.
|
||||
//
|
||||
// The files are read in 4 MB pieces by a few threads that take the pieces strictly in order: on the
|
||||
// disk measured here one thread streams at 125 MB/s, eight such threads at 190 MB/s (and 32 at
|
||||
// 157 MB/s), so eight keep the stream ahead of an image loop that is itself reading from disk.
|
||||
//
|
||||
// It changes only WHEN bytes arrive, never what is read: the loops still read every image themselves,
|
||||
// and a frame the read-ahead has not reached yet is read cold, as it always was.
|
||||
//
|
||||
// It never runs more than a window - a quarter of the memory available when it starts - ahead of
|
||||
// what the loops have read. A dataset larger than that is not streamed in only to push out its own
|
||||
// beginning before the first loop gets there; and a page the read-ahead brought in is then read a
|
||||
// second time by the loop, which is what keeps it in the page cache for the loops after.
|
||||
class ReadAhead {
|
||||
struct Piece {
|
||||
size_t file;
|
||||
uint64_t offset; // in the file
|
||||
uint64_t position; // in the whole stream
|
||||
};
|
||||
|
||||
std::vector<std::string> files_;
|
||||
std::vector<Piece> pieces_;
|
||||
double window_ = 0;
|
||||
double bytes_per_image_ = 0;
|
||||
std::atomic<size_t> next_piece_{0};
|
||||
std::atomic<uint64_t> images_read_{0};
|
||||
std::atomic<bool> stop_{false};
|
||||
std::vector<std::thread> threads_;
|
||||
|
||||
void Run();
|
||||
public:
|
||||
// files: in the order the images sit in them; n_images: how many images they hold between them.
|
||||
ReadAhead(std::vector<std::string> files, uint64_t n_images);
|
||||
~ReadAhead();
|
||||
ReadAhead(const ReadAhead &) = delete;
|
||||
ReadAhead &operator=(const ReadAhead &) = delete;
|
||||
|
||||
// An image loop has read one image; the window moves on by one image's worth of bytes.
|
||||
void ImageRead() { images_read_++; }
|
||||
};
|
||||
@@ -4,10 +4,13 @@
|
||||
#include "SweepLayout.h"
|
||||
|
||||
#include <algorithm>
|
||||
#include <atomic>
|
||||
#include <cmath>
|
||||
#include <cctype>
|
||||
#include <filesystem>
|
||||
#include <future>
|
||||
#include <optional>
|
||||
#include <thread>
|
||||
|
||||
#include "../common/JFJochException.h"
|
||||
#include "../common/Logger.h"
|
||||
@@ -281,4 +284,17 @@ Layout Place(const std::vector<Frame> &frames, const std::string &logger_name) {
|
||||
return out;
|
||||
}
|
||||
|
||||
void ForEachInOrder(size_t n, const std::function<void(size_t)> &fn) {
|
||||
const size_t nthreads = std::min<size_t>(std::max(1u, std::thread::hardware_concurrency()), 8);
|
||||
std::atomic<size_t> next{0};
|
||||
std::vector<std::future<void>> futures;
|
||||
for (size_t t = 0; t < nthreads; t++)
|
||||
futures.push_back(std::async(std::launch::async, [&] {
|
||||
for (size_t i = next++; i < n; i = next++)
|
||||
fn(i);
|
||||
}));
|
||||
for (auto &f : futures)
|
||||
f.get();
|
||||
}
|
||||
|
||||
} // namespace sweep
|
||||
|
||||
@@ -3,6 +3,8 @@
|
||||
|
||||
#pragma once
|
||||
|
||||
#include <cstddef>
|
||||
#include <functional>
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
||||
@@ -57,4 +59,10 @@ struct Layout {
|
||||
// logger_name is the reader's own logger, so the gap report names the format the user passed.
|
||||
Layout Place(const std::vector<Frame> &frames, const std::string &logger_name);
|
||||
|
||||
// Call fn(i) for every i in [0, n) on up to 8 threads, handing the indices out in order. For reading
|
||||
// every frame's header: the threads then read neighbouring files at any moment, so a spinning disk
|
||||
// sees one sweep through the directory rather than eight positions drifting apart (2400 CBF headers
|
||||
// from a cold spinning disk: 18 s with each thread striding through its own share, 12 s in order).
|
||||
void ForEachInOrder(size_t n, const std::function<void(size_t)> &fn);
|
||||
|
||||
} // namespace sweep
|
||||
|
||||
@@ -1456,6 +1456,10 @@ static int RunRugnux(int argc, char **argv) {
|
||||
exit(EXIT_FAILURE);
|
||||
}
|
||||
JFJochReader &reader = *reader_ptr;
|
||||
// Stream the frames in from disk from now on, so a cold or slow disk works through the setup and
|
||||
// the pre-scan instead of sitting idle until the first image loop. --mode scale reads no images.
|
||||
if (mode != RugnuxMode::Scale)
|
||||
reader.StartReadAhead();
|
||||
|
||||
const auto dataset = reader.GetDataset();
|
||||
if (!dataset) {
|
||||
|
||||
@@ -4128,6 +4128,21 @@ TEST_CASE("JFJochCBFReader_incomplete_header_is_refused_not_misread", "[HDF5][Fu
|
||||
}
|
||||
}
|
||||
|
||||
TEST_CASE("MiniCBF_header_longer_than_the_first_probe", "[HDF5][Full]") {
|
||||
// The header is read with a small first probe, and again with the large one when it does not end
|
||||
// inside it; what stands past the first 16 kB must come out exactly as before.
|
||||
std::string filler;
|
||||
for (int i = 0; i < 1000; i++)
|
||||
filler += "# Comment line padding the header out past the first probe\n";
|
||||
WriteRawMiniCBF("cbflong_0001.cbf", "# Pixel_size 172e-6 m x 172e-6 m\n" + filler
|
||||
+ "# Wavelength 0.96864 A\n# Detector_distance 0.33161 m\n", 24, 16);
|
||||
const auto h = minicbf::ReadHeader("cbflong_0001.cbf");
|
||||
CHECK(h.pixel_x_m == Catch::Approx(172e-6));
|
||||
CHECK(h.wavelength_A == Catch::Approx(0.96864));
|
||||
CHECK(h.distance_m == Catch::Approx(0.33161));
|
||||
remove("cbflong_0001.cbf");
|
||||
}
|
||||
|
||||
// marCCD: a TIFF whose instrument header sits in the gap between the TIFF header and the pixels.
|
||||
// The fixtures below are written byte for byte rather than through libtiff, because the layout IS
|
||||
// what the reader has to cope with - the private tag that states where the instrument header
|
||||
|
||||
Reference in New Issue
Block a user