Waiting for Single Messages in Flow Graph

Waiting for Single Messages in Flow Graph#

[flow_graph.wait_for_single_message]

Note

To enable this preview feature, define the TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT or TBB_PREVIEW_FLOW_GRAPH_FEATURES macro to 1.

This feature adds a new try_put_and_wait interface to the receiving nodes in the Flow Graph. This function puts a message as an input into a Flow Graph and waits until all work related to that message is complete. try_put_and_wait may reduce latency compared to calling graph::wait_for_all since graph::wait_for_all waits for all work, including work that is unrelated to the input message, to complete.

node.try_put_and_wait(msg) performs node.try_put(msg) on the node and waits until the work on msg is completed. Therefore, the following conditions are true:

  • Any task initiated by any node in the Flow Graph that involves working with msg or any other intermediate result computed from msg is completed.

  • No intermediate results computed from msg remain in any buffers in the graph.

Caution

To prevent try_put_and_wait calls from infinite waiting, avoid using buffering nodes at the end of the Flow Graph since the final result will not be automatically consumed by the Flow Graph.

Caution

The multifunction_node and async_node classes are not currently supported by this feature. Including one of these nodes in the Flow Graph may cause try_put_and_wait to exit early, even if the computations on the initial input message are still in progress.

The following APIs are extended by this feature:

  • continue_node

  • function_node

  • overwrite_node

  • write_once_node

  • buffer_node

  • queue_node

  • priority_queue_node

  • sequencer_node

  • limiter_node

  • broadcast_node

  • split_node

  • Input ports of join_node

  • Input ports of indexer_node

Member Functions#

bool node::try_put_and_wait(const Input& input);

Effects: Performs try_put(input). Additionally, waits for the completion of the computations related to input in the Flow Graph, meaning all tasks created by each node and related to input are executed, and no related objects remain in any buffer within the graph.

Example#

#define TBB_PREVIEW_FLOW_GRAPH_TRY_PUT_AND_WAIT 1
#include <oneapi/tbb/flow_graph.h>
#include <oneapi/tbb/parallel_for.h>
#include <tuple>

struct f1_body;
struct f2_body;
struct f3_body;
struct f4_body;

int main() {
    using namespace oneapi::tbb;

    flow::graph g;
    flow::broadcast_node<int> start_node(g);

    flow::function_node<int, int> f1(g, flow::unlimited, f1_body{});
    flow::function_node<int, int> f2(g, flow::unlimited, f2_body{});
    flow::function_node<int, int> f3(g, flow::unlimited, f3_body{});

    flow::join_node<std::tuple<int, int>> join(g);

    flow::function_node<std::tuple<int, int>, int> f4(g, flow::serial, f4_body{});

    flow::make_edge(start_node, f1);
    flow::make_edge(f1, f2);

    flow::make_edge(start_node, f3);

    flow::make_edge(f2, flow::input_port<0>(join));
    flow::make_edge(f3, flow::input_port<1>(join));

    flow::make_edge(join, f4);

    // Submit work into the graph
    parallel_for(0, 100, [&](int input) {
        start_node.try_put_and_wait(input);

        // Post processing the result of input
    });
}

Each iteration of parallel_for submits an input into the Flow Graph. After returning from try_put_and_wait(input), it is guaranteed that all of the work related to the completion of input is done by all of the nodes in the graph. Tasks related to inputs submitted by other calls are not guaranteed to be completed.