|
MADNESS 0.10.1
|
cloud class More...
#include <cloud.h>
Classes | |
| struct | cloudtimer |
| struct | DistributeFunctor |
| functor to distribute/rank/node-replicate a function, passed in as a pointer to WorldObjectBase More... | |
| struct | is_tuple |
| struct | is_tuple< std::tuple< T... > > |
| struct | is_vector |
| struct | is_vector< std::vector< Q > > |
Public Types | |
| typedef std::any | cached_objT |
| typedef std::map< keyT, cached_objT > | cacheT |
| template<typename T > | |
| using | has_cloud_serialize = madness::meta::is_detected< member_cloud_serialize_t, T > |
| using | keyT = madness::archive::ContainerRecordOutputArchive::keyT |
| template<typename T > | |
| using | member_cloud_serialize_t = decltype(std::declval< T >().cloud_store(std::declval< World & >(), std::declval< Cloud & >())) |
| typedef Recordlist< keyT > | recordlistT |
| enum | StoragePolicy { StoreFunction , StoreFunctionPointer } |
| using | valueT = std::vector< unsigned char > |
Public Member Functions | |
| Cloud (madness::World &universe) | |
| ~Cloud () | |
| ProcessID | batch_owner (const keyT record) const |
| the owner of a batch record; a pmap lookup, no communication | |
| void | clear () |
| void | clear_cache (World &subworld) |
| void | clear_timings () |
| template<typename T > | |
| T | consuming_load (madness::World &world, recordlistT &recordlist) const |
| similar to load, but will consume the recordlist | |
| template<typename T , std::size_t NDIM> | |
| std::vector< Function< T, NDIM > > | deserialize_batch_p2p (madness::World &subworld, Future< batch_bytesT > fut, const keyT record, const bool cache_result=false) const |
| turn the bytes of a p2p transfer into the batch of functions | |
| void | distribute_targets (const DistributionType dt=Distributed) |
| distribute/node/rank replicate the targets of all world objects stored in the cloud | |
| template<typename T > | |
| std::enable_if< is_vector< T >::value, T >::type | do_load (World &world, recordlistT &recordlist) const |
| template<typename T > | |
| std::enable_if<!is_vector< T >::value, T >::type | do_load (World &world, recordlistT &recordlist) const |
| template<typename T , std::size_t NDIM> | |
| std::vector< Function< T, NDIM > > | fetch_batch_p2p (madness::World &subworld, const keyT record, const bool cache_result=false) const |
| fetch a batch stored by store_batch; resolves without MPI when this rank owns it | |
| template<typename T > | |
| T | forward_load (madness::World &world, recordlistT &recordlist) const |
| load a single object from the cloud, recordlist is consumed while loading elements | |
| nlohmann::json | gather_memory_statistics (World &universe) const |
| get size of the cloud container | |
| nlohmann::json | gather_timings (World &universe) const |
| DistributionType | get_replication_policy () const |
| is the cloud container replicated: per rank, per node, or distributed | |
| nlohmann::json | get_statistics (World &world) const |
| return a json object with the cloud settings and statistics | |
| StoragePolicy | get_storing_policy () const |
| storing policy refers to storing functions or pointers to functions | |
| template<typename T > | |
| T | load (madness::World &world, const recordlistT recordlist) const |
| load a single object from the cloud, recordlist is kept unchanged | |
| template<typename T > | |
| T | load_tuple (madness::World &world, recordlistT &recordlist) const |
| void | print_batch_owner_map (World &universe, const std::string &tag="") const |
| void | print_size (World &universe) |
| void | print_timings (World &universe) const |
| backwards compatibility | |
| void | register_batch_owner (const keyT record, const ProcessID owner) |
| Register the owner of a batch record; local map insert, no communication. | |
| void | replicate (const std::size_t chunk_size=INT_MAX) |
| void | replicate_according_to_policy (const std::size_t chunk_size=INT_MAX) |
| void | replicate_per_node (const std::size_t chunk_size=INT_MAX) |
| Future< batch_bytesT > | request_batch_bytes_async (const keyT record) const |
start fetching record from its owner; the trigger is in flight on return | |
| void | set_debug (bool value) |
| void | set_fence (bool value) |
| void | set_force_load_from_cache (bool value) |
| void | set_replication_policy (const DistributionType value) |
| is the cloud container replicated: per rank, per node, or distributed | |
| void | set_storing_policy (const StoragePolicy value) |
| storing policy refers to storing functions or pointers to functions | |
| template<typename T > | |
| recordlistT | store (madness::World &world, const T &source) |
| template<typename T , std::size_t NDIM> | |
| keyT | store_batch (madness::World &world, const std::vector< Function< T, NDIM > > &batch, const ProcessID owner, const keyT record, const bool fence=true) |
| Store a batch of functions as one owner-pinned record. | |
| template<typename T > | |
| recordlistT | store_other (madness::World &world, const std::vector< T > &source) |
| template<typename... Ts> | |
| recordlistT | store_tuple (World &world, const std::tuple< Ts... > &input) |
| store a tuple in multiple records | |
| bool | validate_replication_policy () const |
Static Public Member Functions | |
| static void | print_memory_statistics (const nlohmann::json stats) |
| static void | print_timings (const nlohmann::json timings) |
Public Attributes | |
| std::atomic< long > | copy_time =0l |
| std::atomic< long > | target_replication_time =0l |
| std::list< WorldObjectBase * > | world_object_base_list |
Private Types | |
| template<typename T > | |
| using | is_parallel_serializable_object = std::is_base_of< archive::ParallelSerializableObject, T > |
| template<typename T > | |
| using | is_world_constructible = std::is_constructible< T, World & > |
Private Member Functions | |
| template<typename T > | |
| T | allocator (World &world) const |
| template<typename T > | |
| void | cache (madness::World &world, const T &obj, const keyT &record) const |
| bool | is_cached (const keyT &key) const |
| bool | is_in_container (const keyT &key) const |
| checks if a (universe) container record is used | |
| template<typename T > | |
| T | load_from_cache (madness::World &world, const keyT &record) const |
| load an object from the cache, record is unchanged | |
| template<typename T > | |
| recordlistT | store_other (madness::World &world, const T &source) |
| valueT | try_get_local_batch_bytes (const keyT record) const |
| bytes of a batch record held by this rank, empty if it holds none | |
| std::pair< const unsigned char *, std::size_t > | try_get_local_batch_ptr (const keyT record) const |
| stable pointer and size of a local batch record, {nullptr,0} if this rank holds none | |
Private Attributes | |
| madness::WorldContainer< keyT, valueT > | batch_container |
| std::atomic< long > | batch_deserialize_time =0l |
| deserializing the bytes, microseconds | |
| std::atomic< long > | batch_find_time =0l |
| waiting on the p2p transfer, microseconds | |
| std::shared_ptr< CloudOwnerPmap< keyT > > | batch_pmap |
| std::atomic< long > | batch_store_time =0l |
| store_batch wall time, microseconds | |
| std::unique_ptr< BatchTransport > | batch_transport_ |
| constructed after batch_container so it is destroyed first, as WorldObject lifetimes require | |
| std::atomic< long > | cache_reads =0l |
| std::atomic< long > | cache_stores =0l |
| cacheT | cached_objects |
| DistributionType | cloud_replication_policy = Distributed |
| cloud is a container: replication policy for the cloud container: distributed, node-replicated, rank-replicated | |
| madness::WorldContainer< keyT, valueT > | container |
| bool | debug = false |
| prints debug output | |
| bool | dofence = true |
| fences after load/store | |
| bool | force_load_from_cache = false |
| forces load from cache (mainly for debugging) | |
| bool | is_replicated =false |
| if contents of the container are replicated | |
| recordlistT | local_list_of_container_keys |
| std::atomic< long > | reading_time =0l |
| std::atomic< long > | replication_time =0l |
| StoragePolicy | storage_policy = StoreFunctionPointer |
| are the functions (WorldObjects) stored in the cloud or only pointers to them | |
| bool | use_cache =true |
| std::atomic< long > | writing_time =0l |
| std::atomic< long > | writing_time1 =0l |
Friends | |
| class | BatchTransport |
| only the transport reads owner-local bytes; callers go through fetch_batch_p2p | |
| std::ostream & | operator<< (std::ostream &os, const StoragePolicy &sp) |
| std::string | to_string (const StoragePolicy sp) |
cloud class
store and load data to/from the cloud into arbitrary worlds
Distributed data is always bound to a certain world. If it needs to be present in another world it can be serialized to the cloud and deserialized from there again. For an example see test_cloud.cc
Data is stored into a distributed container living in the universe. During storing a (replicated) list of records is returned that can be used to find the data in the container. If a combined object (a vector, tuple, etc) is stored a list of records will be generated. When loading the data from the world the record list will be used to deserialize all stored objects.
Note that there must be a fence after the destruction of subworld containers, as in:
create subworlds { dcT(subworld) do work } subworld.gop.fence();
| typedef std::any madness::Cloud::cached_objT |
| typedef std::map<keyT, cached_objT> madness::Cloud::cacheT |
| using madness::Cloud::has_cloud_serialize = madness::meta::is_detected<member_cloud_serialize_t, T> |
|
private |
|
private |
| using madness::Cloud::member_cloud_serialize_t = decltype(std::declval<T>().cloud_store(std::declval<World&>(), std::declval<Cloud&>())) |
| typedef Recordlist<keyT> madness::Cloud::recordlistT |
| using madness::Cloud::valueT = std::vector<unsigned char> |
|
inline |
| [in] | universe | the universe world |
|
inline |
References cached_objects, local_list_of_container_keys, and madness::print().
|
inlineprivate |
the owner of a batch record; a pmap lookup, no communication
References batch_pmap.
Referenced by madness::Exchange< T, NDIM >::ExchangeImpl< T, NDIM >::MacroTaskExchangeSimple::fetch_batch(), madness::BatchTransport::request(), madness::Exchange< T, NDIM >::ExchangeImpl< T, NDIM >::MacroTaskExchangeSimple::sym_pipeline_advance(), and test_batch_store_and_fetch().
|
inlineprivate |
References cached_objects.
|
inline |
References batch_container, and container.
Referenced by madness::MacroTaskQ::run_all(), and test_custom_serialization().
|
inline |
References cached_objects, madness::WorldGopInterface::fence(), madness::World::gop, and local_list_of_container_keys.
Referenced by chunk_example(), main(), madness::MacroTaskQ::run_all(), simple_example(), test_batch_store_and_fetch(), test_copy_function_from_other_world_through_cloud(), test_pointer_to_funcimpl(), test_replication_policy(), test_tuple(), and test_twice().
|
inline |
References cache_reads, cache_stores, copy_time, reading_time, replication_time, target_replication_time, writing_time, and writing_time1.
Referenced by test_twice().
|
inline |
similar to load, but will consume the recordlist
| [in] | world | the subworld the objects are loaded to |
| [in] | recordlist | the list of records where the objects are stored |
References reading_time.
Referenced by madness::MacroTask< taskT >::MacroTaskInternal::get_output().
|
inline |
turn the bytes of a p2p transfer into the batch of functions
Blocks on fut only if the transfer has not landed yet. Runs in a task, which is where a missing record is reported so the failure is attributable.
| [in] | cache_result | default false: the cloud-side cache is not safe to keep across changes of the calling world, so opting in is the caller's decision |
References batch_deserialize_time, batch_find_time, madness::cache, madness::Future< T >::get(), is_cached(), MADNESS_CHECK_THROW, reading_time, use_cache, and madness::wall_time().
Referenced by test_batch_store_and_fetch().
|
inline |
distribute/node/rank replicate the targets of all world objects stored in the cloud
References madness::WorldGopInterface::fence(), madness::World::gop, and world_object_base_list.
Referenced by madness::MacroTaskQ::run_all().
|
inline |
load a vector from the cloud, pop records from recordlist
| [in,out] | world | destination world |
| [in,out] | recordlist | list of records to load from (reduced by the first few elements) |
References target().
|
inline |
load a single object from the cloud, pop record from recordlist
| [in,out] | world | destination world |
| [in,out] | recordlist | list of records to load from (reduced by the first element) |
References madness::cache, container, debug, force_load_from_cache, madness::World::id(), is_cached(), is_replicated, MADNESS_CHECK, MADNESS_EXCEPTION, madness::print(), storage_policy, StoreFunctionPointer, target(), and use_cache.
|
inline |
fetch a batch stored by store_batch; resolves without MPI when this rank owns it
References is_cached(), and request_batch_bytes_async().
Referenced by test_batch_fetch_chunked(), test_batch_fetch_of_missing_record(), and test_batch_store_and_fetch().
|
inline |
load a single object from the cloud, recordlist is consumed while loading elements
References target().
Referenced by madness::CCIntermediatePotentials::cloud_load(), madness::Info::cloud_load(), madness::CCPair::cloud_load(), custom_serialize_tester::cloud_load(), and load_tuple().
|
inline |
get size of the cloud container
References batch_container, container, madness::get_rss_usage_in_GB(), madness::World::gop, madness::WorldGopInterface::max(), madness::WorldGopInterface::min(), madness::World::size(), and madness::WorldGopInterface::sum().
Referenced by get_statistics(), print_size(), test_copy_function_from_other_world_through_cloud(), and test_replication_policy().
|
inline |
References cache_reads, cache_stores, copy_time, madness::World::gop, madness::WorldGopInterface::max(), reading_time, replication_time, madness::World::size(), madness::WorldGopInterface::sum(), target_replication_time, and writing_time.
Referenced by get_statistics(), print_timings(), madness::MacroTaskQ::run_all(), and test_twice().
|
inline |
is the cloud container replicated: per rank, per node, or distributed
References cloud_replication_policy.
Referenced by madness::MacroTaskQ::replicate_inputs().
|
inline |
return a json object with the cloud settings and statistics
References cached_objects, cloud_replication_policy, gather_memory_statistics(), gather_timings(), is_replicated, storage_policy, and to_string.
Referenced by madness::MacroTaskQ::run_all(), and test_batch_fetch_chunked().
|
inline |
storing policy refers to storing functions or pointers to functions
References storage_policy.
|
inlineprivate |
References cached_objects.
Referenced by deserialize_batch_p2p(), do_load(), and fetch_batch_p2p().
|
inlineprivate |
checks if a (universe) container record is used
currently implemented with a local copy of the recordlist, might be reimplemented with container.find(), which would include blocking communication.
References local_list_of_container_keys.
Referenced by store_other().
|
inline |
load a single object from the cloud, recordlist is kept unchanged
| [in] | world | the subworld the objects are loaded to |
| [in] | recordlist | the list of records where the objects are stored |
References reading_time.
Referenced by chunk_example(), main(), madness::MacroTask< taskT >::MacroTaskInternal::run(), simple_example(), test_copy_function_from_other_world_through_cloud(), test_custom_serialization(), test_custom_worldobject(), test_pointer_to_funcimpl(), test_replication_policy(), test_tuple(), and test_twice().
|
inlineprivate |
load an object from the cache, record is unchanged
References cache_reads, cached_objects, debug, madness::World::id(), MADNESS_EXCEPTION, madness::print(), madness::World::rank(), and target().
|
inline |
load a tuple from the cloud, pop records from recordlist
| [in,out] | world | destination world |
| [in,out] | recordlist | list of records to load from (reduced by the first few elements) |
References debug, forward_load(), madness::World::id(), madness::name(), target(), and madness::type().
|
inline |
dump the record->owner table on rank 0; pair with the caller's own task-to-rank print to check that task assignment and batch routing agree
References batch_pmap, and madness::World::rank().
|
inlinestatic |
References madness::print(), and stats.
Referenced by madness::MacroTaskQ::run_all().
|
inline |
References gather_memory_statistics(), is_replicated, madness::print(), madness::World::rank(), madness::World::size(), and stats.
Referenced by madness::MacroTaskQ::run_all().
|
inlinestatic |
References madness::print().
|
inline |
backwards compatibility
References gather_timings(), and madness::print_timings.
Referenced by madness::MacroTaskQ::run_all(), and test_twice().
Register the owner of a batch record; local map insert, no communication.
Collective in the same sense as store_batch: every rank must call it with an identical (record, owner) pair or fetches will route inconsistently. Separating registration from the payload lets all ranks replicate the routing while each owner stores only its own bytes, over a size-1 subworld.
References batch_pmap.
Referenced by madness::Exchange< T, NDIM >::ExchangeImpl< T, NDIM >::MacroTaskExchangeSimple::store_batches(), and test_batch_fetch_of_missing_record().
|
inline |
References madness::WorldMpiInterface::Bcast(), container, madness::cpu_time(), debug, madness::WorldGopInterface::fence(), madness::WorldContainer< keyT, valueT, hashfunT >::find(), madness::World::gop, is_replicated, MADNESS_CHECK, MADNESS_CHECK_THROW, madness::World::mpi, MPI_BYTE, madness::print(), madness::World::rank(), replication_time, and madness::World::size().
Referenced by chunk_example(), replicate_according_to_policy(), and madness::MacroTaskQ::replicate_inputs().
|
inline |
References cloud_replication_policy, container, madness::Distributed, MADNESS_EXCEPTION, madness::NodeReplicated, madness::RankReplicated, replicate(), and replicate_per_node().
Referenced by test_replication_policy().
|
inline |
References container, madness::cpu_time(), debug, madness::WorldGopInterface::fence(), madness::World::gop, is_replicated, MADNESS_CHECK_THROW, MADNESS_EXCEPTION, madness::print(), madness::World::rank(), and replication_time.
Referenced by replicate_according_to_policy(), and madness::MacroTaskQ::replicate_inputs().
|
inline |
start fetching record from its owner; the trigger is in flight on return
References batch_transport_.
Referenced by fetch_batch_p2p(), madness::Exchange< T, NDIM >::ExchangeImpl< T, NDIM >::MacroTaskExchangeSimple::sym_pipeline_advance(), and test_batch_store_and_fetch().
|
inline |
References debug.
Referenced by test_custom_serialization(), test_custom_worldobject(), and test_tuple().
|
inline |
References dofence.
|
inline |
References force_load_from_cache.
Referenced by main(), test_custom_worldobject(), and test_tuple().
|
inline |
is the cloud container replicated: per rank, per node, or distributed
References cloud_replication_policy, madness::RankReplicated, and use_cache.
Referenced by madness::MacroTaskQ::MacroTaskQ(), and test_replication_policy().
|
inline |
storing policy refers to storing functions or pointers to functions
References storage_policy.
Referenced by madness::MacroTaskQ::MacroTaskQ(), simple_example(), test_copy_function_from_other_world_through_cloud(), and test_replication_policy().
|
inline |
| [in] | world | presumably the universe |
References dofence, madness::WorldGopInterface::fence(), madness::World::gop, is_replicated, MADNESS_EXCEPTION, madness::print(), source(), store_other(), store_tuple(), and writing_time.
Referenced by chunk_example(), madness::CCIntermediatePotentials::cloud_store(), madness::Info::cloud_store(), madness::CCPair::cloud_store(), custom_serialize_tester::cloud_store(), main(), madness::MacroTask< taskT >::prepare_output_records(), simple_example(), store_tuple(), test_copy_function_from_other_world_through_cloud(), test_custom_serialization(), test_custom_worldobject(), test_pointer_to_funcimpl(), test_replication_policy(), test_tuple(), and test_twice().
|
inline |
Store a batch of functions as one owner-pinned record.
The whole vector – its size and each function – is serialized into a single record in the batch container and routed to owner, so one batch is one record with one owner. Must be called with identical (owner, record) on every rank of world.
| [in] | fence | false lets a caller storing many batches emit one fence for all of them; the collective gather inside the archive self-synchronizes |
References batch_container, batch_pmap, batch_store_time, dofence, madness::WorldGopInterface::fence(), madness::World::gop, is_replicated, MADNESS_EXCEPTION, madness::print(), and madness::archive::BaseParallelArchive< Archive >::set_dofence().
Referenced by madness::Exchange< T, NDIM >::ExchangeImpl< T, NDIM >::MacroTaskExchangeSimple::store_batches(), test_batch_fetch_chunked(), and test_batch_store_and_fetch().
|
inline |
|
inlineprivate |
References cache_stores, madness::Recordlist< keyT >::compute_record(), container, debug, dofence, madness::WorldGopInterface::fence(), madness::World::gop, is_in_container(), local_list_of_container_keys, madness::print(), madness::World::rank(), source(), storage_policy, StoreFunctionPointer, madness::type_name< T >::value(), world_object_base_list, and writing_time1.
Referenced by store(), and store_other().
|
inline |
bytes of a batch record held by this rank, empty if it holds none
Never throws: the callers are BatchTransport's comm-thread handlers, where an escaping exception bypasses every task-level handler and surfaces as an unattributable abort. A miss is reported to the requester instead, and raised there in task context.
Emptiness is a sound "not found" marker because a stored batch always begins with its serialized element count, so a present record is never zero bytes.
References batch_container.
Referenced by madness::BatchTransport::request().
|
inlineprivate |
stable pointer and size of a local batch record, {nullptr,0} if this rank holds none
Lets the comm thread Isend without copying the payload. The accessor lock is released on return, but the address stays valid because batch records are neither erased nor mutated between the store and the end of the consuming operation. Never throws, for the reason given on try_get_local_batch_bytes.
References batch_container, and madness::WorldContainer< keyT, valueT, hashfunT >::size().
Referenced by madness::BatchTransport::on_trigger().
|
inline |
References cloud_replication_policy, container, and madness::validate_distribution_type().
Referenced by test_replication_policy().
|
friend |
only the transport reads owner-local bytes; callers go through fetch_batch_p2p
|
friend |
|
friend |
Referenced by get_statistics().
|
mutableprivate |
Referenced by clear(), gather_memory_statistics(), store_batch(), try_get_local_batch_bytes(), and try_get_local_batch_ptr().
|
mutableprivate |
deserializing the bytes, microseconds
Referenced by deserialize_batch_p2p().
|
mutableprivate |
waiting on the p2p transfer, microseconds
Referenced by deserialize_batch_p2p().
|
private |
dedicated container for owner-pinned function batches; see store_batch / fetch_batch_p2p. Uses CloudOwnerPmap so each batch record lives on an explicitly chosen rank.
Referenced by batch_owner(), print_batch_owner_map(), register_batch_owner(), and store_batch().
|
mutableprivate |
store_batch wall time, microseconds
Referenced by store_batch().
|
private |
constructed after batch_container so it is destroyed first, as WorldObject lifetimes require
Referenced by request_batch_bytes_async().
|
mutableprivate |
Referenced by clear_timings(), gather_timings(), and load_from_cache().
|
mutableprivate |
Referenced by clear_timings(), gather_timings(), and store_other().
|
private |
Referenced by ~Cloud(), cache(), clear_cache(), get_statistics(), is_cached(), and load_from_cache().
|
private |
cloud is a container: replication policy for the cloud container: distributed, node-replicated, rank-replicated
Referenced by get_replication_policy(), get_statistics(), replicate_according_to_policy(), set_replication_policy(), and validate_replication_policy().
|
mutableprivate |
|
mutable |
|
private |
prints debug output
Referenced by do_load(), load_from_cache(), load_tuple(), replicate(), replicate_per_node(), set_debug(), store_other(), and store_other().
|
private |
fences after load/store
Referenced by set_fence(), store(), store_batch(), store_other(), and store_other().
|
private |
forces load from cache (mainly for debugging)
Referenced by do_load(), and set_force_load_from_cache().
|
private |
if contents of the container are replicated
Referenced by do_load(), get_statistics(), print_size(), replicate(), replicate_per_node(), store(), and store_batch().
|
private |
Referenced by ~Cloud(), clear_cache(), is_in_container(), and store_other().
|
mutableprivate |
Referenced by clear_timings(), consuming_load(), deserialize_batch_p2p(), gather_timings(), and load().
|
mutableprivate |
Referenced by clear_timings(), gather_timings(), replicate(), and replicate_per_node().
|
private |
are the functions (WorldObjects) stored in the cloud or only pointers to them
Referenced by do_load(), get_statistics(), get_storing_policy(), set_storing_policy(), and store_other().
|
mutable |
Referenced by clear_timings(), gather_timings(), and madness::MacroTaskQ::replicate_inputs().
|
private |
Referenced by deserialize_batch_p2p(), do_load(), and set_replication_policy().
| std::list<WorldObjectBase*> madness::Cloud::world_object_base_list |
Referenced by distribute_targets(), madness::MacroTaskQ::replicate_inputs(), and store_other().
|
mutableprivate |
Referenced by clear_timings(), gather_timings(), and store().
|
mutableprivate |
Referenced by clear_timings(), and store_other().