Updated Branches: refs/heads/0.2-dev 6723b7ffa -> 6a667cbc8
Fixes for upgrade to lucene 4.1.0. Also added an inputindex reference watcher to try and close old indexinputs that are not properly closed by lucene itself. Project: http://git-wip-us.apache.org/repos/asf/incubator-blur/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-blur/commit/6a667cbc Tree: http://git-wip-us.apache.org/repos/asf/incubator-blur/tree/6a667cbc Diff: http://git-wip-us.apache.org/repos/asf/incubator-blur/diff/6a667cbc Branch: refs/heads/0.2-dev Commit: 6a667cbc8a150ce54b5b83beb3c57599a62eb0b6 Parents: 6723b7f Author: Aaron McCurry <[email protected]> Authored: Sun Mar 3 22:25:49 2013 -0500 Committer: Aaron McCurry <[email protected]> Committed: Sun Mar 3 22:25:49 2013 -0500 ---------------------------------------------------------------------- .../indexserver/DistributedIndexServer.java | 6 + .../apache/blur/manager/writer/BlurNRTIndex.java | 12 ++- .../org/apache/blur/lucene/search/FacetQuery.java | 2 +- .../org/apache/blur/lucene/search/SlowQuery.java | 2 +- .../org/apache/blur/lucene/search/SuperQuery.java | 2 +- .../refcounter/DirectoryReferenceCounter.java | 50 +++++---- .../store/refcounter/DirectoryReferenceFileGC.java | 3 +- .../lucene/store/refcounter/IndexInputCloser.java | 80 +++++++++++++++ .../apache/blur/store/blockcache/BlockLocks.java | 5 +- src/pom.xml | 4 +- 10 files changed, 129 insertions(+), 37 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-core/src/main/java/org/apache/blur/manager/indexserver/DistributedIndexServer.java ---------------------------------------------------------------------- diff --git a/src/blur-core/src/main/java/org/apache/blur/manager/indexserver/DistributedIndexServer.java b/src/blur-core/src/main/java/org/apache/blur/manager/indexserver/DistributedIndexServer.java index 881f7f7..b6009dc 100644 --- a/src/blur-core/src/main/java/org/apache/blur/manager/indexserver/DistributedIndexServer.java +++ b/src/blur-core/src/main/java/org/apache/blur/manager/indexserver/DistributedIndexServer.java @@ -43,6 +43,7 @@ import org.apache.blur.concurrent.Executors; import org.apache.blur.log.Log; import org.apache.blur.log.LogFactory; import org.apache.blur.lucene.store.refcounter.DirectoryReferenceFileGC; +import org.apache.blur.lucene.store.refcounter.IndexInputCloser; import org.apache.blur.manager.BlurFilterCache; import org.apache.blur.manager.clusterstatus.ClusterStatus; import org.apache.blur.manager.clusterstatus.ZookeeperPathConstants; @@ -107,6 +108,7 @@ public class DistributedIndexServer extends AbstractIndexServer { private DirectoryReferenceFileGC _gc; private WatchChildren _watchOnlineShards; private SharedMergeScheduler _mergeScheduler; + private IndexInputCloser _closer; public static interface ReleaseReader { void release() throws IOException; @@ -118,6 +120,8 @@ public class DistributedIndexServer extends AbstractIndexServer { _gc = new DirectoryReferenceFileGC(); _gc.init(); _mergeScheduler = new SharedMergeScheduler(); + _closer = new IndexInputCloser(); + _closer.init(); setupFlushCacheTimer(); String lockPath = BlurUtil.lockForSafeMode(_zookeeper, getNodeName(), _cluster); try { @@ -377,6 +381,7 @@ public class DistributedIndexServer extends AbstractIndexServer { _timerTableWarmer.purge(); _timerTableWarmer.cancel(); _gc.close(); + _closer.close(); _openerService.shutdownNow(); } } @@ -469,6 +474,7 @@ public class DistributedIndexServer extends AbstractIndexServer { writer.setDirectory(dir); writer.setGc(_gc); writer.setMergeScheduler(_mergeScheduler); + writer.setCloser(_closer); writer.init(); index = writer; } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-core/src/main/java/org/apache/blur/manager/writer/BlurNRTIndex.java ---------------------------------------------------------------------- diff --git a/src/blur-core/src/main/java/org/apache/blur/manager/writer/BlurNRTIndex.java b/src/blur-core/src/main/java/org/apache/blur/manager/writer/BlurNRTIndex.java index d0bfced..9decdef 100644 --- a/src/blur-core/src/main/java/org/apache/blur/manager/writer/BlurNRTIndex.java +++ b/src/blur-core/src/main/java/org/apache/blur/manager/writer/BlurNRTIndex.java @@ -30,6 +30,7 @@ import org.apache.blur.log.Log; import org.apache.blur.log.LogFactory; import org.apache.blur.lucene.store.refcounter.DirectoryReferenceCounter; import org.apache.blur.lucene.store.refcounter.DirectoryReferenceFileGC; +import org.apache.blur.lucene.store.refcounter.IndexInputCloser; import org.apache.blur.server.IndexSearcherClosable; import org.apache.blur.server.ShardContext; import org.apache.blur.server.TableContext; @@ -37,7 +38,6 @@ import org.apache.blur.thrift.generated.Document; import org.apache.blur.thrift.generated.Query; import org.apache.blur.thrift.generated.Term; import org.apache.blur.thrift.generated.UpdatePackage; -import org.apache.lucene.codecs.appending.AppendingCodec; import org.apache.lucene.index.CorruptIndexException; import org.apache.lucene.index.IndexReader; import org.apache.lucene.index.IndexWriterConfig; @@ -73,7 +73,7 @@ public class BlurNRTIndex extends BlurIndex { // created private TransactionRecorder _recorder; private TrackingIndexWriter _trackingWriter; - + private IndexInputCloser _closer; public void init() throws IOException { tableContext = shardContext.getTableContext(); @@ -82,11 +82,11 @@ public class BlurNRTIndex extends BlurIndex { conf.setWriteLockTimeout(TimeUnit.MINUTES.toMillis(5)); conf.setSimilarity(tableContext.getSimilarity()); conf.setIndexDeletionPolicy(tableContext.getIndexDeletionPolicy()); - conf.setCodec(new AppendingCodec()); + // conf.setCodec(new AppendingCodec()); TieredMergePolicy mergePolicy = (TieredMergePolicy) conf.getMergePolicy(); mergePolicy.setUseCompoundFile(false); conf.setMergeScheduler(mergeScheduler); - DirectoryReferenceCounter referenceCounter = new DirectoryReferenceCounter(_directory, _gc); + DirectoryReferenceCounter referenceCounter = new DirectoryReferenceCounter(_directory, _gc, _closer); _writer = new IndexWriter(referenceCounter, conf); _recorder = new TransactionRecorder(); _recorder.setContext(shardContext); @@ -264,4 +264,8 @@ public class BlurNRTIndex extends BlurIndex { public void setMergeScheduler(SharedMergeScheduler mergeScheduler) { this.mergeScheduler = mergeScheduler; } + + public void setCloser(IndexInputCloser closer) { + _closer = closer; + } } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-query/src/main/java/org/apache/blur/lucene/search/FacetQuery.java ---------------------------------------------------------------------- diff --git a/src/blur-query/src/main/java/org/apache/blur/lucene/search/FacetQuery.java b/src/blur-query/src/main/java/org/apache/blur/lucene/search/FacetQuery.java index 38cc866..1841fe9 100644 --- a/src/blur-query/src/main/java/org/apache/blur/lucene/search/FacetQuery.java +++ b/src/blur-query/src/main/java/org/apache/blur/lucene/search/FacetQuery.java @@ -193,7 +193,7 @@ public class FacetQuery extends AbstractWrapperQuery { } @Override - public float freq() throws IOException { + public int freq() throws IOException { return baseScorer.freq(); } } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-query/src/main/java/org/apache/blur/lucene/search/SlowQuery.java ---------------------------------------------------------------------- diff --git a/src/blur-query/src/main/java/org/apache/blur/lucene/search/SlowQuery.java b/src/blur-query/src/main/java/org/apache/blur/lucene/search/SlowQuery.java index 46ee4f5..d333e30 100644 --- a/src/blur-query/src/main/java/org/apache/blur/lucene/search/SlowQuery.java +++ b/src/blur-query/src/main/java/org/apache/blur/lucene/search/SlowQuery.java @@ -124,7 +124,7 @@ public class SlowQuery extends Query { return scorer.score(); } - public float freq() throws IOException { + public int freq() throws IOException { return scorer.freq(); } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-query/src/main/java/org/apache/blur/lucene/search/SuperQuery.java ---------------------------------------------------------------------- diff --git a/src/blur-query/src/main/java/org/apache/blur/lucene/search/SuperQuery.java b/src/blur-query/src/main/java/org/apache/blur/lucene/search/SuperQuery.java index cc0d484..72be513 100644 --- a/src/blur-query/src/main/java/org/apache/blur/lucene/search/SuperQuery.java +++ b/src/blur-query/src/main/java/org/apache/blur/lucene/search/SuperQuery.java @@ -305,7 +305,7 @@ public class SuperQuery extends AbstractWrapperQuery { } @Override - public float freq() throws IOException { + public int freq() throws IOException { return scorer.freq(); } } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceCounter.java ---------------------------------------------------------------------- diff --git a/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceCounter.java b/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceCounter.java index 5ac1633..7e03974 100644 --- a/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceCounter.java +++ b/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceCounter.java @@ -36,12 +36,23 @@ public class DirectoryReferenceCounter extends Directory { private final static Log LOG = LogFactory.getLog(DirectoryReferenceCounter.class); private Directory directory; - private Map<String, AtomicInteger> refs = new ConcurrentHashMap<String, AtomicInteger>(); + private Map<String, AtomicInteger> refCounters = new ConcurrentHashMap<String, AtomicInteger>(); private DirectoryReferenceFileGC gc; + private IndexInputCloser closer; - public DirectoryReferenceCounter(Directory directory, DirectoryReferenceFileGC gc) { + public DirectoryReferenceCounter(Directory directory, DirectoryReferenceFileGC gc, IndexInputCloser closer) { this.directory = directory; this.gc = gc; + this.closer = closer; + } + + private IndexInput wrap(String name, IndexInput input) { + AtomicInteger counter = refCounters.get(name); + if (counter == null) { + counter = new AtomicInteger(); + refCounters.put(name, counter); + } + return new RefIndexInput(input.toString(), input, counter, closer); } public void deleteFile(String name) throws IOException { @@ -49,7 +60,7 @@ public class DirectoryReferenceCounter extends Directory { deleteFile(name); return; } - AtomicInteger counter = refs.get(name); + AtomicInteger counter = refCounters.get(name); if (counter != null && counter.get() > 0) { addToFileGC(name); } else { @@ -64,13 +75,12 @@ public class DirectoryReferenceCounter extends Directory { return directory.createOutput(name, context); } LOG.debug("Create file [{0}]", name); - AtomicInteger counter = refs.get(name); + AtomicInteger counter = refCounters.get(name); if (counter != null) { LOG.error("Unknown error while trying to create ref counter for [{0}] reference exists.", name); throw new IOException("Reference exists [" + name + "]"); } - counter = new AtomicInteger(0); - refs.put(name, counter); + refCounters.put(name, new AtomicInteger(0)); return directory.createOutput(name, context); } @@ -88,31 +98,33 @@ public class DirectoryReferenceCounter extends Directory { private IndexInput input; private AtomicInteger ref; private boolean closed = false; - private Throwable throwable; + private IndexInputCloser closer; - public RefIndexInput(String resourceDescription, IndexInput input, AtomicInteger ref) { + public RefIndexInput(String resourceDescription, IndexInput input, AtomicInteger ref, IndexInputCloser closer) { super(resourceDescription); this.input = input; this.ref = ref; + this.closer = closer; ref.incrementAndGet(); + closer.add(this); } @Override protected void finalize() throws Throwable { // Seems like not all the clones are being closed... if (!closed) { - LOG.warn("[" + input.toString() + "] Seems like not all the clones are being closed. Printing stack trace.", throwable); + LOG.debug("[{0}] Last resort closing.", input.toString()); close(); } } @Override public RefIndexInput clone() { - RefIndexInput ref = (RefIndexInput) super.clone(); - ref.input = (IndexInput) input.clone(); - ref.ref.incrementAndGet(); - ref.throwable = new Throwable(); - return ref; + RefIndexInput refIndexInput = (RefIndexInput) super.clone(); + closer.add(refIndexInput); + refIndexInput.input = (IndexInput) input.clone(); + refIndexInput.ref.incrementAndGet(); + return refIndexInput; } @Override @@ -259,16 +271,8 @@ public class DirectoryReferenceCounter extends Directory { private void addToFileGC(String name) { if (gc != null) { LOG.debug("Add file [{0}] to be GCed once refs are closed.", name); - gc.add(directory, name, refs); + gc.add(directory, name, refCounters); } } - private IndexInput wrap(String name, IndexInput input) { - AtomicInteger counter = refs.get(name); - if (counter == null) { - counter = new AtomicInteger(); - refs.put(name, counter); - } - return new RefIndexInput(input.toString(), input, counter); - } } \ No newline at end of file http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceFileGC.java ---------------------------------------------------------------------- diff --git a/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceFileGC.java b/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceFileGC.java index ec95633..9db6256 100644 --- a/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceFileGC.java +++ b/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/DirectoryReferenceFileGC.java @@ -28,7 +28,6 @@ import org.apache.blur.log.Log; import org.apache.blur.log.LogFactory; import org.apache.lucene.store.Directory; - public class DirectoryReferenceFileGC extends TimerTask { private static final Log LOG = LogFactory.getLog(DirectoryReferenceFileGC.class); @@ -56,7 +55,7 @@ public class DirectoryReferenceFileGC extends TimerTask { directory.deleteFile(name); return true; } else { - LOG.info("File [{0}] had too many refs [{1}]", name, counter.get()); + LOG.debug("File [{0}] had too many refs [{1}]", name, counter.get()); } return false; } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/IndexInputCloser.java ---------------------------------------------------------------------- diff --git a/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/IndexInputCloser.java b/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/IndexInputCloser.java new file mode 100644 index 0000000..e6cd5a0 --- /dev/null +++ b/src/blur-store/src/main/java/org/apache/blur/lucene/store/refcounter/IndexInputCloser.java @@ -0,0 +1,80 @@ +package org.apache.blur.lucene.store.refcounter; + +import java.io.Closeable; +import java.io.IOException; +import java.lang.ref.ReferenceQueue; +import java.lang.ref.WeakReference; +import java.util.Collection; +import java.util.HashSet; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.apache.blur.log.Log; +import org.apache.blur.log.LogFactory; +import org.apache.lucene.store.IndexInput; +import org.apache.lucene.util.IOUtils; + +public class IndexInputCloser implements Runnable { + + private static final Log LOG = LogFactory.getLog(IndexInputCloser.class); + + private ReferenceQueue<IndexInput> referenceQueue = new ReferenceQueue<IndexInput>(); + private Thread daemon; + private AtomicBoolean running = new AtomicBoolean(); + private Collection<IndexInputCloserRef> refs = new HashSet<IndexInputCloserRef>(); + + static class IndexInputCloserRef extends WeakReference<IndexInput> implements Closeable { + public IndexInputCloserRef(IndexInput referent, ReferenceQueue<? super IndexInput> q) { + super(referent, q); + } + + @Override + public void close() throws IOException { + IndexInput input = get(); + if (input != null) { + LOG.debug("Closing indexinput [{0}]", input); + input.close(); + } + } + } + + public void init() { + running.set(true); + daemon = new Thread(this); + daemon.setDaemon(true); + daemon.setName("IndexIndexCloser"); + daemon.start(); + } + + public void add(IndexInput indexInput) { + LOG.debug("Adding [{0}]", indexInput); + IndexInputCloserRef ref = new IndexInputCloserRef(indexInput, referenceQueue); + synchronized (refs) { + refs.add(ref); + } + } + + public void close() { + running.set(false); + refs.clear(); + daemon.interrupt(); + } + + @Override + public void run() { + while (running.get()) { + try { + IndexInputCloserRef ref = (IndexInputCloserRef) referenceQueue.remove(); + LOG.debug("Closing [{0}]", ref); + IOUtils.closeWhileHandlingException(ref); + synchronized (refs) { + refs.remove(ref); + } + } catch (InterruptedException e) { + LOG.info("Interrupted"); + running.set(false); + return; + } + } + } + +} http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/blur-store/src/main/java/org/apache/blur/store/blockcache/BlockLocks.java ---------------------------------------------------------------------- diff --git a/src/blur-store/src/main/java/org/apache/blur/store/blockcache/BlockLocks.java b/src/blur-store/src/main/java/org/apache/blur/store/blockcache/BlockLocks.java index 502d786..b25cd12 100644 --- a/src/blur-store/src/main/java/org/apache/blur/store/blockcache/BlockLocks.java +++ b/src/blur-store/src/main/java/org/apache/blur/store/blockcache/BlockLocks.java @@ -18,7 +18,6 @@ package org.apache.blur.store.blockcache; */ import java.util.concurrent.atomic.AtomicLongArray; -import org.apache.lucene.util.BitUtil; import org.apache.lucene.util.OpenBitSet; public class BlockLocks { @@ -46,12 +45,12 @@ public class BlockLocks { long word = ~bits.get(i) >> subIndex; // skip all the bits to the right of // index if (word != 0) { - return (i << 6) + subIndex + BitUtil.ntz(word); + return (i << 6) + subIndex + Long.numberOfTrailingZeros(word); } while (++i < wlen) { word = ~bits.get(i); if (word != 0) { - return (i << 6) + BitUtil.ntz(word); + return (i << 6) + Long.numberOfTrailingZeros(word); } } return -1; http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/6a667cbc/src/pom.xml ---------------------------------------------------------------------- diff --git a/src/pom.xml b/src/pom.xml index eb88760..833d404 100644 --- a/src/pom.xml +++ b/src/pom.xml @@ -42,12 +42,12 @@ under the License. <zookeeper.version>3.4.5</zookeeper.version> <log4j.version>1.2.15</log4j.version> <jersey.version>1.14</jersey.version> - <lucene.version>4.0.0</lucene.version> + <lucene.version>4.1.0</lucene.version> <junit.version>4.7</junit.version> <thrift.version>0.9.0</thrift.version> <slf4j.version>1.6.1</slf4j.version> <commons-cli.version>1.2</commons-cli.version> - <concurrentlinkedhashmap-lru.version>1.3.1</concurrentlinkedhashmap-lru.version> + <concurrentlinkedhashmap-lru.version>1.3.2</concurrentlinkedhashmap-lru.version> <jline.version>2.7</jline.version> <guava.version>12.0</guava.version> </properties>
