diff --git a/modules/dmrpp_module/Chunk.cc b/modules/dmrpp_module/Chunk.cc index 2078411654..cbe331031c 100644 --- a/modules/dmrpp_module/Chunk.cc +++ b/modules/dmrpp_module/Chunk.cc @@ -152,7 +152,7 @@ size_t chunk_write_data(void *buffer, size_t size, size_t nmemb, void *data) { BESDEBUG(MODULE, prolog << "BEGIN " << endl); size_t nbytes = size * nmemb; auto chunk = reinterpret_cast(data); - +//cerr<<"coming to chunk_write_data"<get_data_url(); BESDEBUG(MODULE, prolog << "chunk->get_data_url():" << data_url << endl); @@ -189,6 +189,8 @@ size_t chunk_write_data(void *buffer, size_t size, size_t nmemb, void *data) { // | bytes_read unsigned long long bytes_read = chunk->get_bytes_read(); +//cerr<<"bytes_read: "< chunk->get_rbuf_size()) { diff --git a/modules/dmrpp_module/CurlHandlePool.h b/modules/dmrpp_module/CurlHandlePool.h index 5528bb8d46..635c3a058b 100644 --- a/modules/dmrpp_module/CurlHandlePool.h +++ b/modules/dmrpp_module/CurlHandlePool.h @@ -57,6 +57,7 @@ class dmrpp_easy_handle { public: dmrpp_easy_handle(); + CURL * get_curl_handle() { return d_handle;} ~dmrpp_easy_handle(); diff --git a/modules/dmrpp_module/DMZ.cc b/modules/dmrpp_module/DMZ.cc index cbf65294bf..7faf4d869b 100644 --- a/modules/dmrpp_module/DMZ.cc +++ b/modules/dmrpp_module/DMZ.cc @@ -1060,7 +1060,7 @@ void DMZ::set_up_direct_io_flag_phase_2(D4Group * grp, BaseType *btp) { if (!chunk_less_dim) return; - // Another special case is that some chunks are only filled with the fvalues. This case cannot be handled by direct IO. + // Another special case is that some chunks are only filled with the fvalues. This case can be conditionally supported, so need to check. // First calculate the number of logical chunks. // Also up to this step, the size of chunk_dim_sizes must be the same as the size of dim_sizes. No need to double check. @@ -1072,10 +1072,13 @@ void DMZ::set_up_direct_io_flag_phase_2(D4Group * grp, BaseType *btp) { has_filled_chunks = true; // Filled chunks can be supported for the whole variable case. However, we also need to check if _FillValue attribute is - // defined in this variable. + // defined in this variable. This works since in fileout netCDF, NC_NOFILL is set. So for DIO, we don't need to care about + // the data values in those filled chunks. if (has_filled_chunks) { BESDEBUG(PARSER, prolog << "has_filled_chunks: " <name() << endl); + // To be consistent with the current dmrpp's fillvalue handling, check if the _FillValue attribute exists. + if (btp->attributes()->find("_FillValue")==nullptr) return; } diff --git a/modules/dmrpp_module/DmrppArray.cc b/modules/dmrpp_module/DmrppArray.cc index 2643d352b2..5e3ece1976 100644 --- a/modules/dmrpp_module/DmrppArray.cc +++ b/modules/dmrpp_module/DmrppArray.cc @@ -31,8 +31,11 @@ #include #include #include +#include #include +#include + #include #include #include @@ -62,6 +65,8 @@ #include "byteswap_compat.h" #include "float_byteswap.h" #include "vlsa_util.h" +#include "ThreadCount.h" +#include "HttpNames.h" // Used with BESDEBUG #define dmrpp_3 "dmrpp:3" @@ -657,6 +662,363 @@ void DmrppArray::read_one_chunk_dio() { memcpy(target_buffer, source_buffer, the_one_chunk->get_size()); } +static ThreadCount &transfer_thread_count() { + static ThreadCount tc(DmrppRequestHandler::d_max_transfer_threads); + return tc; +} + +struct CurlMultiTransfer { + shared_ptr super_chunk; + unique_ptr super_chunk_internal; + unique_ptr easy_handle{nullptr, [](dmrpp_easy_handle *h){ CurlHandlePool::release_handle(h); }}; + unsigned int re_try = 0; +}; + + +static unique_ptr prepare_super_chunk_transfer(const shared_ptr &sc, bool dio) { + + if (sc->get_d_is_read()) + return nullptr; + + auto tf = make_unique(); + tf->super_chunk = sc; + + sc->set_read_buffer(sc->get_size()); + + if (sc->get_non_contiguous_chunk_flag()) + sc->map_non_contiguous_chunks_to_buffer(); + else + sc->map_chunks_to_buffer(); + //if (!dio && sc->get_uses_fill_value()) + // sc->read_fill_value_chunk(); + //else { + tf->super_chunk_internal = make_unique(sc->get_data_url(), "NOT_USED", sc->get_size(), sc->get_offset()); + tf->super_chunk_internal->set_read_buffer(sc->get_read_buffer(), sc->get_size(),0, false); + auto *curl_handle = DmrppRequestHandler::curl_handle_pool->get_easy_handle(tf->super_chunk_internal.get()); + if (!curl_handle) + throw BESInternalError(prolog + "No more libcurl handles.", __FILE__, __LINE__); + tf->easy_handle.reset(curl_handle); + //} + + return tf; +} + +// There may be 0.1% S3 failure rate, so we need to retry to see if we can obtain the data for a SuperChunk. +enum class ParallelTransferStatus { PT_SUCCESS, PT_RETRYABLE, PT_FAILURE }; + +static constexpr unsigned int MAX_ATTEMPTS = 3; +static constexpr std::chrono::microseconds INITIAL_RETRY_BACKOFF{250000}; // 0.25s + +// Follow eval_curl_easy_perform_code() + eval_http_get_response() from CurlUtils.cc to check the retry results. +static ParallelTransferStatus PT_result(CURL *easy, const shared_ptr &data_url, CURLcode curl_code) { + // Mirrors dmrpp_easy_handle::read_data()'s protocol branch: only + // HTTP/HTTPS gets retried; everything else is single-attempt, matching + // the plain curl_easy_perform() call in the `else` branch of read_data(). + if (data_url->protocol() != HTTPS_PROTOCOL && data_url->protocol() != HTTP_PROTOCOL) + return curl_code == CURLE_OK ? ParallelTransferStatus::PT_SUCCESS : ParallelTransferStatus::PT_FAILURE; + + // --- mirrors eval_curl_easy_perform_code(): every non-OK curl-level + // code is treated as retryable there (just logged differently per + // case), so we do the same here. --- + if (curl_code != CURLE_OK) + return ParallelTransferStatus::PT_RETRYABLE; + + // --- mirrors eval_http_get_response() / process_http_code_helper() --- + long http_code = 0; + if (curl_easy_getinfo(easy, CURLINFO_RESPONSE_CODE, &http_code) != CURLE_OK) + return ParallelTransferStatus::PT_FAILURE; // couldn't even read the response code back + + if (http_code == 200 || http_code == 206) + return ParallelTransferStatus::PT_SUCCESS; + + switch (http_code) { + case 400: case 401: case 402: case 403: case 404: case 408: + return ParallelTransferStatus::PT_FAILURE; + case 422: case 500: case 502: case 503: case 504: + // May need to check if the url is listed as retryable url. + return ParallelTransferStatus::PT_RETRYABLE; + default: + return ParallelTransferStatus::PT_FAILURE; + } +} + + +template +static void finish_super_chunk_transfer(CurlMultiTransfer &transfer, SendOffFn send_off) { + + + auto &sc = transfer.super_chunk; + if (sc->get_size() != transfer.super_chunk_internal->get_bytes_read()) { + ostringstream oss; + oss << prolog << "Wrong number of bytes read for SuperChunk " << sc->id() + << "; read: " << transfer.super_chunk_internal->get_bytes_read() << ", expected: " << sc->get_size(); + throw BESInternalError(oss.str(), __FILE__, __LINE__); + } + send_off(sc); + +} + +struct RetryCandidate { + unique_ptr tf; + chrono::steady_clock::time_point not_before; +}; + + +// Pass the parameter dio because for dio read the filled chunks won't be called. +template +void read_super_chunks_concurrent_curl_multi(queue> &super_chunks, + SendOffFn sof, bool dio) { + +struct timeval tv,tv2; +gettimeofday(&tv,NULL); + + const unsigned long max_threads = DmrppRequestHandler::d_max_transfer_threads; + CURLM *curl_multi = curl_multi_init(); + if (!curl_multi) + throw BESInternalError("curl_multi_init() failed.", __FILE__, __LINE__); + + unordered_map> multi_transfer_maps; + multi_transfer_maps.reserve(max_threads); + vector retries; + + auto cleanup = [&]() { + for (auto &handle : multi_transfer_maps) + curl_multi_remove_handle(curl_multi, handle.first); + multi_transfer_maps.clear(); + retries.clear(); + curl_multi_cleanup(curl_multi); + }; + + try { + int still_running = 0; + while (!super_chunks.empty() || !multi_transfer_maps.empty() || !retries.empty()) { + + // Mark the current time + auto time_now = chrono::steady_clock::now(); + + //Check retries + for (auto it = retries.begin(); it != retries.end(); ) { + if (multi_transfer_maps.size() >= max_threads) + break; + if (it->not_before <= time_now) { + CURL *curl_handle = it->tf->easy_handle->get_curl_handle(); + CURLMcode mc = curl_multi_add_handle(curl_multi, curl_handle); + if (mc != CURLM_OK) + throw BESInternalError(prolog + "curl_multi_add_handle() failed: " + + curl_multi_strerror(mc), __FILE__, __LINE__); + + multi_transfer_maps.emplace(curl_handle, std::move(it->tf)); + + it = retries.erase(it); + } + else { + ++it; + } + } + + while (multi_transfer_maps.size() < max_threads && !super_chunks.empty()) { + + auto sc = super_chunks.front(); + super_chunks.pop(); + + + if(!dio && sc->get_uses_fill_value()) { +//cerr<<"coming to filled chunks"<read_fill_value_chunk(); + sof(sc); + continue; + } + auto transfer = prepare_super_chunk_transfer(sc, dio); + + CURL *curl_handle = transfer->easy_handle->get_curl_handle(); + + //if (curl_handle) { +//cerr<<"coming to curl_handle check" < 0 || !retries.empty()) { + + int wait_ms = 1000; + if (!retries.empty()) { + auto earliest = min_element(retries.begin(), retries.end(), + [](const RetryCandidate &a, const RetryCandidate &b) { return a.not_before < b.not_before; }); + auto until = chrono::duration_cast(earliest->not_before - time_now).count(); + wait_ms = static_cast(max(0, min(until, 1000))); + } + + int numfds = 0; +#if LIBCURL_VERSION_NUM >= 0x074200 + mc = curl_multi_poll(curl_multi, nullptr, 0, wait_ms, &numfds); +#else + mc = curl_multi_wait(curl_multi, nullptr, 0, wait_ms, &numfds); +#endif + if (mc != CURLM_OK) + throw BESInternalError(prolog + "curl_multi_poll() failed: " + + curl_multi_strerror(mc), __FILE__, __LINE__); + } + + + // Wrap up + int msgs_left = 0; + CURLMsg *msg = nullptr; + while ((msg = curl_multi_info_read(curl_multi, &msgs_left)) != nullptr) { + if (msg->msg != CURLMSG_DONE) + continue; + + CURL *easy = msg->easy_handle; + auto it = multi_transfer_maps.find(easy); + if (it == multi_transfer_maps.end()) + continue; + + CURLcode result = msg->data.result; + curl_multi_remove_handle(curl_multi, easy); + + //finish_super_chunk_transfer(*it->second, sof, result); + unique_ptr transfer = std::move(it->second); + multi_transfer_maps.erase(it); + ParallelTransferStatus pt_status = PT_result(easy, transfer->super_chunk->get_data_url(), result); + + switch (pt_status) { + case ParallelTransferStatus::PT_SUCCESS: + // May itself throw (bytes-read mismatch) -- let it + // propagate to the outer catch, same as before. + finish_super_chunk_transfer(*transfer, sof); + break; + case ParallelTransferStatus::PT_RETRYABLE: { + ++transfer->re_try; + stringstream msg; + msg <<"Attempt to retry the data transfer, re_try "<re_try <<" times."<re_try >= MAX_ATTEMPTS) { + throw BESInternalError(prolog + "Made " + std::to_string(transfer->re_try) + + " failed attempts to retrieve SuperChunk " + + transfer->super_chunk->id() + ". Giving up.", + __FILE__, __LINE__); + } + // reset the bytes read before this handle is re-used. + transfer->super_chunk_internal->set_bytes_read(0); + + // This reuses the same CURL* (Range header, callbacks, + // etc. are already set correctly for the same request) + // by re-adding it to the multi handle -- curl_multi + // supports removing and re-adding a handle to retry a + // transfer. If you'd rather not trust that the handle's + // internal state is fully clean after a failed transfer, + // the more conservative alternative is to discard it and + // call DmrppRequestHandler::curl_handle_pool->get_easy_handle() + // again for a fresh one, same as begin_super_chunk_fetch() does. + auto backoff = INITIAL_RETRY_BACKOFF * (1u << (transfer->re_try - 1)); // 0.25s, 0.5s, ... + retries.push_back({std::move(transfer), time_now + backoff}); + break; + } + case ParallelTransferStatus::PT_FAILURE: + throw BESInternalError(prolog + "Data transfer error for SuperChunk " + + transfer->super_chunk->id(), __FILE__, __LINE__); + } + } + } + } catch (...) { + cleanup(); + throw; + } + + + cleanup(); + +gettimeofday(&tv2,NULL); + long seconds = tv2.tv_sec - tv.tv_sec; + long useconds = tv2.tv_usec -tv.tv_usec; + double elapsed = seconds *1000.0 + useconds/1000.0; + stringstream msg; +msg <<"Parallel data transfer Execution time: " << elapsed <<" ms"< +void read_super_chunks_concurrent_internal(queue> &super_chunks, + SendOffFn sof) { +struct timeval tv,tv2; +gettimeofday(&tv,NULL); + + auto &thread_store = transfer_thread_count(); + + list> futures; + try { + while (!super_chunks.empty() || !futures.empty()) { + while (!super_chunks.empty()) { + auto super_chunk = super_chunks.front(); + bool started = thread_store.start_future(futures, [super_chunk, sof]() -> bool { + sof(super_chunk); + return true; + }); + if (!started) + break; // The store is full; wait for a free spot. + + super_chunks.pop(); + } + + if (!futures.empty()) + thread_store.wait_for_one(futures, std::chrono::milliseconds(DMRPP_WAIT_FOR_FUTURE_MS)); + } + } catch (...) { + // Before exception, clean up all the running threads. + thread_store.release_all_threads(futures); + throw; + } +gettimeofday(&tv2,NULL); + long seconds = tv2.tv_sec - tv.tv_sec; + long useconds = tv2.tv_usec -tv.tv_usec; + double elapsed = seconds *1000.0 + useconds/1000.0; + stringstream msg; +msg <<"Parallel data transfer Execution time: " << elapsed <<" ms"<> &super_chunks) { + read_super_chunks_concurrent_curl_multi(super_chunks,[](const shared_ptr &sc) { sc->read_curl_multi(); }, + false); +} + +void read_super_chunks_unconstrained_concurrent(queue> &super_chunks) { + read_super_chunks_concurrent_curl_multi(super_chunks,[](const shared_ptr &sc) { sc->read_unconstrained_curl_multi(); }, + false); +} + + +void read_super_chunks_dio_concurrent(queue> &super_chunks) { + read_super_chunks_concurrent_curl_multi(super_chunks,[](const shared_ptr &sc) { sc->read_dio_curl_multi(); }, + true); +} + + +#if 0 +void read_super_chunks_concurrent(queue> &super_chunks) { + read_super_chunks_concurrent_internal(super_chunks,[](const shared_ptr &sc) { sc->read(); }); +} + +void read_super_chunks_unconstrained_concurrent(queue> &super_chunks) { + read_super_chunks_concurrent_internal(super_chunks,[](const shared_ptr &sc) { sc->read_unconstrained(); }); +} + +void read_super_chunks_dio_concurrent(queue> &super_chunks) { + read_super_chunks_concurrent_internal(super_chunks, [](const shared_ptr &sc) { sc->read_dio(); }); +} +#endif + /** * @brief Insert a chunk into an unconstrained Array * @@ -798,6 +1160,7 @@ void DmrppArray::read_chunks_unconstrained() { // The size, in elements, of each of the chunk's dimensions const vector chunk_shape = get_chunk_dimension_sizes(); + if (!DmrppRequestHandler::d_use_transfer_threads || super_chunks.size() == 1) { #if DMRPP_ENABLE_THREAD_TIMERS BES_STOPWATCH_START(dmrpp_3, prolog + "Serial SuperChunk Processing."); #endif @@ -807,6 +1170,10 @@ void DmrppArray::read_chunks_unconstrained() { BESDEBUG(dmrpp_3, prolog << super_chunk->to_string(true) << endl); super_chunk->read_unconstrained(); } + } + else { + read_super_chunks_unconstrained_concurrent(super_chunks); + } if (is_readable_struct) read_array_of_structure(d_structure_array_buf); @@ -830,6 +1197,11 @@ void DmrppArray::read_chunks_dio_unconstrained() { // The size, in elements, of each of the chunk's dimensions const vector chunk_shape = get_chunk_dimension_sizes(); + if (!DmrppRequestHandler::d_use_transfer_threads || super_chunks.size() == 1) { +struct timeval tv,tv2; +gettimeofday(&tv,NULL); + + #if DMRPP_ENABLE_THREAD_TIMERS BES_STOPWATCH_START(dmrpp_3, prolog + "Serial SuperChunk Processing."); #endif @@ -841,6 +1213,17 @@ void DmrppArray::read_chunks_dio_unconstrained() { // Call direct IO routine super_chunk->read_dio(); } +gettimeofday(&tv2,NULL); + long seconds = tv2.tv_sec - tv.tv_sec; + long useconds = tv2.tv_usec -tv.tv_usec; + double elapsed = seconds *1000.0 + useconds/1000.0; + stringstream msg; +msg <<" Sequential data transfer Execution time: " << elapsed <<" ms"<read_dio(); } + } + else + read_super_chunks_dio_concurrent(super_chunks); set_read_p(true); } @@ -1667,6 +2054,7 @@ void DmrppArray::read_chunks() { // This version is the 'serial' version of the code. It reads a chunk, inserts it, // reads the next one, and so on. + if (!DmrppRequestHandler::d_use_transfer_threads || super_chunks.size() == 1) { #if DMRPP_ENABLE_THREAD_TIMERS BES_STOPWATCH_START(dmrpp_3, prolog + "Serial SuperChunk Processing."); #endif @@ -1676,6 +2064,10 @@ void DmrppArray::read_chunks() { BESDEBUG(dmrpp_3, prolog << super_chunk->to_string(true) << endl); super_chunk->read(); } + } + else { + read_super_chunks_concurrent(super_chunks); + } if (is_readable_struct) read_array_of_structure(d_structure_array_buf); @@ -1722,6 +2114,7 @@ void DmrppArray::read_chunks_dio_constrained() { // Use the same approach as the non-direct chunk IO code for parallel transfer for now. // The parallel transfer part is the same as the non-direct chunk IO code. + if (!DmrppRequestHandler::d_use_transfer_threads || super_chunks.size() == 1) { #if DMRPP_ENABLE_THREAD_TIMERS BES_STOPWATCH_START(dmrpp_3, prolog + "Serial SuperChunk Processing."); #endif @@ -1732,6 +2125,9 @@ void DmrppArray::read_chunks_dio_constrained() { // For the direct IO, the unconstrained and constrained cases are the same. super_chunk->read_dio(); } + } + else + read_super_chunks_dio_concurrent(super_chunks); set_read_p(true); } @@ -1830,11 +2226,16 @@ void DmrppArray::read_buffer_chunks() { reserve_value_capacity_ll(get_size(true)); + if (!DmrppRequestHandler::d_use_transfer_threads || super_chunks.size() == 1) { while (!super_chunks.empty()) { auto super_chunk = super_chunks.front(); super_chunks.pop(); super_chunk->read(); } + } + else { + read_super_chunks_concurrent(super_chunks); + } set_read_p(true); } @@ -1922,12 +2323,18 @@ void DmrppArray::read_buffer_chunks_dio_constrained() { reserve_value_capacity_ll(get_var_chunks_storage_size()); + if (!DmrppRequestHandler::d_use_transfer_threads || super_chunks.size() == 1) { while (!super_chunks.empty()) { auto super_chunk = super_chunks.front(); super_chunks.pop(); // For the direct IO, the unconstrained and constrained cases are the same. super_chunk->read_dio(); } + } + else { + read_super_chunks_dio_concurrent(super_chunks); + + } set_read_p(true); } @@ -2942,11 +3349,16 @@ void DmrppArray::read_buffer_chunks_unconstrained() { reserve_value_capacity_ll(get_size()); + if (!DmrppRequestHandler::d_use_transfer_threads || super_chunks.size() == 1) { while (!super_chunks.empty()) { auto super_chunk = super_chunks.front(); super_chunks.pop(); super_chunk->read_unconstrained(); } + } + else { + read_super_chunks_unconstrained_concurrent(super_chunks); + } set_read_p(true); } diff --git a/modules/dmrpp_module/DmrppRequestHandler.cc b/modules/dmrpp_module/DmrppRequestHandler.cc index 369b2b4f99..90ad035a63 100644 --- a/modules/dmrpp_module/DmrppRequestHandler.cc +++ b/modules/dmrpp_module/DmrppRequestHandler.cc @@ -113,6 +113,10 @@ namespace dmrpp int DmrppRequestHandler::d_object_cache_entries = 100; double DmrppRequestHandler::d_object_cache_purge_level = 0.2; + bool DmrppRequestHandler::d_use_transfer_threads = false; + unsigned long DmrppRequestHandler::d_max_transfer_threads = 8UL; + unsigned long DmrppRequestHandler::d_default_max_transfer_threads = 4UL; + bool DmrppRequestHandler::d_use_compute_threads = true; unsigned long DmrppRequestHandler::d_max_compute_threads = 8UL; unsigned long DmrppRequestHandler::d_default_max_compute_threads = 4UL; @@ -155,6 +159,22 @@ namespace dmrpp stringstream msg; // Using TheBESKeys calls instead of custom functions kln 4/1/26 + d_use_transfer_threads = TheBESKeys::read_bool_key(DMRPP_USE_TRANSFER_THREADS_KEY, d_use_transfer_threads); + d_max_transfer_threads = std::max(d_default_max_transfer_threads, TheBESKeys::read_ulong_key(DMRPP_MAX_TRANSFER_THREADS_KEY, d_max_transfer_threads)); + msg << prolog << "Concurrent Transfer Threads: "; + if (DmrppRequestHandler::d_use_transfer_threads) + { + msg << "Enabled. max_transfer_threads: " << DmrppRequestHandler::d_max_transfer_threads << endl; + } + else + { + msg << "Disabled." << endl; + } + + INFO_LOG(msg.str()); + msg.str(std::string()); + + d_use_compute_threads = TheBESKeys::read_bool_key(DMRPP_USE_COMPUTE_THREADS_KEY, d_use_compute_threads); d_max_compute_threads = std::max(d_default_max_compute_threads, TheBESKeys::read_ulong_key(DMRPP_MAX_COMPUTE_THREADS_KEY, d_max_compute_threads)); msg << prolog << "Concurrent Compute Threads: "; diff --git a/modules/dmrpp_module/DmrppRequestHandler.h b/modules/dmrpp_module/DmrppRequestHandler.h index 2b51d2ed67..34a4ceb48b 100644 --- a/modules/dmrpp_module/DmrppRequestHandler.h +++ b/modules/dmrpp_module/DmrppRequestHandler.h @@ -77,6 +77,10 @@ class DmrppRequestHandler: public BESRequestHandler { static CurlHandlePool *curl_handle_pool; + static bool d_use_transfer_threads; + static unsigned long d_max_transfer_threads; + static unsigned long d_default_max_transfer_threads; + static bool d_use_compute_threads; static unsigned long d_max_compute_threads; static unsigned long d_default_max_compute_threads; diff --git a/modules/dmrpp_module/Makefile.am b/modules/dmrpp_module/Makefile.am index 241274b568..e01f8835de 100644 --- a/modules/dmrpp_module/Makefile.am +++ b/modules/dmrpp_module/Makefile.am @@ -50,7 +50,7 @@ DmrppFloat32.cc DmrppFloat64.cc DmrppInt16.cc DmrppInt32.cc DmrppInt64.cc \ DmrppInt8.cc DmrppUInt16.cc DmrppUInt32.cc DmrppUInt64.cc DmrppStr.cc \ DmrppStructure.cc DmrppUrl.cc DmrppD4Enum.cc DmrppD4Group.cc DmrppD4Opaque.cc \ DmrppD4Sequence.cc DmrppTypeFactory.cc DmrppMetadataStore.cc \ -SuperChunk.cc DMZ.cc vlsa_util.cc float_byteswap.cc +SuperChunk.cc DMZ.cc vlsa_util.cc float_byteswap.cc BES_HDRS = DMRpp.h DmrppCommon.h Chunk.h CurlHandlePool.h DmrppByte.h \ DmrppArray.h DmrppFloat32.h DmrppFloat64.h DmrppInt16.h DmrppInt32.h \ @@ -59,7 +59,7 @@ DmrppStr.h DmrppStructure.h DmrppUrl.h DmrppD4Enum.h DmrppD4Group.h \ DmrppD4Opaque.h DmrppD4Sequence.h DmrppTypeFactory.h \ DmrppMetadataStore.h DmrppNames.h byteswap_compat.h \ SuperChunk.h Base64.h DMZ.h DmrppChunkOdometer.h UnsupportedTypeException.h \ -vlsa_util.h float_byteswap.h +vlsa_util.h float_byteswap.h ThreadCount.h DMRPP_MODULE = DmrppModule.cc DmrppRequestHandler.cc DmrppModule.h DmrppRequestHandler.h diff --git a/modules/dmrpp_module/SuperChunk.cc b/modules/dmrpp_module/SuperChunk.cc index 4f7f6dfdc2..0d1a87f1bd 100644 --- a/modules/dmrpp_module/SuperChunk.cc +++ b/modules/dmrpp_module/SuperChunk.cc @@ -958,4 +958,122 @@ void SuperChunk::read_dio() { } } +void SuperChunk::read_curl_multi() { + + for (const auto &chunk : d_chunks) { + chunk->set_is_read(true); + chunk->set_bytes_read(chunk->get_size()); + } + BES_PROFILE_TIMING( + string("Handle SuperChunk data constrained - ") + (get_data_url() ? get_data_url()->get_url_no_query() : "") + + string(" - ") + ::to_string(get_size()) + string(" byte(s) - ") + ::to_string(get_chunk_count()) + + " chunk(s) - Using multithreading: " + (DmrppRequestHandler::d_use_compute_threads ? "true" : "false")); + + vector constrained_array_shape = d_parent_array->get_shape(true); + BESDEBUG(SUPER_CHUNK_MODULE, prolog << "d_use_compute_threads: " + << (DmrppRequestHandler::d_use_compute_threads ? "true" : "false") << endl); + BESDEBUG(SUPER_CHUNK_MODULE, + prolog << "d_max_compute_threads: " << DmrppRequestHandler::d_max_compute_threads << endl); + + if (!DmrppRequestHandler::d_use_compute_threads) { +#if DMRPP_ENABLE_THREAD_TIMERS + BES_STOPWATCH_START(SUPER_CHUNK_MODULE, prolog + "Serial Chunk Processing. id: " + d_id); +#endif + for (const auto &chunk : d_chunks) { + process_one_chunk(chunk, d_parent_array, constrained_array_shape); + } + } else { +#if DMRPP_ENABLE_THREAD_TIMERS + stringstream timer_tag; + timer_tag << prolog << "Concurrent Chunk Processing. id: " << d_id; + BES_STOPWATCH_START(SUPER_CHUNK_MODULE, timer_tag.str()); +#endif + queue> chunks_to_process; + for (const auto &chunk : d_chunks) + chunks_to_process.push(chunk); + + process_chunks_concurrent(d_id, chunks_to_process, d_parent_array, constrained_array_shape); + } + +} + +void SuperChunk::read_unconstrained_curl_multi() { + + for (const auto &chunk : d_chunks) { + chunk->set_is_read(true); + chunk->set_bytes_read(chunk->get_size()); + } + + BES_PROFILE_TIMING( + string("Handle SuperChunk data unconstrained - ") + (get_data_url() ? get_data_url()->get_url_no_query() : "") + + string(" - ") + ::to_string(get_size()) + string(" byte(s) - ") + ::to_string(get_chunk_count()) + + " chunk(s) - Using multithreading: " + (DmrppRequestHandler::d_use_compute_threads ? "true" : "false")); + + // The size in element of each of the array's dimensions + const vector array_shape = d_parent_array->get_shape(true); + // The size, in elements, of each of the chunk's dimensions + const vector chunk_shape = d_parent_array->get_chunk_dimension_sizes(); + + if (!DmrppRequestHandler::d_use_compute_threads) { +#if DMRPP_ENABLE_THREAD_TIMERS + BES_STOPWATCH_START(SUPER_CHUNK_MODULE, prolog + "Serial Chunk Processing. sc_id: " + d_id); +#endif + for (const auto &chunk : d_chunks) { + process_one_chunk_unconstrained(chunk, chunk_shape, d_parent_array, array_shape); + } + } else { +#if DMRPP_ENABLE_THREAD_TIMERS + stringstream timer_tag; + timer_tag << prolog << "Concurrent Chunk Processing. sc_id: " << d_id; + BES_STOPWATCH_START(SUPER_CHUNK_MODULE, timer_tag.str()); +#endif + queue> chunks_to_process; + for (const auto &chunk : d_chunks) { + chunks_to_process.push(chunk); + } + process_chunks_unconstrained_concurrent(d_id, chunks_to_process, chunk_shape, d_parent_array, array_shape); + } + +} + +void SuperChunk::read_dio_curl_multi() { + + for (const auto &chunk : d_chunks) { + chunk->set_is_read(true); + chunk->set_bytes_read(chunk->get_size()); + } + + BES_PROFILE_TIMING( + string("Handle SuperChunk data dio - ") + (get_data_url() ? get_data_url()->get_url_no_query() : "") + + string(" - ") + ::to_string(get_size()) + string(" byte(s) - ") + ::to_string(get_chunk_count()) + + " chunk(s) - Using multithreading: " + (DmrppRequestHandler::d_use_compute_threads ? "true" : "false")); + + // The size in element of each of the array's dimensions + const vector array_shape = d_parent_array->get_shape(true); + // The size, in elements, of each of the chunk's dimensions + const vector chunk_shape = d_parent_array->get_chunk_dimension_sizes(); + + if (!DmrppRequestHandler::d_use_compute_threads) { +#if DMRPP_ENABLE_THREAD_TIMERS + BES_STOPWATCH_START(SUPER_CHUNK_MODULE, prolog + "Serial Chunk Processing. sc_id: " + d_id); +#endif + for (const auto &chunk : d_chunks) { + process_one_chunk_unconstrained_dio(chunk, chunk_shape, d_parent_array, array_shape); + } + } else { +#if DMRPP_ENABLE_THREAD_TIMERS + stringstream timer_tag; + timer_tag << prolog << "Concurrent Chunk Processing. sc_id: " << d_id; + BES_STOPWATCH_START(SUPER_CHUNK_MODULE, timer_tag.str()); +#endif + queue> chunks_to_process; + for (const auto &chunk : d_chunks) + chunks_to_process.push(chunk); + + process_chunks_unconstrained_concurrent_dio(d_id, chunks_to_process, chunk_shape, d_parent_array, array_shape); + } + + +} + } // namespace dmrpp diff --git a/modules/dmrpp_module/SuperChunk.h b/modules/dmrpp_module/SuperChunk.h index 8b50a1a279..b4f23fd820 100644 --- a/modules/dmrpp_module/SuperChunk.h +++ b/modules/dmrpp_module/SuperChunk.h @@ -59,13 +59,16 @@ class SuperChunk { bool non_contiguous_chunk{false}; +public: + bool is_contiguous(std::shared_ptr candidate_chunk); void map_chunks_to_buffer(); void map_non_contiguous_chunks_to_buffer(); void read_aggregate_bytes(); void read_fill_value_chunk(); + bool get_uses_fill_value() {return d_uses_fill_value; } + -public: // Make the sc_id an uint64 and not a string - the code uses sstream to make the value. jhrg 5/7/22 explicit SuperChunk(const std::string &sc_id, DmrppArray *parent = nullptr) : d_id(sc_id), d_parent_array(parent) {} @@ -80,6 +83,11 @@ class SuperChunk { virtual unsigned long long get_offset() const { return d_offset; } virtual size_t get_chunk_count() const { return d_chunks.size(); } std::vector> get_chunks() const { return d_chunks; } + char * get_read_buffer() { return d_read_buffer;} + bool get_d_is_read() { return d_is_read;} + void set_read_buffer(unsigned long long size) { + if(!d_read_buffer) + d_read_buffer = new char[size];} virtual void read() { process_child_chunks(); } @@ -93,6 +101,10 @@ class SuperChunk { virtual void process_child_chunks(); virtual void process_child_chunks_unconstrained(); + virtual void read_curl_multi(); + virtual void read_unconstrained_curl_multi(); + virtual void read_dio_curl_multi(); + virtual bool empty() const { return d_chunks.empty(); } void set_non_contiguous_chunk_flag(bool flag) { non_contiguous_chunk = flag; } diff --git a/modules/dmrpp_module/ThreadCount.h b/modules/dmrpp_module/ThreadCount.h new file mode 100644 index 0000000000..197eb3c6fe --- /dev/null +++ b/modules/dmrpp_module/ThreadCount.h @@ -0,0 +1,86 @@ +#ifndef _thread_count_h +#define _thread_count_h 1 + +#include +#include +#include +#include "BESInternalError.h" + +namespace dmrpp { + +class ThreadCount { + +private: + + unsigned int tc_count = 0; + unsigned int tc_max_num_threads; + std::mutex tc_mtx; + + void release_thread_slot() { + + unique_lock tc_lock(tc_mtx); + if (tc_count > 0) + tc_count--; + + } + +public: + + explicit ThreadCount(unsigned int max_num_threads): tc_max_num_threads(max_num_threads) {} + template bool start_future(std::list> &futures, ThreadTask task){ + + unique_lock tc_lock(tc_mtx); + if (tc_count >= tc_max_num_threads) + return false; + tc_count++; + tc_lock.unlock(); + + // Use a lambda and try to handle exception gracefully;tested and it works for a sample file. + futures.push_back(async(launch::async, + [this](ThreadTask t) -> bool { + struct ThreadExceptionGuard { + ThreadCount *tc; + ~ThreadExceptionGuard() { tc->release_thread_slot(); } + } finish_exit{this}; + return t(); + }, + std::move(task))); + return true; + + } + template bool wait_for_one(std::list> &futures, const std::chrono::duration &timeout) { + + while (!futures.empty()) { + for (auto it = futures.begin(); it != futures.end(); ++it) { + if (!it->valid()) { + futures.erase(it); + return true; + } + if (it->wait_for(timeout) != std::future_status::timeout) { + bool success = it->get(); + futures.erase(it); + if (!success) + throw BESInternalError("A parallel thread task failed.", __FILE__, __LINE__); + return true; + } + } + } + return false; + } + void release_all_threads(std::list> &futures) { + + while (!futures.empty()) { + try { + if (futures.back().valid()) + futures.back().get(); + } catch (...) { + // Thread is released already by ThreadExceptionGuard. + } + futures.pop_back(); + } + } + +}; + +} +#endif