MCPcopy Create free account
hub / github.com/WinVector/Logistic / reduce

Method reduce

src/com/winvector/util/ThreadedReducer.java:61–101  ·  view source on GitHub ↗

@param dat data source @param serialObserver (option can be null) a cheap observer that will be applied to all data serially @param parallelObserver an expensive observer that will be applied in parallel

(final Iterable<? extends T> dat, final S serialObserver, final Z parallelObserver)

Source from the content-addressed store, hash-verified

59 * @param parallelObserver an expensive observer that will be applied in parallel
60 */
61 public void reduce(final Iterable<? extends T> dat, final S serialObserver, final Z parallelObserver) {
62 final Ticker ticker = new Ticker(logString);
63 if(parallelism>1) {
64 final BlockingQueue<Runnable> workQueue = new ArrayBlockingQueue<Runnable>(2*parallelism + 10);
65 final ThreadPoolExecutor executor = new ThreadPoolExecutor(parallelism,parallelism,1000L,TimeUnit.SECONDS,workQueue);
66 ArrayList<T> al = new ArrayList<T>(gulpSize);
67 for(final T di: dat) {
68 ticker.tick();
69 if(serialObserver!=null) {
70 serialObserver.observe(di);
71 }
72 al.add(di);
73 if(al.size()>=gulpSize) {
74 if((executor==null)||(executor.getTaskCount()-executor.getCompletedTaskCount()>parallelism)) {
75 new EJob(parallelObserver,al).run();
76 } else {
77 executor.execute(new EJob(parallelObserver,al));
78 }
79 al = new ArrayList<T>(gulpSize);
80 }
81 }
82 if(!al.isEmpty()) {
83 new EJob(parallelObserver,al).run();
84 }
85 al = null;
86 executor.shutdown();
87 while(!executor.isTerminated()) {
88 try {
89 Thread.sleep(200L);
90 } catch (InterruptedException e) {
91 }
92 }
93 } else {
94 for(final T di: dat) {
95 if(serialObserver!=null) {
96 serialObserver.observe(di);
97 }
98 parallelObserver.observe(di);
99 }
100 }
101 }
102
103 /**
104 *

Callers 3

evalMethod · 0.95
accuracyMethod · 0.95

Calls 4

tickMethod · 0.95
observeMethod · 0.65
addMethod · 0.45
runMethod · 0.45

Tested by

no test coverage detected