tag_matching
……と、持ち上げておいて落とすことになるんですけど、実はこのコード、バグっています。実行結果をながめてみると "adder(1,1000)" なんてのが混じっています。が、そんなバカな。x^2とx^3とを組にしてadderに渡すのだから"adder(1,1)"もしくは"adder(100,1000)"でなくてはならんはずです。
タスク・スケジューラは効率最優先でタスクの実行スケジュールを組むので、投入されたデータの順とその結果が後段nodeに流れ込む順とは一致しないのです。join_nodeは全入力ポートにデータが揃えばお構いなしにそれらを束ねて後段nodeに送り出しちゃいますからね。タスクの同時投入数を1に制限すれば解消するのでしょうが、それだとスレッドが遊んでいる時間が増えてせっかくのスピードを殺してしまいます。
解決策は用意されています。join_nodeにはtag_matchingという動作モードがあり、データにsize_t型のタグ(合札)を埋め込んでおけば、タグの一致するデータどうしを束ねてくれます。
#include <atomic>
#include <tbb/tbb.h>
#include "utils.h"
using namespace std;
using namespace tbb::flow;
// 扱うデータの型とformat
typedef slow_item<int> item;
#define ITEM_FORMAT "%d"
// tag付けされた型 tagged<T>
template<typename T> using tagged = tuple<size_t,T>;
// tagged<T> から タグ/値 を取り出す関数
template<typename T> inline size_t tag(const tagged<T>& t) { return get<0>(t); }
template<typename T> inline T val(const tagged<T>& t) { return get<1>(t); }
item calc(int lo, int hi) {
item result = 0;
atomic<size_t> tagval = 0U;
graph g;
// 入力にタグ付けする : x -> (t,x)
function_node<item,tagged<item>>
input( g, unlimited,
[&tagval](item v) {
return tagged<item>(++tagval,v);
}
);
broadcast_node<tagged<item>>
forker( g);
// 二乗を求める : (t,x) -> (t,x^2)
function_node<tagged<item>,tagged<item>>
squarer( g, unlimited,
[](const tagged<item>& t) {
item v = val(t);
trace("squarer(" ITEM_FORMAT ")\n", value(v));
return tagged<item>(tag(t), v * v);
}
);
// 三乗を求める : (t,x) -> (t,x^3)
function_node<tagged<item>,tagged<item>>
cuber( g, unlimited,
[](const tagged<item>& t) {
item v = val(t);
trace("cuber(" ITEM_FORMAT ")\n", value(v));
return tagged<item>(tag(t), v * v * v);
}
);
// 入力データからtagを抽出する関数オブジェクトをコンストラクタに与える
join_node<tuple<tagged<item>,tagged<item>>,tag_matching>
join( g, tag<item>, tag<item>);
// 和を求め、resultに積み上げる : ((t,a),(t,b)) -> Σa+b
function_node<tuple<tagged<item>,tagged<item>>,item>
summer( g, serial,
[&result](const tuple<tagged<item>,tagged<item>>& t) {
item v0 = val(get<0>(t));
item v1 = val(get<1>(t));
trace("summer(" ITEM_FORMAT "," ITEM_FORMAT ")\n", value(v0), value(v1));
return result += v0 + v1;
}
);
make_edge( input, forker);
make_edge( forker, squarer );
make_edge( forker, cuber );
make_edge( squarer, get<0>(join.input_ports()) );
make_edge( cuber, get<1>(join.input_ports()) );
make_edge( join, summer );
for (int i = lo; i <= hi; ++i) {
input.try_put(i);
}
g.wait_for_all();
return result;
}
int main() {
tbb::task_scheduler_init tbb_init(2);
item result;
measure([&]() { result = calc(1,10);});
cout << value(result) << endl;
}

