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>

Reply via email to