Performance comparison (C#, C++, Go)
Разбирая примеры из книги Адама Фримана Pro .NET 4 Parallel Programming in C#, нашел там следующее описание алгоритма быстрой сортировки:
Quicksort is a sorting algorithm that is well suited to parallelization. It has three steps:
- Pick an element, called a pivot, from the data to be sorted
- Reorder the data so that all of the elements that are less than the pivot come before the pivot in the data and all the elements that are greater than the pivot come after the pivot.
- Recursively process the subset of lesser elements and the subset of greater elements.
internalstaticintpartitionBlock(T[]data,intstartIndex,intendIndex,IComparer<T>comparer){// get the pivot value - we will be comparing all of the other items against this valueTpivot=data[startIndex];// put the pivot value at the end of blockswapValues(data,startIndex,endIndex);// index used to store values smaller than the pivotintstoreIndex=startIndex;// iterate through the items in the blockfor(inti=startIndex;i<endIndex;i++){// look for items that are smaller or equal to the pivotif(comparer.Compare(data[i],pivot)<=0){// move the value and increment the indexswapValues(data,i,storeIndex);storeIndex++;}}swapValues(data,storeIndex,endIndex);returnstoreIndex;}Пример кода из книги показался мне недостаточно эффективным в случае наличия большого числа повторяющихся элементов. Не выполнялся 2-й шаг алгоритма, среди "меньших" элементов могли встречаться равные pivot, а вычленять из цепи только по одному элементу - неэффективно, в итоге, работает более чем 2 раза медленнее. Переписав всё по-своему, получил код, который в синхронном варианте сортирует 5 000 000 псевдослучайных чисел (со значениями от 0 до 255) за ~206 ms (при использовании библиотечного Array.Sort результат ~570 ms). При параллельном выполнении ~85 ms, ускорение данного алгоритма всего в 2.5 раза при использовании 8 ядер.
privateintDoPartByPivot(ShortStatework){inteqlCount=0,fstIndex=work.fstIndex,lstIndex=work.lstIndex;shortvalue,pivot=data[fstIndex++];// iterate through comparing with pivotwhile(fstIndex<=lstIndex){value=data[fstIndex];// look for smaller or equal to the pivotif(value<=pivot){// increment the indexfstIndex++;}else{// move greater to the right sideSwapValues(fstIndex,lstIndex);lstIndex--;}}fstIndex=work.fstIndex;// iterate through comparing with pivotwhile(fstIndex<=lstIndex){value=data[fstIndex];// look for equal to the pivotif(value==pivot){// move equal to the right sideSwapValues(fstIndex,lstIndex);lstIndex--;eqlCount++;}else{// increment the indexfstIndex++;}}//count of repeated valueswork.eqlCount=eqlCount;//values greater or equal to the pivotreturnfstIndex;}Большое спасибо dot.net за ConcurrentQueue, который позволяет обходиться нам без мьютексов. Скорее всего, он написан с использованием неблокирующей технологии <atomic>.
Аналогичный код на C++ (см. CppSort.cpp) в синхронном исполнении сортирует 5 000 000 псевдослучайных чисел (со значениями от 0 до 255) за ~205 ms, против ~190 ms при использовании std::sort(vec.begin(), vec.end());. Параллельная сортировка на C++ в моем исполнении только сравнялась с C# и составила ~85 ms. Сколько я ни бился, быстрее не получилось. Сортировка производится на 8-ми потоках, когда нет блоков для обработки - поток засыпает. Работающий поток посылает уведомления о новых блоках. Подсчет "уснувших" потоков дает возможность разбудить их все и завершить программу.
voidWorker()
{
//std::cout << std::this_thread::get_id() << " started " << std::endl;
std::unique_lock<std::mutex> uni(mux);
while (true)
{
//entered workersif (!works.empty())
{
ShortState *work = works.front();
works.pop();
uni.unlock();
//work obtainedDoSort(work);
delete work;
uni.lock();
}
else
{
++counter;
if (counter > 7)
{
exit = true;
worker_cnd.notify_all();
main_cnd.notify_one();
}
else
{
worker_cnd.wait(uni, [this] { returnUnLock(); });
--counter;
}
if (exit) break;
}
}
//std::cout << std::this_thread::get_id() << " done " << std::endl;
}В С++17 STL я не нашел готовых к применению неблокирующих контейнеров. Работа с библиотекой <atomic> для меня остаётся чёрной магией. Лучшая книга по данной теме - в книге издательства Manning "C++ Concurrency in Action" by Antony Williams.
По прошествии мучительных исканий, я всё же немного разобрался с применением неблокирующей библиотеки <atomic>. Пример кода в CppSortAtom.cpp, фрагмент неблокирующего стека LIFO приведён ниже. Особенностью решения является использование относительно небольшого буфера для хранения подготовленных к обработке блоков данных, при использовании LIFO можно обойтись порядком Log2(5 000 000), я взял с запасом - 101.
При формировании нового задания вызывается метод SafeDeque::push, в котором из головы стека свободных блоков берём первый. Указатель aheap безопасно, с точки зрения гонки потоков, сдвигается на одну позицию вперёд, оставляя выбранный позади. Далее, заполняем выбранный элемент данными для обработки и помещаем его в голову стека подготовленных заданий. Атомарный указатель awork безопасно сдвигается на одну позицию к началу, указывая на новый элемент.
Рабочий поток запрашивает новый блок данных, вызывая метод SafeDeque::pop, при этом выбирается блок с головы рабочего стека. Атомарный указатель awork безопасно сдвигается на одну позицию вперёд, оставляя выбранный позади. После обработки блока, он возвращается в голову стека свободных блоков посредством вызова метода SafeDeque::free. Указатель aheap безопасно сдвигается на одну позицию к началу и указывает на возвращенный элемент.
Скорость сортировки не улучшилась, в лучшем случае 88 ms. Зато было интересно и память программа потребляет экономно.
structWorkNode
{
WorkNode *next;
int fstIndex, lstIndex, eqlCount;
WorkNode() {
next = NULL; fstIndex = lstIndex = eqlCount = 0;
}
};
classSafeDeque
{
private:
WorkNode* root;
std::atomic<WorkNode*> aheap{ NULL };
std::atomic<WorkNode*> awork{ NULL };
public:SafeDeque() {
//allocate memory on heap
root = new WorkNode[101];
aheap.store(&root[0]);
//chain of heap nodesfor (int i{ 0 }; i < 100; ++i) {
root[i].next = &root[i + 1];
}
}
~SafeDeque() {
delete root;
}
voidpush(int fst, int lst) {
//obtain node from heap
WorkNode* heap = aheap.load();
//exlude from heapwhile (!aheap.compare_exchange_weak(heap, heap->next));
//obtain last work node
WorkNode* work = awork.load();
//add node to lifo stackwhile (!awork.compare_exchange_weak(work, heap));
heap->next = work;
heap->fstIndex = fst;
heap->lstIndex = lst;
}
WorkNode* pop() {
//obtain last (head) node from stack
WorkNode* work = awork.load();
if (!work) returnNULL;
//trim last node from chainwhile (!awork.compare_exchange_weak(work, work->next) && work);
return work;
}
voidfree(WorkNode* work) {
//obtain last heap node
WorkNode* heap = aheap.load();
//return work node to heapwhile (!aheap.compare_exchange_weak(heap, work));
//add on head of chain
work->next = heap;
}
};Для синхронизации потоков Go полагается на горутины и каналы. Настораживает тот факт, что стандартный шаблон sort.Sort(Int16Slice(data[:])) отрабатывает аналогичный массив за 620 ms, однако, параллельный вариант программы сортирует всего за 78 ms! См. qsortpar.go Также как в C# не пришлось извращаться, чтобы подсчитать, когда мои горутины все завершились. Чтобы канал works не заблокировал все горутины, когда закончатся неотсортированные блоки, применил неблокирующий код, позоляющий подсчитывать количество действующих горутин. Очень интересный инструмент на основе каналов и горутин, логика работы с потоками стала нагляднее, а количество строк кода уменьшилось почти вдвое!
func (dataInt16Slice) workLoop(workschan*Work, counter*sync.WaitGroup) {
counter.Add(1)
hasWork:=trueforhasWork {
select {
casework:=<-works:
//add goroutineifdata.doSort(works, work) {
godata.workLoop(works, counter)
}
default:
hasWork=false
}
}
//fmt.Printf("Goroutine done\n")counter.Done()
}