Repository: incubator-blur
Updated Branches:
  refs/heads/apache-blur-0.2 d166c5b32 -> 3fc741cae


Adding a lock per facet during the execution phase to reduce extra work if 
minimums are used.


Project: http://git-wip-us.apache.org/repos/asf/incubator-blur/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-blur/commit/3fc741ca
Tree: http://git-wip-us.apache.org/repos/asf/incubator-blur/tree/3fc741ca
Diff: http://git-wip-us.apache.org/repos/asf/incubator-blur/diff/3fc741ca

Branch: refs/heads/apache-blur-0.2
Commit: 3fc741caea593ef8ebaf47e15abe26671f44c398
Parents: d166c5b
Author: Aaron McCurry <[email protected]>
Authored: Thu Mar 6 16:48:38 2014 -0500
Committer: Aaron McCurry <[email protected]>
Committed: Thu Mar 6 16:48:38 2014 -0500

----------------------------------------------------------------------
 .../blur/lucene/search/FacetExecutor.java       | 48 +++++++++++++++++---
 1 file changed, 41 insertions(+), 7 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/3fc741ca/blur-query/src/main/java/org/apache/blur/lucene/search/FacetExecutor.java
----------------------------------------------------------------------
diff --git 
a/blur-query/src/main/java/org/apache/blur/lucene/search/FacetExecutor.java 
b/blur-query/src/main/java/org/apache/blur/lucene/search/FacetExecutor.java
index 587695b..7359857 100644
--- a/blur-query/src/main/java/org/apache/blur/lucene/search/FacetExecutor.java
+++ b/blur-query/src/main/java/org/apache/blur/lucene/search/FacetExecutor.java
@@ -23,10 +23,14 @@ import java.util.Comparator;
 import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
+import java.util.concurrent.ArrayBlockingQueue;
+import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.AtomicLongArray;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 import org.apache.blur.log.Log;
 import org.apache.blur.log.LogFactory;
@@ -90,19 +94,21 @@ public class FacetExecutor {
     final AtomicReader _reader;
     final String _readerStr;
     final int _maxDoc;
+    final Lock[] _locks;
 
     @Override
     public String toString() {
       return "Info scorers length [" + _scorers.length + "] reader[" + _reader 
+ "]";
     }
 
-    Info(AtomicReaderContext context, Scorer[] scorers) {
+    Info(AtomicReaderContext context, Scorer[] scorers, Lock[] locks) {
       AtomicReader reader = context.reader();
       _bitSet = new OpenBitSet(reader.maxDoc());
       _scorers = scorers;
       _reader = reader;
       _readerStr = _reader.toString();
       _maxDoc = _reader.maxDoc();
+      _locks = locks;
     }
 
     void process(AtomicLongArray counts, long[] minimumsBeforeReturning) 
throws IOException {
@@ -118,16 +124,39 @@ public class FacetExecutor {
           trace.done();
         }
       } else {
-        for (int i = 0; i < _scorers.length; i++) {
-          long min = minimumsBeforeReturning[i];
-          long currentCount = counts.get(i);
-          if (currentCount < min) {
-            runFacet(counts, col, i);
+        BlockingQueue<Integer> ids = new 
ArrayBlockingQueue<Integer>(_scorers.length + 1);
+        try {
+          populate(ids);
+          while (!ids.isEmpty()) {
+            int id = ids.take();
+            Lock lock = _locks[id];
+            if (lock.tryLock()) {
+              try {
+                long min = minimumsBeforeReturning[id];
+                long currentCount = counts.get(id);
+                if (currentCount < min) {
+                  runFacet(counts, col, id);
+                }
+              } finally {
+                lock.unlock();
+              }
+            } else {
+              ids.put(id);
+            }
           }
+        } catch (Exception e) {
+          throw new IOException(e);
         }
       }
     }
 
+    private void populate(BlockingQueue<Integer> ids) throws 
InterruptedException {
+      for (int i = 0; i < _scorers.length; i++) {
+        ids.put(i);
+      }
+
+    }
+
     private void runFacet(AtomicLongArray counts, SimpleCollector col, int i) 
throws IOException {
       Scorer scorer = _scorers[i];
       if (scorer != null) {
@@ -148,6 +177,7 @@ public class FacetExecutor {
   private final int _length;
   private final AtomicLongArray _counts;
   private final long[] _minimumsBeforeReturning;
+  private final Lock[] _locks;
   private boolean _processed;
 
   public FacetExecutor(int length) {
@@ -158,6 +188,10 @@ public class FacetExecutor {
     _length = length;
     _counts = counts;
     _minimumsBeforeReturning = minimumsBeforeReturning;
+    _locks = new Lock[_length];
+    for (int i = 0; i < _length; i++) {
+      _locks[i] = new ReentrantReadWriteLock().writeLock();
+    }
   }
 
   public FacetExecutor(int length, long[] minimumsBeforeReturning) {
@@ -171,7 +205,7 @@ public class FacetExecutor {
     Object key = getKey(context);
     Info info = _infoMap.get(key);
     if (info == null) {
-      info = new Info(context, scorers);
+      info = new Info(context, scorers, _locks);
       _infoMap.put(key, info);
     } else {
       AtomicReader reader = context.reader();

Reply via email to