forked from mbien/jocl
/
CLForkJoinPool.java
209 lines (178 loc) · 7.25 KB
/
CLForkJoinPool.java
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
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
/*
* Copyright (c) 2011, Michael Bien
* All rights reserved.
*
* Redistribution and use in source and binary forms, with or without modification, are
* permitted provided that the following conditions are met:
*
* 1. Redistributions of source code must retain the above copyright notice, this list of
* conditions and the following disclaimer.
*
* 2. Redistributions in binary form must reproduce the above copyright notice, this list
* of conditions and the following disclaimer in the documentation and/or other materials
* provided with the distribution.
*
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND ANY EXPRESS OR IMPLIED
* WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND
* FITNESS FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
* CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR
* CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR
* SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON
* ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
* NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF
* ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
*
*/
/*
* Created on Tuesday, August 02 2011 01:53
*/
package com.jogamp.opencl.util.concurrent;
import com.jogamp.opencl.CLCommandQueue;
import com.jogamp.opencl.CLDevice;
import com.jogamp.opencl.util.CLMultiContext;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.ForkJoinWorkerThread;
import java.util.concurrent.Future;
/**
* A multithreaded, fixed size pool of OpenCL command queues supporting fork-join tasks.
* <p>
* The usage is similar to {@link ForkJoinPool} but uses {@link CLRecursiveTask}s.
* </p>
* @see CLRecursiveTask
* @author Michael Bien
*/
public class CLForkJoinPool extends CLAbstractExecutorService {
private CLForkJoinPool(ExecutorService executor, List<CLCommandQueue> queues) {
super(executor, queues);
}
public static CLForkJoinPool create(CLMultiContext mc, CLCommandQueue.Mode... modes) {
return create(mc.getDevices(), modes);
}
public static CLForkJoinPool create(Collection<? extends CLDevice> devices, CLCommandQueue.Mode... modes) {
List<CLCommandQueue> queues = new ArrayList<CLCommandQueue>(devices.size());
for (CLDevice device : devices) {
queues.add(device.createCommandQueue(modes));
}
return create(queues);
}
public static CLForkJoinPool create(Collection<CLCommandQueue> queues) {
List<CLCommandQueue> list = new ArrayList<CLCommandQueue>(queues);
CLThreadFactory factory = new CLThreadFactory(list);
int size = list.size();
ExecutorService service = new ForkJoinPool(size, factory, null, false);
return new CLForkJoinPool(service, list);
}
/**
* Performs the given task, returning its result upon completion.
* @see ForkJoinPool#invoke(java.util.concurrent.ForkJoinTask)
*/
public <R> R invoke(CLRecursiveTask<? extends CLQueueContext, R> task) {
// shortcut, prevents redundant wrapping
return getExcecutor().invoke(task);
}
/**
* Submits this task to the pool for execution returning its {@link Future}.
* @see ForkJoinPool#submit(java.util.concurrent.ForkJoinTask)
*/
public <R> Future<R> submit(CLRecursiveTask<? extends CLQueueContext, R> task) {
// shortcut, prevents redundant wrapping
return getExcecutor().submit(task);
}
@Override
ForkJoinPool getExcecutor() {
return (ForkJoinPool) super.getExcecutor();
}
/**
* Returns an estimate of the total number of tasks stolen from
* one thread's work queue by another.
* @see ForkJoinPool#getStealCount()
*/
public long getStealCount() {
return getExcecutor().getStealCount();
}
/**
* Returns an estimate of the number of tasks submitted to this
* pool that have not yet begun executing. This method may take
* time proportional to the number of submissions.
* @see ForkJoinPool#getQueuedSubmissionCount()
*/
public int getQueuedSubmissionCount() {
return getExcecutor().getQueuedSubmissionCount();
}
/**
* Returns an estimate of the total number of tasks currently held
* in queues by worker threads (but not including tasks submitted
* to the pool that have not begun executing). This value is only
* an approximation, obtained by iterating across all threads in
* the pool. This method may be useful for tuning task
* granularities.
* @see ForkJoinPool#getQueuedTaskCount()
*/
public long getQueuedTaskCount() {
return getExcecutor().getQueuedTaskCount();
}
/**
* Returns {@code true} if there are any tasks submitted to this
* pool that have not yet begun executing.
* @see ForkJoinPool#hasQueuedSubmissions()
*/
public boolean hasQueuedSubmissions() {
return getExcecutor().hasQueuedSubmissions();
}
/**
* Returns {@code true} if all worker threads are currently idle.
* An idle worker is one that cannot obtain a task to execute
* because none are available to steal from other threads, and
* there are no pending submissions to the pool. This method is
* conservative; it might not return {@code true} immediately upon
* idleness of all threads, but will eventually become true if
* threads remain inactive.
* @see ForkJoinPool#isQuiescent()
*/
public boolean isQuiescent() {
return getExcecutor().isQuiescent();
}
private static class CLThreadFactory implements ForkJoinPool.ForkJoinWorkerThreadFactory {
private int index = 0;
private final List<CLCommandQueue> queues;
private CLThreadFactory(List<CLCommandQueue> queues) {
this.queues = queues;
}
@Override
public synchronized ForkJoinWorkerThread newThread(ForkJoinPool pool) {
return new ForkJoinQueueWorkerThread(pool, queues.get(index++));
}
}
final static class ForkJoinQueueWorkerThread extends ForkJoinWorkerThread implements CommandQueueThread {
private final CLCommandQueue queue;
private final Map<Object, CLQueueContext> contextMap;
public ForkJoinQueueWorkerThread(ForkJoinPool pool, CLCommandQueue queue) {
super(pool);
this.queue = queue;
this.contextMap = new HashMap<Object, CLQueueContext>();
}
@Override
public void run() {
super.run();
//release threadlocal contexts
queue.finish();
for (CLQueueContext context : contextMap.values()) {
context.release();
}
}
@Override
public Map<Object, CLQueueContext> getContextMap() {
return contextMap;
}
@Override
public CLCommandQueue getQueue() {
return queue;
}
}
}