-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathProcessorAction.java
More file actions
79 lines (64 loc) · 2.34 KB
/
Copy pathProcessorAction.java
File metadata and controls
79 lines (64 loc) · 2.34 KB
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
package processor;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.RecursiveTask;
/**
* 多线程任务调度处理
* Created by focus on 2018/3/16.
*/
public class ProcessorAction extends RecursiveTask{
private ProcessorQueue queue;
public ProcessorAction(ProcessorQueue queue){
this.queue = queue;
}
@Override
protected Object compute() {
IProcessor processor;
Object result = null;
if(queue.isProcessEnable()){
processor = queue.nextProcessor();
return processor.process();
}else{
List<IProcessor> asyncProcessors = new CopyOnWriteArrayList<>();
while((processor = queue.nextProcessor()) != null){
if(processor.isAsyn()){
asyncProcessors.add(processor);
}else{
if(!asyncProcessors.isEmpty()){
invokeAllAsyncProcessors(asyncProcessors);
}
// 执行同步程序
ProcessorAction action = new ProcessorAction(queue.buildNewQueue(processor));
result = action.invoke();
queue.processorService.addResult(processor.id(),result);
asyncProcessors.clear();
}
}
if(!asyncProcessors.isEmpty()){
invokeAllAsyncProcessors(asyncProcessors);
}
}
return result;
}
/**
* 执行所有异步程序
* @param asyncProcessors
*/
public void invokeAllAsyncProcessors(Collection<IProcessor> asyncProcessors){
Map<String, ProcessorAction> actionMap = new HashMap<String, ProcessorAction>();
for(IProcessor processor: asyncProcessors){
actionMap.put(processor.id(), new ProcessorAction(queue.buildNewQueue(processor)));
}
Collection<ProcessorAction> actions = actionMap.values();
for(ProcessorAction action : actions){
action.fork();
}
ProcessorExecuteService result = queue.processorService;
for(Map.Entry<String, ProcessorAction> entry : actionMap.entrySet()){
result.addResult(entry.getKey(),entry.getValue().join());
}
}
}