parallel_invoke(と同期オブジェクト/コンテナ)
もっとも単純でお手軽な並列アルゴリズムがparallel_invoke()。関数オブジェクトa, b, c……があるとき、parallel_invoke(a, b, c ……)するとa(), b(), c()……が並列に実行され、全部終わるのを待ってくれます(引数に与える関数オブジェクトは最大10個)。parallel_invokeを使って素数探しを実装すると……
#include <ppl.h>
using namespace std;
using namespace concurrency;
void prime_parallel_invoke(unsigned int limit, vector<unsigned int>& primes) {
limit /= 8;
auto prime_task = [&primes, limit](unsigned int bias) {
for ( unsigned int i = 0; i < limit; ++i) {
unsigned int n = nth_prime_candidate(i, bias);
if ( is_prime(n) ) primes.push_back(n);
}
};
// 8つのlambda式を並列に評価
parallel_invoke(
[&]() { prime_task(0); },
[&]() { prime_task(1); },
[&]() { prime_task(2); },
[&]() { prime_task(3); },
[&]() { prime_task(4); },
[&]() { prime_task(5); },
[&]() { prime_task(6); },
[&]() { prime_task(7); }
);
}
……申し訳ない。このコード、ちゃんと動いてくれません。問題は得られた素数をvector<unsigned int>primesにpush_back()しているトコロ。vector<T>などのSTLコンテナは、複数のスレッドから同時に操作したときの動作が保証されていないのです。なので1つのタスクがprimes.push_back()している間、他のタスクがprimes.push_back()しないよう、排他制御してあげにゃなりません。
// ちゃんと動く(1) : critical_section で排他制御を行う
void prime_parallel_invoke(unsigned int limit, vector<unsigned int>& primes) {
critical_section cs;
limit /= 8;
auto prime_task = [&primes, &cs, limit](unsigned int bias) {
for ( unsigned int i = 0; i < limit; ++i) {
unsigned int n = nth_prime_candidate(i, bias);
if ( is_prime(n) ) {
// guardが有効な間は手出し無用
critical_section::scoped_lock guard(cs);
primes.push_back(n);
}
}
};
parallel_invoke(
[&]() { prime_task(0); },
[&]() { prime_task(1); },
[&]() { prime_task(2); },
[&]() { prime_task(3); },
[&]() { prime_task(4); },
[&]() { prime_task(5); },
[&]() { prime_task(6); },
[&]() { prime_task(7); }
);
}
critical_section::scoped_lockインスタンスが有効である間、他のcritical_section::scoped_lockインスタンスの生成がブロックされ、それによって排他制御されるってカラクリです。
PPLの同期オブジェクトにはもう一つ、reader_writer_lockがあります。読み手(reader)用と書き手(writer)用の2種類のlockを作ることができ、
- 書き手がいなければ(他に読み手がいても)読み手が作れる
- 書き手/読み手の両方がいなければ書き手を作れる
という動作を行います。読み手は値を変更しないから複数作れてもいいじゃない、と。
同期オブジェクトを使った排他制御だけでなく、並列動作が可能なコンテナ:
- concurrent_vector
- concurrent_queue
- concurrent_unordered_map
- concurrent_unordered_multimap
- concurrent_unordered_set
- concurrent_unordered_multiset
が用意されています。これらコンテナは基本操作の際に排他制御を必要としません。parallel_invoke()とconcurrent_vector<>を使った版がコチラ。
#include <concurrent_vector.h>
using namespace concurrency;
void prime_parallel_invoke(unsigned int limit, vector<unsigned int>& primes) {
concurrent_vector<unsigned int> cprimes;
limit /= 8;
auto prime_task = [&cprimes, limit](unsigned int bias) {
for ( unsigned int i = 0; i < limit; ++i) {
unsigned int n = 30*i + bias;
if ( is_prime(n) ) cprimes.push_back(n);
}
};
parallel_invoke(
[&]() { prime_task( 1); },
[&]() { prime_task( 7); },
[&]() { prime_task(11); },
[&]() { prime_task(13); },
[&]() { prime_task(17); },
[&]() { prime_task(19); },
[&]() { prime_task(23); },
[&]() { prime_task(29); }
);
primes.assign(begin(cprimes), end(cprimes));
}
並列コンテナ使用上の注意:並列コンテナは基本操作に関してスレッド・セーフです。例えば1つのスレッドがイテレータを介して要素を読み出している最中に他のスレッドがpush_back()しても構いません。言い換えれば、並列コンテナは要素の再配置を行いません。std::vector<>とconcurrency::concurrent_vector<>に要素を追加し、各要素のアドレスを出力してみましょう。
#include <iostream>
#include <vector>
#include <concurrent_vector.h>
using namespace std;
using namespace concurrency;
int main() {
cout << "----- concurrent_vector is NOT linear" << endl;
cout << "std::vetctor<char>\tconcurrentcy::concurrent_vector<char>" << endl;
vector<char> v;
concurrent_vector<char> cv;
for ( int i = 0; i < 16; ++i ) {
v.push_back('!');
cv.push_back('!');
}
for ( int i = 0; i < 16; ++i ) {
cout << reinterpret_cast<void*>(&v.at(i)) << "\t\t" << reinterpret_cast<void*>(&cv.at(i)) << endl;
}
}
ごらんのとおり、vector<>の要素がきれいに等間隔で並んでいるのに対し、concurrent_vector<>ではアドレスが飛び飛びになっています。要素数の増加にともなって領域を拡張する際、要素の再配置を行わないからです。concurrent_vector<>の先頭要素のアドレスを配列代わりに使ってはいけないことになります。
