1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
|
#include "parallel.h"
#include <stdlib.h>
#if defined(__unix__) || defined(__APPLE__)
#include <pthread.h>
#include <unistd.h>
size_t psi_thread_count(void)
{
long n = sysconf(_SC_NPROCESSORS_ONLN);
return n > 1 ? (size_t)n : 1;
}
struct PsiParallelTask
{
void (*body)(size_t, size_t, void*);
void* ctx;
size_t start;
size_t end;
};
static void* worker(void* arg)
{
struct PsiParallelTask* task = arg;
task->body(task->start, task->end, task->ctx);
return NULL;
}
void psi_parallel_for(size_t count, void (*body)(size_t, size_t, void*), void* ctx)
{
size_t threads = psi_thread_count();
if (threads <= 1 || count <= 1)
{
body(0, count, ctx);
return;
}
if (threads > count)
threads = count;
pthread_t* handles = malloc(threads * sizeof *handles);
struct PsiParallelTask* tasks = malloc(threads * sizeof *tasks);
if (handles == NULL || tasks == NULL)
{
free(handles);
free(tasks);
body(0, count, ctx);
return;
}
size_t chunk = count / threads;
size_t remainder = count % threads;
size_t start = 0;
for (size_t t = 0; t < threads; t++)
{
size_t len = chunk + (t < remainder ? 1 : 0);
tasks[t] = (struct PsiParallelTask){ body, ctx, start, start + len };
start += len;
pthread_create(&handles[t], NULL, worker, &tasks[t]);
}
for (size_t t = 0; t < threads; t++)
pthread_join(handles[t], NULL);
free(handles);
free(tasks);
}
#else
size_t psi_thread_count(void)
{
return 1;
}
void psi_parallel_for(size_t count, void (*body)(size_t, size_t, void*), void* ctx)
{
body(0, count, ctx);
}
#endif
|