From 2317eed841e32034235d87584fba44e3191e886e Mon Sep 17 00:00:00 2001 From: leonarski_f Date: Wed, 29 Apr 2026 19:21:19 +0200 Subject: [PATCH] jfjoch_process: Increase performance by decompressing image outside of global lock, avoiding extra processing and avoiding extra compression (if writing) --- docs/CHANGELOG.md | 1 + reader/JFJochHDF5Reader.cpp | 53 +++++++++++++++++++++----------- reader/JFJochHDF5Reader.h | 5 +++- reader/JFJochHttpReader.cpp | 37 +++++++++++++++++++++++ reader/JFJochHttpReader.h | 2 ++ reader/JFJochReader.h | 2 ++ reader/JFJochReaderImage.h | 5 ++++ tests/JFJochReaderTest.cpp | 60 ++++++++++++++++++++++++++++++++++++- tools/jfjoch_process.cpp | 40 +++++++------------------ 9 files changed, 155 insertions(+), 50 deletions(-) diff --git a/docs/CHANGELOG.md b/docs/CHANGELOG.md index 9df9e967..2cb97a08 100644 --- a/docs/CHANGELOG.md +++ b/docs/CHANGELOG.md @@ -5,6 +5,7 @@ This is an UNSTABLE release. The release has significant modifications and bug f * jfjoch_broker: For DECTRIS detectors, ZeroMQ link is persistent, to save time for establishing new connection * jfjoch_broker: Minor bug fixes for rare conditions +* jfjoch_process: Significantly improve performance ### 1.0.0-rc.139 This is an UNSTABLE release. The release has significant modifications and bug fixes, if things go wrong, it is better to revert to 1.0.0-rc.132. diff --git a/reader/JFJochHDF5Reader.cpp b/reader/JFJochHDF5Reader.cpp index e7265a72..2668b2e8 100644 --- a/reader/JFJochHDF5Reader.cpp +++ b/reader/JFJochHDF5Reader.cpp @@ -101,7 +101,7 @@ void JFJochHDF5Reader::ReadFile(const std::string& filename) { std::unique_lock ul(hdf5_mutex); try { auto dataset = std::make_shared(); - master_file = std::make_unique(filename); + master_file = std::make_shared(filename); dataset->experiment = default_experiment; dataset->arm_date = master_file->GetString("/entry/start_time"); @@ -493,11 +493,43 @@ CompressedImage JFJochHDF5Reader::LoadImageDataset(std::vector &tmp, HD if (datatype.IsFloat()) throw JFJochException(JFJochExceptionCategory::InputParameterInvalid,"Float datasets not supported at this time"); - return {tmp, dim[1], dim[2], + return {tmp, dim[2], dim[1], CalcImageMode(datatype.GetElemSize(), datatype.IsFloat(), datatype.IsSigned()), algorithm}; } +std::pair, uint32_t> JFJochHDF5Reader::GetImageLocation(int64_t image_number) { + if (image_number >= number_of_images || image_number < 0) + throw JFJochException(JFJochExceptionCategory::HDF5, "Image out of bounds"); + + uint32_t image_id; + std::shared_ptr data_file; + + if (format == FileWriterFormat::NXmxLegacy) { + uint32_t file_id = image_number / images_per_file; + image_id = image_number % images_per_file; + data_file = std::make_shared(legacy_format_files.at(file_id)); + } else { + image_id = image_number; + data_file = master_file; + } + + return {std::move(data_file), image_id}; +} + +std::shared_ptr JFJochHDF5Reader::GetRawImage(int64_t image_number) { + std::unique_lock ul(hdf5_mutex); + + if (!master_file) + throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, + "Cannot load image if file not loaded"); + + auto [source_file, image_id] = GetImageLocation(image_number); + auto ret = std::make_shared(); + ret->image = LoadImageDataset(ret->image_buffer, *source_file, image_id); + return ret; +} + bool JFJochHDF5Reader::LoadImage_i(std::shared_ptr &dataset, DataMessage &message, std::vector &buffer, @@ -512,22 +544,7 @@ bool JFJochHDF5Reader::LoadImage_i(std::shared_ptr &dataset throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "Cannot load image if file not loaded"); - if (image_number >= number_of_images) - throw JFJochException(JFJochExceptionCategory::HDF5, "Image out of bounds"); - - std::unique_ptr tmp_data_file; - uint32_t image_id; - HDF5Object *source_file; - - if (format == FileWriterFormat::NXmxLegacy) { - uint32_t file_id = image_number / images_per_file; - image_id = image_number % images_per_file; - tmp_data_file = std::make_unique(legacy_format_files.at(file_id)); - source_file = tmp_data_file.get(); - } else { - image_id = image_number; - source_file = master_file.get(); - } + auto [source_file, image_id] = GetImageLocation(image_number); message.image = LoadImageDataset(buffer, *source_file, image_id); message.number = image_number; diff --git a/reader/JFJochHDF5Reader.h b/reader/JFJochHDF5Reader.h index a3455838..3068e243 100644 --- a/reader/JFJochHDF5Reader.h +++ b/reader/JFJochHDF5Reader.h @@ -10,7 +10,7 @@ class JFJochHDF5Reader : public JFJochReader { FileWriterFormat format = FileWriterFormat::NoFile; - std::unique_ptr master_file; + std::shared_ptr master_file; std::vector legacy_format_files; @@ -32,6 +32,7 @@ class JFJochHDF5Reader : public JFJochReader { size_t image0, size_t nimages); + std::pair, uint32_t> GetImageLocation(int64_t image_number); std::optional ReadAxis(HDF5Object *file, const std::string &name); public: ~JFJochHDF5Reader() override = default; @@ -42,6 +43,8 @@ public: void Close() override; CompressedImage ReadCalibration(std::vector &tmp, const std::string &name) const; + + std::shared_ptr GetRawImage(int64_t image_number) override; }; diff --git a/reader/JFJochHttpReader.cpp b/reader/JFJochHttpReader.cpp index 119089bf..08935fc5 100644 --- a/reader/JFJochHttpReader.cpp +++ b/reader/JFJochHttpReader.cpp @@ -269,6 +269,43 @@ bool JFJochHttpReader::LoadImage_i(std::shared_ptr &dataset } } +std::shared_ptr JFJochHttpReader::GetRawImage(int64_t image_number) { + if (addr.empty()) + return {}; + + httplib::Client cli_cmd(addr); + auto res = cli_cmd.Get("/image_buffer/image.cbor?id=" + std::to_string(image_number)); + + if (!res || res->status != httplib::StatusCode::OK_200 || res->body.empty()) + return {}; + + try { + auto msg = + CBORStream2Deserialize(reinterpret_cast(res->body.data()), res->body.size()); + + if (msg->msg_type != CBORImageType::IMAGE) + return {}; + + std::shared_ptr image = std::make_shared(); + image->image_buffer.resize(msg->data_message->image.GetCompressedSize()); + memcpy(image->image_buffer.data(), + msg->data_message->image.GetCompressed(), + msg->data_message->image.GetCompressedSize()); + image->image = CompressedImage( + image->image_buffer.data(), + image->image_buffer.size(), + msg->data_message->image.GetWidth(), + msg->data_message->image.GetHeight(), + msg->data_message->image.GetMode(), + msg->data_message->image.GetCompressionAlgorithm() + ); + + return image; + } catch (std::exception &e) { + return {}; + } +} + std::vector JFJochHttpReader::GetPlot_i(const std::string &plot_type, float fill_value) const { if (addr.empty()) return {}; diff --git a/reader/JFJochHttpReader.h b/reader/JFJochHttpReader.h index 8a322cd2..1e4dfb77 100644 --- a/reader/JFJochHttpReader.h +++ b/reader/JFJochHttpReader.h @@ -34,6 +34,8 @@ public: void UploadUserMask(const std::vector& mask); BrokerStatus GetBrokerStatus() const; + + std::shared_ptr GetRawImage(int64_t image_number) override; }; diff --git a/reader/JFJochReader.h b/reader/JFJochReader.h index 52a8224a..4c4ec5cd 100644 --- a/reader/JFJochReader.h +++ b/reader/JFJochReader.h @@ -47,6 +47,8 @@ public: void UpdateGeomMetadata(const DiffractionExperiment& experiment); void UpdateUserMask(const std::vector& mask); + + virtual std::shared_ptr GetRawImage(int64_t image_number) = 0; }; #endif //JFJOCHIMAGEREADER_H diff --git a/reader/JFJochReaderImage.h b/reader/JFJochReaderImage.h index f9d159a1..c0968fe8 100644 --- a/reader/JFJochReaderImage.h +++ b/reader/JFJochReaderImage.h @@ -19,6 +19,11 @@ constexpr static int32_t GAP_PXL_VALUE = INT32_MIN + 1; constexpr static int32_t ERROR_PXL_VALUE = INT32_MIN; constexpr static int32_t SATURATED_PXL_VALUE = INT32_MAX; +struct JFJochReaderRawImage { + std::vector image_buffer; + CompressedImage image; +}; + class JFJochReaderImage { std::shared_ptr dataset; diff --git a/tests/JFJochReaderTest.cpp b/tests/JFJochReaderTest.cpp index ec53eac9..221d43fe 100644 --- a/tests/JFJochReaderTest.cpp +++ b/tests/JFJochReaderTest.cpp @@ -163,7 +163,6 @@ TEST_CASE("JFJochReader_MasterFile_Calibration", "[HDF5][Full]") { REQUIRE(H5Fget_obj_count(H5F_OBJ_ALL, H5F_OBJ_ALL) == 0); } - TEST_CASE("JFJochReader_DefaultExperiment", "[HDF5][Full]") { DiffractionExperiment x(DetJF(1)); @@ -1495,5 +1494,64 @@ TEST_CASE("JFJochReader_NXmxIntegrated", "[HDF5][Full]") { } remove("test_reader_integrated_master.h5"); + REQUIRE(H5Fget_obj_count(H5F_OBJ_ALL, H5F_OBJ_ALL) == 0); +} + +TEST_CASE("JFJochReader_GetRawImage", "[HDF5][Full]") { + DiffractionExperiment x(DetJF(1)); + + x.FilePrefix("test_read_raw_image").ImagesPerTrigger(4).OverwriteExistingFiles(true); + x.BitDepthImage(16).ImagesPerFile(2).SetFileWriterFormat(FileWriterFormat::NXmxLegacy).PixelSigned(true) + .IndexingAlgorithm(IndexingAlgorithmEnum::FFT); + x.Compression(CompressionAlgorithm::BSHUF_ZSTD); + + std::vector image(x.GetPixelsNum()); + for (int i = 0; i < image.size(); i++) + image[i] = static_cast((i * 7 + 33) % UINT16_MAX); + + RegisterHDF5Filter(); + JFJochBitShuffleCompressor compressor(CompressionAlgorithm::BSHUF_ZSTD); + auto compressed_image = compressor.Compress(image); + { + StartMessage start_message; + x.FillMessage(start_message); + FileWriter file_set(start_message); + + for (int i = 0; i < x.GetImageNum(); i++) { + DataMessage message{}; + message.image = CompressedImage(compressed_image, x.GetXPixelsNum(), x.GetYPixelsNum(), + CompressedImageMode::Int16, CompressionAlgorithm::BSHUF_ZSTD); + message.number = i; + REQUIRE_NOTHROW(file_set.WriteHDF5(message)); + } + EndMessage end_message; + end_message.max_image_number = x.GetImageNum(); + file_set.WriteHDF5(end_message); + file_set.Finalize(); + } + { + JFJochHDF5Reader reader; + REQUIRE_NOTHROW(reader.ReadFile("test_read_raw_image_master.h5")); + auto dataset = reader.GetDataset(); + CHECK(dataset->experiment.GetImageNum() == 4); + + std::shared_ptr reader_image; + for (int i = 0; i < 4; i++) { + REQUIRE_NOTHROW(reader_image = reader.GetRawImage(i)); + + CHECK(reader_image->image.GetMode() == CompressedImageMode::Int16); + CHECK(reader_image->image.GetCompressionAlgorithm() == CompressionAlgorithm::BSHUF_ZSTD); + CHECK(reader_image->image.GetWidth() == x.GetXPixelsNum()); + CHECK(reader_image->image.GetHeight() == x.GetYPixelsNum()); + CHECK(reader_image->image.GetCompressedSize() == compressed_image.size()); + CHECK(reader_image->image.GetCompressed() == reader_image->image_buffer.data()); + REQUIRE(reader_image->image_buffer.size() == compressed_image.size()); + CHECK(memcmp(reader_image->image_buffer.data(), compressed_image.data(), compressed_image.size()) == 0); + } + } + remove("test_read_raw_image_master.h5"); + remove("test_read_raw_image_data_000001.h5"); + remove("test_read_raw_image_data_000002.h5"); +// No leftover HDF5 objects REQUIRE(H5Fget_obj_count(H5F_OBJ_ALL, H5F_OBJ_ALL) == 0); } \ No newline at end of file diff --git a/tools/jfjoch_process.cpp b/tools/jfjoch_process.cpp index 7d42d4f2..82d8c491 100644 --- a/tools/jfjoch_process.cpp +++ b/tools/jfjoch_process.cpp @@ -308,7 +308,8 @@ int main(int argc, char **argv) { experiment.PixelSigned(true); experiment.OverwriteExistingFiles(true); experiment.PolarizationFactor(0.99); - + experiment.SetFileWriterFormat(FileWriterFormat::NXmxIntegrated); + if (fixed_reference_unit_cell.has_value()) experiment.SetUnitCell(*fixed_reference_unit_cell); @@ -404,13 +405,6 @@ int main(int argc, char **argv) { std::atomic finished_count = 0; auto worker = [&](int thread_id) { - JFJochBitShuffleCompressor compressor(experiment.GetCompressionAlgorithm()); - std::vector compressed_buffer; - - compressed_buffer.resize(MaxCompressedSize(experiment.GetCompressionAlgorithm(), - experiment.GetPixelsNum(), - experiment.GetByteDepthImage())); - // Thread-local analysis resources MXAnalysisWithoutFPGA analysis(experiment, mapping, pixel_mask, indexer); @@ -423,9 +417,9 @@ int main(int argc, char **argv) { if (image_idx >= end_image) break; // Load Image - std::shared_ptr img; + std::shared_ptr img; try { - img = reader.LoadImage(image_idx); + img = reader.GetRawImage(image_idx); } catch (const std::exception &e) { logger.Error("Failed to load image {}: {}", image_idx, e.what()); continue; @@ -433,15 +427,12 @@ int main(int argc, char **argv) { if (!img) continue; - DataMessage msg; + DataMessage msg{}; + + msg.image = img->image; + msg.number = image_idx; + msg.image_collection_efficiency = dataset->efficiency[image_idx]; - msg.image = img->ImageData().image; - msg.number = img->ImageData().number; - msg.error_pixel_count = img->ImageData().error_pixel_count; - msg.image_collection_efficiency = img->ImageData().image_collection_efficiency; - msg.storage_cell = img->ImageData().storage_cell; - msg.user_data = img->ImageData().user_data; - msg.compression_time_s = img->ImageData().compression_time_s; total_uncompressed_bytes += msg.image.GetUncompressedSize(); auto image_start_time = std::chrono::high_resolution_clock::now(); @@ -464,19 +455,8 @@ int main(int argc, char **argv) { plots.Add(msg, profile); // Write Result - if (writer) { - auto size = compressor.Compress(compressed_buffer.data(), - img->Image().data(), - experiment.GetPixelsNum(), - sizeof(int32_t)); - - msg.image = CompressedImage(compressed_buffer.data(), - size, experiment.GetXPixelsNum(), - experiment.GetYPixelsNum(), - CompressedImageMode::Int32, - experiment.GetCompressionAlgorithm()); + if (writer) writer->Write(msg); - } // Update max sent tracking uint64_t current_max = max_image_number_sent.load();