ConcurrentQueue: Update to HEAD

This commit is contained in:
Kyle Schwarz
2022-09-09 15:46:46 -04:00
parent b35fab754c
commit 2e296dc8d3
208 changed files with 42592 additions and 274 deletions
@@ -28,6 +28,7 @@
#include "boostqueue.h"
#include "tbbqueue.h"
#include "stdqueue.h"
#include "dlibqueue.h"
#include "../tests/common/simplethread.h"
#include "../tests/common/systemtime.h"
#include "cpuid.h"
@@ -192,7 +193,7 @@ enum queue_id_t
queue_simplelockfree,
queue_lockbased,
queue_std,
queue_dlib,
QUEUE_COUNT
};
@@ -204,6 +205,7 @@ const char QUEUE_NAMES[QUEUE_COUNT][64] = {
"SimpleLockFreeQueue",
"LockBasedQueue",
"std::queue",
"dlib::pipe"
};
const char QUEUE_SUMMARY_NOTES[QUEUE_COUNT][128] = {
@@ -214,6 +216,7 @@ const char QUEUE_SUMMARY_NOTES[QUEUE_COUNT][128] = {
"",
"",
"single thread only",
""
};
const bool QUEUE_TOKEN_SUPPORT[QUEUE_COUNT] = {
@@ -224,6 +227,7 @@ const bool QUEUE_TOKEN_SUPPORT[QUEUE_COUNT] = {
false,
false,
false,
false
};
const int QUEUE_MAX_THREADS[QUEUE_COUNT] = {
@@ -234,6 +238,7 @@ const int QUEUE_MAX_THREADS[QUEUE_COUNT] = {
-1,
-1,
1,
-1
};
const bool QUEUE_BENCH_SUPPORT[QUEUE_COUNT][BENCHMARK_TYPE_COUNT] = {
@@ -244,6 +249,7 @@ const bool QUEUE_BENCH_SUPPORT[QUEUE_COUNT][BENCHMARK_TYPE_COUNT] = {
{ 1, 1, 1, 0, 0, 1, 0, 1, 0, 1, 0, 1, 1, 1, 1, 1, 1 },
{ 1, 1, 1, 0, 0, 1, 0, 1, 0, 1, 0, 1, 1, 1, 1, 1, 1 },
{ 0, 1, 1, 0, 0, 1, 0, 0, 0, 0, 0, 0, 0, 0, 1, 1, 0 },
{ 1, 1, 1, 0, 0, 1, 0, 1, 0, 1, 0, 1, 1, 1, 1, 1, 1 }
};
@@ -254,6 +260,9 @@ struct Traits : public moodycamel::ConcurrentQueueDefaultTraits
// block size will improve throughput (which is mostly what
// we're after with these benchmarks).
static const size_t BLOCK_SIZE = 64;
// Reuse blocks once allocated.
static const bool RECYCLE_ALLOCATED_BLOCKS = true;
};
@@ -295,7 +304,7 @@ counter_t adjustForThreads(counter_t suggestedOps, int nthreads)
}
template<typename TQueue>
template<typename TQueue, typename item_t>
counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, unsigned int randSeed)
{
switch (benchmark) {
@@ -306,9 +315,10 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
std::uniform_int_distribution<int> rand(0, 20);
double total = 0;
SystemTime start;
item_t item = 1;
for (counter_t i = 0; i != ops; ++i) {
start = getSystemTime();
q.enqueue(i);
q.enqueue(item);
total += getTimeDelta(start);
}
return total;
@@ -319,9 +329,10 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
case bench_mostly_enqueue: {
return adjustForThreads(rampUpToMeasurableNumberOfMaxOps([](counter_t ops) {
TQueue q;
item_t item = 1;
auto start = getSystemTime();
for (counter_t i = 0; i != ops; ++i) {
q.enqueue(i);
q.enqueue(item);
}
return getTimeDelta(start);
}), nthreads);
@@ -333,13 +344,14 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
case bench_mpsc: {
return adjustForThreads(rampUpToMeasurableNumberOfMaxOps([](counter_t ops) {
TQueue q;
item_t item = 1;
for (counter_t i = 0; i != ops; ++i) {
q.enqueue(i);
q.enqueue(item);
}
int item;
item_t item_rec;
auto start = getSystemTime();
for (counter_t i = 0; i != ops; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
return getTimeDelta(start);
}), nthreads);
@@ -347,7 +359,7 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
case bench_only_enqueue_bulk:
case bench_only_enqueue_bulk_prealloc:
case bench_mostly_enqueue_bulk: {
std::vector<counter_t> data;
std::vector<item_t> data;
for (counter_t i = 0; i != BULK_BATCH_SIZE; ++i) {
data.push_back(i);
}
@@ -364,7 +376,7 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
case bench_mostly_dequeue_bulk: {
return adjustForThreads(rampUpToMeasurableNumberOfMaxOps([](counter_t ops) {
TQueue q;
std::vector<int> data(BULK_BATCH_SIZE);
std::vector<item_t> data(BULK_BATCH_SIZE);
for (counter_t i = 0; i != ops; ++i) {
q.enqueue_bulk(data.cbegin(), data.size());
}
@@ -379,10 +391,10 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
case bench_empty_dequeue: {
return adjustForThreads(rampUpToMeasurableNumberOfMaxOps([](counter_t ops) {
TQueue q;
int item;
item_t item_rec;
auto start = getSystemTime();
for (counter_t i = 0; i != ops; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
return getTimeDelta(start);
}), nthreads);
@@ -390,11 +402,12 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
case bench_enqueue_dequeue_pairs: {
return adjustForThreads(rampUpToMeasurableNumberOfMaxOps([](counter_t ops) {
TQueue q;
int item;
item_t item = 1;
item_t item_rec;
auto start = getSystemTime();
for (counter_t i = 0; i != ops; ++i) {
q.enqueue(i);
q.try_dequeue(item);
q.enqueue(item);
q.try_dequeue(item_rec);
}
return getTimeDelta(start);
}), nthreads);
@@ -403,11 +416,12 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
case bench_heavy_concurrent: {
return adjustForThreads(rampUpToMeasurableNumberOfMaxOps([](counter_t ops) {
TQueue q;
int item;
item_t item=1;
item_t item_rec;
auto start = getSystemTime();
for (counter_t i = 0; i != ops; ++i) {
q.enqueue(i);
q.try_dequeue(item);
q.enqueue(item);
q.try_dequeue(item_rec);
}
return getTimeDelta(start);
}), nthreads);
@@ -421,7 +435,7 @@ counter_t determineMaxOpsForBenchmark(benchmark_type_t benchmark, int nthreads,
// Returns time elapsed, in (fractional) milliseconds
template<typename TQueue>
template<typename TQueue, typename item_t>
double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, unsigned int randSeed, counter_t maxOps, int maxThreads, counter_t& out_opCount)
{
double result = 0;
@@ -435,13 +449,14 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
std::vector<counter_t> ops(nthreads);
std::vector<double> times(nthreads);
std::atomic<int> ready(0);
item_t item_rec;
item_t item = 1;
for (int tid = 0; tid != nthreads; ++tid) {
threads[tid] = SimpleThread([&](int id) {
ready.fetch_add(1, std::memory_order_relaxed);
while (ready.load(std::memory_order_relaxed) != nthreads)
continue;
int item;
SystemTime start;
RNG_t rng(randSeed * (id + 1));
std::uniform_int_distribution<int> rand(0, 20);
@@ -449,24 +464,25 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
times[id] = 0;
typename TQueue::consumer_token_t consTok(q);
typename TQueue::producer_token_t prodTok(q);
for (counter_t i = 0; i != maxOps; ++i) {
if (rand(rng) == 0) {
start = getSystemTime();
if ((i & 1) == 0) {
if (useTokens) {
q.try_dequeue(consTok, item);
q.try_dequeue(consTok, item_rec);
}
else {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
}
else {
if (useTokens) {
q.enqueue(prodTok, i);
q.enqueue(prodTok, item);
}
else {
q.enqueue(i);
q.enqueue(item);
}
}
times[id] += getTimeDelta(start);
@@ -482,8 +498,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
out_opCount += ops[tid];
result += times[tid];
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
@@ -491,23 +506,26 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
out_opCount = maxOps * nthreads;
TQueue q;
item_t item = 1;
item_t item_rec;
{
// Enqueue opcount elements first, then dequeue them; this
// will "stretch out" the queue, letting implementatations
// that re-use memory internally avoid having to allocate
// more later during the timed enqueue operations.
std::vector<SimpleThread> threads(nthreads);
for (int tid = 0; tid != nthreads; ++tid) {
threads[tid] = SimpleThread([&](int id) {
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
}, tid);
@@ -517,8 +535,8 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
}
// Now empty the queue
int item;
while (q.try_dequeue(item))
while (q.try_dequeue(item_rec))
continue;
}
@@ -528,12 +546,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
result = getTimeDelta(start);
@@ -552,12 +570,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
timings[id] = getTimeDelta(start);
@@ -569,8 +587,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
result += timings[tid];
}
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
@@ -578,18 +595,20 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
out_opCount = maxOps * nthreads;
TQueue q;
item_t item = 1;
item_t item_rec;
if (nthreads == 1) {
// No contention -- measures raw single-item enqueue speed
auto start = getSystemTime();
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
result = getTimeDelta(start);
@@ -608,12 +627,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
timings[id] = getTimeDelta(start);
@@ -625,8 +644,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
result += timings[tid];
}
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
@@ -635,6 +653,8 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
out_opCount = maxOps * nthreads;
TQueue q;
item_t item = 1;
item_t item_rec;
{
// Fill up the queue first
std::vector<SimpleThread> threads(benchmark == bench_spmc_preproduced ? 1 : nthreads);
@@ -644,12 +664,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != itemsPerThread; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != itemsPerThread; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
}, tid);
@@ -661,17 +681,16 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (nthreads == 1) {
// No contention -- measures raw single-item dequeue speed
int item;
auto start = getSystemTime();
if (useTokens) {
typename TQueue::consumer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(tok, item);
q.try_dequeue(tok, item_rec);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
}
result = getTimeDelta(start);
@@ -686,17 +705,16 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
while (ready.load(std::memory_order_relaxed) != nthreads)
continue;
int item;
auto start = getSystemTime();
if (useTokens) {
typename TQueue::consumer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(tok, item);
q.try_dequeue(tok, item_rec);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
}
timings[id] = getTimeDelta(start);
@@ -708,14 +726,15 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
result += timings[tid];
}
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_mostly_enqueue: {
// Measures the average operation speed when most threads are enqueueing
TQueue q;
item_t item = 1;
item_t item_rec;
out_opCount = maxOps * nthreads;
std::vector<SimpleThread> threads(nthreads);
std::vector<double> timings(nthreads);
@@ -731,12 +750,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
timings[id] = getTimeDelta(start);
@@ -748,17 +767,16 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
while (ready.load(std::memory_order_relaxed) != nthreads)
continue;
int item;
auto start = getSystemTime();
if (useTokens) {
typename TQueue::consumer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(tok, item);
q.try_dequeue(tok, item_rec);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
}
timings[id] = getTimeDelta(start);
@@ -769,14 +787,15 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
threads[tid].join();
result += timings[tid];
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_mostly_dequeue: {
// Measures the average operation speed when most threads are dequeueing
TQueue q;
item_t item = 1;
item_t item_rec;
out_opCount = maxOps * nthreads;
std::vector<SimpleThread> threads(nthreads);
std::vector<double> timings(nthreads);
@@ -789,12 +808,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
}, tid);
@@ -810,17 +829,16 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
while (ready.load(std::memory_order_relaxed) != nthreads)
continue;
int item;
auto start = getSystemTime();
if (useTokens) {
typename TQueue::consumer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(tok, item);
q.try_dequeue(tok, item_rec);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
}
timings[id] = getTimeDelta(start);
@@ -836,12 +854,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
timings[id] = getTimeDelta(start);
@@ -852,13 +870,14 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
threads[tid].join();
result += timings[tid];
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_only_enqueue_bulk_prealloc: {
TQueue q;
item_t item = 1;
item_t item_rec;
{
// Enqueue opcount elements first, then dequeue them; this
// will "stretch out" the queue, letting implementatations
@@ -870,12 +889,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
}, tid);
@@ -885,8 +904,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
}
// Now empty the queue
int item;
while (q.try_dequeue(item))
while (q.try_dequeue(item_rec))
continue;
}
@@ -942,13 +960,14 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
result += timings[tid];
}
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_only_enqueue_bulk: {
TQueue q;
item_t item = 1;
item_t item_rec;
std::vector<counter_t> data;
for (counter_t i = 0; i != BULK_BATCH_SIZE; ++i) {
data.push_back(i);
@@ -1001,14 +1020,15 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
result += timings[tid];
}
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_mostly_enqueue_bulk: {
// Measures the average speed of enqueueing in bulk under light contention
TQueue q;
item_t item = 1;
item_t item_rec;
std::vector<counter_t> data;
for (counter_t i = 0; i != BULK_BATCH_SIZE; ++i) {
data.push_back(i);
@@ -1076,14 +1096,15 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
out_opCount += ops[tid];
}
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_only_dequeue_bulk: {
// Measures the average speed of dequeueing in bulk when all threads are consumers
TQueue q;
item_t item = 1;
item_t item_rec;
{
// Fill up the queue first
std::vector<int> data(BULK_BATCH_SIZE);
@@ -1166,14 +1187,15 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
out_opCount += ops[tid];
}
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_mostly_dequeue_bulk: {
// Measures the average speed of dequeueing in bulk under light contention
TQueue q;
item_t item = 1;
item_t item_rec;
auto enqueueThreads = std::max(1, nthreads / 4);
out_opCount = maxOps * BULK_BATCH_SIZE * enqueueThreads;
std::vector<SimpleThread> threads(nthreads);
@@ -1260,8 +1282,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
out_opCount += ops[tid];
}
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
@@ -1269,6 +1290,8 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
counter_t elementsToDequeue = maxOps * (nthreads - 1);
TQueue q;
item_t item = 1;
item_t item_rec;
std::vector<SimpleThread> threads(nthreads - 1);
std::vector<double> timings(nthreads - 1);
std::vector<counter_t> ops(nthreads - 1);
@@ -1297,7 +1320,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
}
else {
while (true) {
if (q.try_dequeue(item)) {
if (q.try_dequeue(item_rec)) {
totalDequeued.fetch_add(1, std::memory_order_relaxed);
}
else if (totalDequeued.load(std::memory_order_relaxed) == elementsToDequeue) {
@@ -1313,7 +1336,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
lynchpin.store(true, std::memory_order_seq_cst);
for (counter_t i = 0; i != elementsToDequeue; ++i) {
q.enqueue(i);
q.enqueue(item);
}
result = 0;
@@ -1323,13 +1346,14 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
result += timings[tid];
out_opCount += ops[tid];
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_mpsc: {
TQueue q;
item_t item = 1;
item_t item_rec;
counter_t elementsToDequeue = maxOps * (nthreads - 1);
std::vector<SimpleThread> threads(nthreads);
std::atomic<int> ready(0);
@@ -1353,7 +1377,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
}
else {
for (counter_t i = 0; i != elementsToDequeue;) {
i += q.try_dequeue(item) ? 1 : 0;
i += q.try_dequeue(item_rec) ? 1 : 0;
++out_opCount;
}
}
@@ -1369,12 +1393,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
}, tid);
@@ -1384,14 +1408,15 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
for (int tid = 0; tid != nthreads; ++tid) {
threads[tid].join();
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
case bench_empty_dequeue: {
// Measures the average speed of attempting to dequeue from an empty queue
TQueue q;
item_t item = 1;
item_t item_rec;
// Fill up then empty the queue first
{
std::vector<SimpleThread> threads(maxThreads > 0 ? maxThreads : 8);
@@ -1400,12 +1425,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t tok(q);
for (counter_t i = 0; i != 10000; ++i) {
q.enqueue(tok, i);
q.enqueue(tok, item);
}
}
else {
for (counter_t i = 0; i != 10000; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
}, tid);
@@ -1415,8 +1440,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
}
// Empty the queue
int item;
while (q.try_dequeue(item))
while (q.try_dequeue(item_rec))
continue;
}
@@ -1433,11 +1457,11 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
}
result = getTimeDelta(start);
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
}
else {
out_opCount = maxOps * nthreads;
@@ -1460,7 +1484,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
}
timings[id] = getTimeDelta(start);
@@ -1471,8 +1495,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
threads[tid].join();
result += timings[tid];
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
}
break;
}
@@ -1482,26 +1505,28 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
// (that eight separate threads had at one point enqueued to)
out_opCount = maxOps * 2 * nthreads;
TQueue q;
item_t item = 1;
item_t item_rec;
if (nthreads == 1) {
// No contention -- measures speed of immediately dequeueing the item that was just enqueued
int item;
auto start = getSystemTime();
if (useTokens) {
typename TQueue::producer_token_t ptok(q);
typename TQueue::consumer_token_t ctok(q);
typename TQueue::producer_token_t prodTok(q);
typename TQueue::consumer_token_t consTok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(ptok, i);
q.try_dequeue(ctok, item);
q.enqueue(prodTok, item);
q.try_dequeue(consTok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.try_dequeue(item);
q.enqueue(item);
q.try_dequeue(item_rec);
}
}
result = getTimeDelta(start);
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
}
else {
std::vector<SimpleThread> threads(nthreads);
@@ -1516,17 +1541,17 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
int item;
auto start = getSystemTime();
if (useTokens) {
typename TQueue::producer_token_t ptok(q);
typename TQueue::consumer_token_t ctok(q);
typename TQueue::producer_token_t prodTok(q);
typename TQueue::consumer_token_t consTok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(ptok, i);
q.try_dequeue(ctok, item);
q.enqueue(prodTok, item);
q.try_dequeue(consTok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.try_dequeue(item);
q.enqueue(item);
q.try_dequeue(item_rec);
}
}
timings[id] = getTimeDelta(start);
@@ -1537,8 +1562,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
threads[tid].join();
result += timings[tid];
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
}
break;
}
@@ -1547,6 +1571,8 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
// Measures the average operation speed with many threads under heavy load
out_opCount = maxOps * nthreads;
TQueue q;
item_t item = 1;
item_t item_rec;
std::vector<SimpleThread> threads(nthreads);
std::vector<double> timings(nthreads);
std::atomic<int> ready(0);
@@ -1566,13 +1592,13 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
for (counter_t i = 0; i != maxOps / 2; ++i) {
q.try_dequeue(consTok, item);
q.enqueue(prodTok, i);
q.enqueue(prodTok, item);
}
}
else {
for (counter_t i = 0; i != maxOps / 2; ++i) {
q.try_dequeue(item);
q.enqueue(i);
q.try_dequeue(item_rec);
q.enqueue(item);
}
}
}
@@ -1582,12 +1608,12 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
if (useTokens) {
typename TQueue::producer_token_t prodTok(q);
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(prodTok, i);
q.enqueue(prodTok, item);
}
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.enqueue(i);
q.enqueue(item);
}
}
}
@@ -1602,7 +1628,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
}
else {
for (counter_t i = 0; i != maxOps; ++i) {
q.try_dequeue(item);
q.try_dequeue(item_rec);
}
}
}
@@ -1615,8 +1641,7 @@ double runBenchmark(benchmark_type_t benchmark, int nthreads, bool useTokens, un
threads[tid].join();
result += timings[tid];
}
int item;
forceNoOptimizeDummy = q.try_dequeue(item) ? 1 : 0;
forceNoOptimizeDummy = q.try_dequeue(item_rec) ? 1 : 0;
break;
}
@@ -1964,25 +1989,28 @@ int main(int argc, char** argv)
counter_t maxOps;
switch ((queue_id_t)queue) {
case queue_moodycamel_ConcurrentQueue:
maxOps = determineMaxOpsForBenchmark<moodycamel::ConcurrentQueue<int, Traits>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
maxOps = determineMaxOpsForBenchmark<moodycamel::ConcurrentQueue<int, Traits>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
break;
case queue_moodycamel_BlockingConcurrentQueue:
maxOps = determineMaxOpsForBenchmark<moodycamel::BlockingConcurrentQueue<int, Traits>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
maxOps = determineMaxOpsForBenchmark<moodycamel::BlockingConcurrentQueue<int, Traits>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
break;
case queue_lockbased:
maxOps = determineMaxOpsForBenchmark<LockBasedQueue<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
maxOps = determineMaxOpsForBenchmark<LockBasedQueue<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
break;
case queue_simplelockfree:
maxOps = determineMaxOpsForBenchmark<SimpleLockFreeQueue<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
maxOps = determineMaxOpsForBenchmark<SimpleLockFreeQueue<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
break;
case queue_boost:
maxOps = determineMaxOpsForBenchmark<BoostQueueWrapper<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
maxOps = determineMaxOpsForBenchmark<BoostQueueWrapper<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
break;
case queue_tbb:
maxOps = determineMaxOpsForBenchmark<TbbQueueWrapper<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
maxOps = determineMaxOpsForBenchmark<TbbQueueWrapper<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
break;
case queue_std:
maxOps = determineMaxOpsForBenchmark<StdQueueWrapper<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
maxOps = determineMaxOpsForBenchmark<StdQueueWrapper<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
break;
case queue_dlib:
maxOps = determineMaxOpsForBenchmark<DlibQueueWrapper<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed);
break;
default:
assert(false && "There should be a case here for every queue in the benchmarks!");
@@ -1997,25 +2025,28 @@ int main(int argc, char** argv)
switch ((queue_id_t)queue) {
case queue_moodycamel_ConcurrentQueue:
elapsed = runBenchmark<moodycamel::ConcurrentQueue<int, Traits>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
elapsed = runBenchmark<moodycamel::ConcurrentQueue<int, Traits>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
break;
case queue_moodycamel_BlockingConcurrentQueue:
elapsed = runBenchmark<moodycamel::BlockingConcurrentQueue<int, Traits>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
elapsed = runBenchmark<moodycamel::BlockingConcurrentQueue<int, Traits>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
break;
case queue_lockbased:
elapsed = runBenchmark<LockBasedQueue<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
elapsed = runBenchmark<LockBasedQueue<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
break;
case queue_simplelockfree:
elapsed = runBenchmark<SimpleLockFreeQueue<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
elapsed = runBenchmark<SimpleLockFreeQueue<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
break;
case queue_boost:
elapsed = runBenchmark<BoostQueueWrapper<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
elapsed = runBenchmark<BoostQueueWrapper<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
break;
case queue_tbb:
elapsed = runBenchmark<TbbQueueWrapper<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
elapsed = runBenchmark<TbbQueueWrapper<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
break;
case queue_std:
elapsed = runBenchmark<StdQueueWrapper<int>>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
elapsed = runBenchmark<StdQueueWrapper<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
break;
case queue_dlib:
elapsed = runBenchmark<DlibQueueWrapper<int>, int>((benchmark_type_t)benchmark, nthreads, (bool)useTokens, seed, maxOps, maxThreads, ops);
break;
default:
assert(false && "There should be a case here for every queue in the benchmarks!");