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();
