Updated Branches: refs/heads/apache-blur-0.2 ccdd40fcd -> ed2da2c05
Removing the need for the reference counter directory. Project: http://git-wip-us.apache.org/repos/asf/incubator-blur/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-blur/commit/b98ad9d8 Tree: http://git-wip-us.apache.org/repos/asf/incubator-blur/tree/b98ad9d8 Diff: http://git-wip-us.apache.org/repos/asf/incubator-blur/diff/b98ad9d8 Branch: refs/heads/apache-blur-0.2 Commit: b98ad9d8a4b20d49791ec07d4b67fe33fba4af22 Parents: ccdd40f Author: Aaron McCurry <[email protected]> Authored: Fri Jan 31 20:03:46 2014 -0500 Committer: Aaron McCurry <[email protected]> Committed: Fri Jan 31 20:03:46 2014 -0500 ---------------------------------------------------------------------- .../manager/writer/BlurIndexSimpleWriter.java | 15 ++-- .../blur/index/IndexDeletionPolicyReader.java | 87 ++++++++++++++++++++ 2 files changed, 97 insertions(+), 5 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/b98ad9d8/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java b/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java index d9093dc..12476b0 100644 --- a/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java +++ b/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java @@ -30,10 +30,10 @@ import java.util.concurrent.locks.ReentrantReadWriteLock; import org.apache.blur.analysis.FieldManager; import org.apache.blur.index.ExitableReader; +import org.apache.blur.index.IndexDeletionPolicyReader; import org.apache.blur.log.Log; import org.apache.blur.log.LogFactory; import org.apache.blur.lucene.codec.Blur022Codec; -import org.apache.blur.lucene.store.refcounter.DirectoryReferenceCounter; import org.apache.blur.lucene.store.refcounter.DirectoryReferenceFileGC; import org.apache.blur.lucene.warmup.TraceableDirectory; import org.apache.blur.manager.indexserver.BlurIndexWarmup; @@ -48,6 +48,7 @@ import org.apache.lucene.index.BlurIndexWriter; import org.apache.lucene.index.DirectoryReader; import org.apache.lucene.index.IndexReader; import org.apache.lucene.index.IndexWriterConfig; +import org.apache.lucene.index.KeepOnlyLastCommitDeletionPolicy; import org.apache.lucene.index.TieredMergePolicy; import org.apache.lucene.store.Directory; @@ -72,6 +73,8 @@ public class BlurIndexSimpleWriter extends BlurIndex { private Thread _optimizeThread; private Thread _writerOpener; + private IndexDeletionPolicyReader _policy; + public BlurIndexSimpleWriter(ShardContext shardContext, Directory directory, SharedMergeScheduler mergeScheduler, DirectoryReferenceFileGC gc, final ExecutorService searchExecutor, BlurIndexCloser indexCloser, BlurIndexRefresher refresher, BlurIndexWarmup indexWarmup) throws IOException { @@ -89,13 +92,15 @@ public class BlurIndexSimpleWriter extends BlurIndex { TieredMergePolicy mergePolicy = (TieredMergePolicy) _conf.getMergePolicy(); mergePolicy.setUseCompoundFile(false); _conf.setMergeScheduler(mergeScheduler.getMergeScheduler()); + _policy = new IndexDeletionPolicyReader(new KeepOnlyLastCommitDeletionPolicy()); + _conf.setIndexDeletionPolicy(_policy); if (!DirectoryReader.indexExists(directory)) { new BlurIndexWriter(directory, _conf).close(); } - DirectoryReferenceCounter referenceCounter = new DirectoryReferenceCounter(directory, gc); + // This directory allows for warm up by adding tracing ability. - TraceableDirectory dir = new TraceableDirectory(referenceCounter); + TraceableDirectory dir = new TraceableDirectory(directory); _directory = dir; _indexCloser = indexCloser; @@ -119,11 +124,11 @@ public class BlurIndexSimpleWriter extends BlurIndex { _writerOpener.start(); } - private DirectoryReader wrap(DirectoryReader reader) { + private DirectoryReader wrap(DirectoryReader reader) throws IOException { if (_makeReaderExitable) { reader = new ExitableReader(reader); } - return reader; + return _policy.register(reader); } private Thread getWriterOpener(ShardContext shardContext) { http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/b98ad9d8/blur-store/src/main/java/org/apache/blur/index/IndexDeletionPolicyReader.java ---------------------------------------------------------------------- diff --git a/blur-store/src/main/java/org/apache/blur/index/IndexDeletionPolicyReader.java b/blur-store/src/main/java/org/apache/blur/index/IndexDeletionPolicyReader.java new file mode 100644 index 0000000..2bca357 --- /dev/null +++ b/blur-store/src/main/java/org/apache/blur/index/IndexDeletionPolicyReader.java @@ -0,0 +1,87 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.blur.index; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + +import org.apache.lucene.index.DirectoryReader; +import org.apache.lucene.index.IndexCommit; +import org.apache.lucene.index.IndexDeletionPolicy; +import org.apache.lucene.index.IndexReader; +import org.apache.lucene.index.IndexReader.ReaderClosedListener; + +public class IndexDeletionPolicyReader extends IndexDeletionPolicy { + + private final IndexDeletionPolicy _base; + private final Set<Long> _gens = Collections.newSetFromMap(new ConcurrentHashMap<Long, Boolean>()); + + public IndexDeletionPolicyReader(IndexDeletionPolicy base) { + _base = base; + } + + @Override + public void onInit(List<? extends IndexCommit> commits) throws IOException { + _base.onInit(commits); + } + + @Override + public void onCommit(List<? extends IndexCommit> commits) throws IOException { + commits = removeCommitsStillInUse(commits); + _base.onCommit(commits); + } + + private List<? extends IndexCommit> removeCommitsStillInUse(List<? extends IndexCommit> commits) { + List<IndexCommit> validForRemoval = new ArrayList<IndexCommit>(); + for (IndexCommit commit : commits) { + if (!isStillInUse(commit)) { + validForRemoval.add(commit); + } + } + return validForRemoval; + } + + public DirectoryReader register(DirectoryReader reader) throws IOException { + final long generation = reader.getIndexCommit().getGeneration(); + register(generation); + reader.addReaderClosedListener(new ReaderClosedListener() { + @Override + public void onClose(IndexReader reader) { + unregister(generation); + } + }); + return reader; + } + + private boolean isStillInUse(IndexCommit commit) { + long generation = commit.getGeneration(); + return _gens.contains(generation); + } + + public void register(long gen) { + _gens.add(gen); + } + + public void unregister(long gen) { + _gens.remove(gen); + } + +}
