@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)
| 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 | * |