Files
Jungfraujoch/reader/JFJochHDF5Reader.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

244 lines
9.2 KiB
C++

// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include "JFJochHDF5Reader.h"
#include "../common/JFJochException.h"
void JFJochHDF5Reader::ReadFile(const std::string &filename) {
std::unique_lock ul(hdf5_mutex);
image_source_.Clear();
snapshots_.clear();
active_metadata_.reset();
active_snapshot_.clear();
number_of_images = 0;
try {
auto metadata = std::make_shared<HDF5MetadataSource>();
auto open_result = metadata->Open(filename, default_experiment);
image_source_.Configure(std::move(open_result.image_layout));
// Original file: per-image metadata is co-located with the pixels.
metadata->UseImageSourceForMetadata(&image_source_);
number_of_images = open_result.number_of_images;
snapshots_["Original"] = metadata;
active_metadata_ = metadata;
active_snapshot_ = "Original";
SetStartMessage(metadata->Dataset());
} catch (const std::exception &e) {
image_source_.Clear();
snapshots_.clear();
active_metadata_.reset();
active_snapshot_.clear();
number_of_images = 0;
SetStartMessage({});
throw;
}
}
uint64_t JFJochHDF5Reader::GetNumberOfImages() const {
std::unique_lock ul(hdf5_mutex);
return number_of_images;
}
void JFJochHDF5Reader::Close() {
std::unique_lock ul(hdf5_mutex);
image_source_.Clear();
snapshots_.clear();
active_metadata_.reset();
active_snapshot_.clear();
number_of_images = 0;
SetStartMessage({});
}
HDF5ImageLocator::Location JFJochHDF5Reader::GetImageLocation(int64_t image_number) const {
if (image_number >= static_cast<int64_t>(number_of_images) || image_number < 0)
throw JFJochException(JFJochExceptionCategory::HDF5, "Image out of bounds");
return image_source_.Resolve(image_number);
}
bool JFJochHDF5Reader::ReadRawImage(int64_t image_number, JFJochReaderRawImage &ret) {
// Every worker thread of an offline run comes through here, and HDF5 lets only one of them in at
// a time. So ask HDF5 only where the image is - a chunk-index lookup - and read the bytes after
// dropping the lock, which is the part that takes any time and the part that parallelises.
std::optional<HDF5ImageSource::DirectChunk> chunk;
{
std::unique_lock ul(hdf5_mutex);
if (!active_metadata_)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Cannot load image if file not loaded");
auto loc = GetImageLocation(image_number);
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,
int64_t image_number,
bool update_dataset) {
std::unique_lock ul(hdf5_mutex);
(void) update_dataset;
if (!dataset)
return false;
if (!active_metadata_)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Cannot load image if file not loaded");
// Pixels from the shared image source, per-image metadata from the active snapshot.
auto loc = GetImageLocation(image_number);
message.image = image_source_.ReadImageAt(buffer, loc);
message.number = image_number;
active_metadata_->FillPerImage(message, image_number, dataset);
return true;
}
std::vector<HDF5DataSourceMessage> JFJochHDF5Reader::GetHDF5DataSource(uint64_t first_image,
std::optional<uint64_t> image_count,
uint64_t stride) const {
std::unique_lock ul(hdf5_mutex);
if (!active_metadata_)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Cannot generate HDF5 source mapping if file not loaded");
return image_source_.GetSourceMapping(first_image, image_count, number_of_images, stride);
}
StoredPixelFormat JFJochHDF5Reader::GetStoredPixelFormat() const {
std::unique_lock ul(hdf5_mutex);
if (!active_metadata_)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Cannot read stored pixel format if file not loaded");
return image_source_.GetStoredPixelFormat();
}
std::vector<IntegrationOutcome> JFJochHDF5Reader::ReadReflections(size_t start_image,
std::optional<size_t> end_image) const {
std::unique_lock ul(hdf5_mutex);
if (!active_metadata_)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Cannot read reflections if file not loaded");
return active_metadata_->ReadReflections(start_image, end_image);
}
std::vector<SpotToSave> JFJochHDF5Reader::ReadSpots(int64_t image) const {
std::unique_lock ul(hdf5_mutex);
if (!active_metadata_)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Cannot read spots if file not loaded");
return active_metadata_->ReadSpots(image);
}
bool JFJochHDF5Reader::HasStoredSpots() const {
std::unique_lock ul(hdf5_mutex);
return active_metadata_ && active_metadata_->HasSpots();
}
CompressedImage JFJochHDF5Reader::ReadCalibration(std::vector<uint8_t> &tmp, const std::string &name) const {
std::unique_lock ul(hdf5_mutex);
if (!active_metadata_)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "Master file not loaded");
return active_metadata_->ReadCalibration(tmp, name);
}
void JFJochHDF5Reader::RegisterSnapshot(const std::string &name, const std::string &master_path) {
std::unique_lock ul(hdf5_mutex);
if (!active_metadata_)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Open a dataset before registering a snapshot");
auto metadata = std::make_shared<HDF5MetadataSource>();
auto open_result = metadata->Open(master_path, default_experiment);
// A snapshot may cover a subset of the images (a sub-range or filtered reprocessing); its
// /entry/detector/number map (read in Open) ties each snapshot image back to an original one.
// It just must not claim more images than the dataset has.
if (open_result.number_of_images > number_of_images)
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,
"Snapshot has more images than the open dataset");
// Snapshot pixels come from the existing image source; its own (integrated) master holds the
// per-image metadata at the global index.
metadata->UseImageSourceForMetadata(nullptr);
snapshots_[name] = metadata;
}
void JFJochHDF5Reader::RemoveSnapshot(const std::string &name) {
std::unique_lock ul(hdf5_mutex);
if (name == "Original")
return; // the original file metadata is always kept
snapshots_.erase(name);
if (active_snapshot_ == name) {
active_metadata_ = snapshots_.at("Original");
active_snapshot_ = "Original";
SetStartMessage(active_metadata_->Dataset());
}
}
void JFJochHDF5Reader::SetActiveSnapshot(const std::string &name) {
std::unique_lock ul(hdf5_mutex);
auto it = snapshots_.find(name);
if (it == snapshots_.end())
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "Unknown snapshot: " + name);
active_metadata_ = it->second;
active_snapshot_ = name;
SetStartMessage(active_metadata_->Dataset());
}
std::vector<std::string> JFJochHDF5Reader::SnapshotNames() const {
std::unique_lock ul(hdf5_mutex);
std::vector<std::string> names;
names.reserve(snapshots_.size());
for (const auto &[name, _]: snapshots_)
names.push_back(name);
return names;
}
std::string JFJochHDF5Reader::ActiveSnapshot() const {
std::unique_lock ul(hdf5_mutex);
return active_snapshot_;
}
std::vector<std::pair<std::string, std::shared_ptr<const JFJochReaderDataset>>>
JFJochHDF5Reader::AllSnapshotDatasets() const {
std::unique_lock ul(hdf5_mutex);
std::vector<std::pair<std::string, std::shared_ptr<const JFJochReaderDataset>>> out;
out.reserve(snapshots_.size());
// "Original" first, then the rest, so overlay colours stay stable across updates.
if (auto it = snapshots_.find("Original"); it != snapshots_.end())
out.emplace_back(it->first, it->second->Dataset());
for (const auto &[name, source]: snapshots_)
if (name != "Original")
out.emplace_back(name, source->Dataset());
return out;
}