Two things every worker thread of an offline run did inside the global HDF5 mutex, per image. It opened /entry/data/data and asked it for its dataspace, its datatype and its creation plist, then asked those for the rank, the dimensions, the chunking and the compression. All of that is a property of the file and identical for all of its images, so it is now resolved once when the file is first touched. And it read the pixels - megabytes of them, with the lock held, which is what turned a worker per hardware thread into a queue. HDF5 can say where a chunk lives instead - address and byte count, a lookup in the chunk index with no read attached - so that is all it is asked for now, and the bytes are fetched after the lock is dropped, with a positional read that any number of threads can make through one handle at once. Chunk addresses count from the end of the user block, so its size is added; zero for anything this project writes, not for every file. A file that is not one chunk per image, or a chunk that was never written and exists only as a fill value, still goes the old way - only HDF5 knows what those read as. On a 16 Mpx rotation dataset with the process file being written, the per-image loop at 48 workers goes 12.4 s -> 6.8 s, and stops getting slower as workers are added: 8 workers were faster than 48 before, and are not now. Where no process file is written the same loop only improves ~1%, because this machine has 1.5 TB of RAM and held the whole 7 GB test set in page cache - the read was never the expensive part here. It is where the cache is cold or the filesystem is remote. Battery 9m45s, space group 21/24, no failures, unchanged. The Windows path uses ReadFile with an OVERLAPPED offset for the same reason pread is used elsewhere: it takes the offset as an argument rather than moving a shared file position, so the viewer keeps building under MSVC and gets the same concurrency. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
229 lines
8.7 KiB
C++
229 lines
8.7 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);
|
|
}
|
|
|
|
std::shared_ptr<JFJochReaderRawImage> JFJochHDF5Reader::GetRawImage(int64_t image_number) {
|
|
auto ret = std::make_shared<JFJochReaderRawImage>();
|
|
|
|
// 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);
|
|
return ret;
|
|
}
|
|
}
|
|
|
|
ret->image = HDF5ImageSource::ReadDirect(ret->image_buffer, *chunk);
|
|
return ret;
|
|
}
|
|
|
|
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);
|
|
}
|
|
|
|
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;
|
|
}
|