map reduce index: make ForwardIndex responsible to create a data diff (allow to optimize diff calculation)

This commit is contained in:
Dmitry Batkovich
2017-01-13 10:46:12 +03:00
parent 67964b82d4
commit a2f9625f21
21 changed files with 423 additions and 431 deletions
@@ -573,12 +573,9 @@ public class StubIndexImpl extends StubIndex implements ApplicationComponentAdap
@NotNull final Map<K, StubIdList> newValues) {
try {
final MyIndex<K> index = (MyIndex<K>)getAsyncState().myIndices.get(key);
final ThrowableComputable<ForwardIndex.InputKeyIterator<K, StubIdList>, IOException>
oldMapGetter = () -> new MapInputKeyIterator<>(oldValues);
index.updateWithMap(fileId,
DiffUpdateData.ourDiffUpdateEnabled
? new DiffUpdateData<>(newValues, oldMapGetter, key, null)
: new SimpleUpdateData<>(newValues, oldMapGetter, key, null));
final ThrowableComputable<InputDataDiffBuilder<K, StubIdList>, IOException>
oldMapGetter = () -> new MapInputDataDiffBuilder<>(fileId, oldValues);
index.updateWithMap(fileId, new UpdateData<>(newValues, oldMapGetter, key, null));
}
catch (StorageException e) {
LOG.info(e);
@@ -39,6 +39,7 @@ import com.intellij.util.indexing.impl.*;
import com.intellij.util.io.*;
import gnu.trove.THashMap;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.io.*;
import java.io.DataOutputStream;
@@ -406,29 +407,29 @@ public class StubUpdatingIndex extends CustomImplementationFileBasedIndexExtensi
@NotNull
@Override
protected UpdateData<Integer, SerializedStubTree> createUpdateData(Map<Integer, SerializedStubTree> data,
ThrowableComputable<ForwardIndex.InputKeyIterator<Integer, SerializedStubTree>, IOException> oldKeys,
ThrowableComputable<InputDataDiffBuilder<Integer, SerializedStubTree>, IOException> oldKeys,
ThrowableRunnable<IOException> forwardIndexUpdate) {
return new StubUpdatingData(data, oldKeys, forwardIndexUpdate);
}
static class StubUpdatingData extends SimpleUpdateData<Integer, SerializedStubTree> {
static class StubUpdatingData extends UpdateData<Integer, SerializedStubTree> {
private Collection<Integer> oldStubIndexKeys;
public StubUpdatingData(@NotNull Map<Integer, SerializedStubTree> newData,
@NotNull ThrowableComputable<ForwardIndex.InputKeyIterator<Integer, SerializedStubTree>, IOException> iterator,
ThrowableRunnable<IOException> forwardIndexUpdate) {
@NotNull ThrowableComputable<InputDataDiffBuilder<Integer, SerializedStubTree>, IOException> iterator,
@Nullable ThrowableRunnable<IOException> forwardIndexUpdate) {
super(newData, iterator, INDEX_ID, forwardIndexUpdate);
}
@Override
protected void iterateKeys(int inputId,
KeyValueUpdateProcessor<Integer, SerializedStubTree> addProcessor,
RemovedKeyProcessor<Integer> removeProcessor,
ForwardIndex.InputKeyIterator<Integer, SerializedStubTree> currentData) throws StorageException {
if (currentData instanceof CollectionInputKeyIterator) {
oldStubIndexKeys = ((CollectionInputKeyIterator<Integer, SerializedStubTree>)currentData).getCollection();
}
super.iterateKeys(inputId, addProcessor, removeProcessor, currentData);
protected ThrowableComputable<InputDataDiffBuilder<Integer, SerializedStubTree>, IOException> getCurrentDataEvaluator() {
return () -> {
final InputDataDiffBuilder<Integer, SerializedStubTree> diffBuilder = super.getCurrentDataEvaluator().compute();
if (diffBuilder instanceof CollectionInputDataDiffBuilder) {
oldStubIndexKeys = ((CollectionInputDataDiffBuilder<Integer, SerializedStubTree>) diffBuilder).getSeq();
}
return diffBuilder;
};
}
public Map<StubIndexKey, Map<Object, StubIdList>> getOldStubIndicesValueMap() {
@@ -458,7 +459,6 @@ public class StubUpdatingIndex extends CustomImplementationFileBasedIndexExtensi
}
}
@Override
protected void updateWithMap(int inputId,
@NotNull UpdateData<Integer, SerializedStubTree> updateData) throws StorageException {
@@ -16,10 +16,7 @@
package com.intellij.util.indexing;
import com.intellij.openapi.vfs.newvfs.persistent.PersistentFS;
import com.intellij.util.indexing.impl.AbstractForwardIndex;
import com.intellij.util.indexing.impl.CollectionInputKeyIterator;
import com.intellij.util.indexing.impl.DebugAssertions;
import com.intellij.util.indexing.impl.MapBasedForwardIndex;
import com.intellij.util.indexing.impl.*;
import com.intellij.util.io.DataExternalizer;
import org.jetbrains.annotations.NotNull;
@@ -39,7 +36,7 @@ class SharedMapBasedForwardIndex<Key, Value> extends AbstractForwardIndex<Key,Va
@NotNull
@Override
public InputKeyIterator<Key, Value> getInputKeys(int inputId) throws IOException {
public InputDataDiffBuilder<Key, Value> getDiffBuilder(int inputId) throws IOException {
Collection<Key> keys;
if (SharedIndicesData.ourFileSharedIndicesEnabled) {
keys = SharedIndicesData.recallFileData(inputId, myIndexId, mySnapshotIndexExternalizer);
@@ -56,9 +53,9 @@ class SharedMapBasedForwardIndex<Key, Value> extends AbstractForwardIndex<Key,Va
}
keys = keysFromInputsIndex;
}
return new CollectionInputKeyIterator<>(keys);
return new CollectionInputDataDiffBuilder<>(inputId, keys);
}
return new CollectionInputKeyIterator<>(myUnderlying.getInputsIndex().get(inputId));
return new CollectionInputDataDiffBuilder<>(inputId, myUnderlying.getInputsIndex().get(inputId));
}
@Override
@@ -98,20 +98,20 @@ public class VfsAwareMapReduceIndex<Key, Value, Input> extends MapReduceIndex<Ke
}
return createUpdateData(data, () -> {
if (mySnapshotInputMappings != null && isContentPhysical) {
return new MapInputKeyIterator<>(mySnapshotInputMappings.readInputKeys(inputId));
return new MapInputDataDiffBuilder<>(inputId, mySnapshotInputMappings.readInputKeys(inputId));
}
if (myInMemoryMode.get()) {
synchronized (myInMemoryKeys) {
Collection<Key> keys = myInMemoryKeys.get(inputId);
if (keys != null) {
return new CollectionInputKeyIterator<>(keys);
return new CollectionInputDataDiffBuilder<>(inputId, keys);
}
}
}
if (myForwardIndex != null) {
return readInputKeys(inputId);
return getKeysDiffBuilder(inputId);
}
return EmptyInputKeyIterator.getInstance();
return new EmptyInputDataDiffBuilder(inputId);
}, () -> {
if (myInMemoryMode.get()) {
synchronized (myInMemoryKeys) {
@@ -35,8 +35,4 @@ public abstract class AbstractForwardIndex<Key, Value> implements ForwardIndex<K
public IndexExtension<Key, Value, ?> getIndexExtension() {
return myIndexExtension;
}
public boolean hasOnlyKeysData() {
return true;
}
}
@@ -0,0 +1,58 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.util.indexing.StorageException;
import org.jetbrains.annotations.ApiStatus;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.util.Collection;
import java.util.Collections;
import java.util.Map;
@ApiStatus.Experimental
public class CollectionInputDataDiffBuilder<Key, Value> extends InputDataDiffBuilder<Key,Value> {
private final Collection<Key> mySeq;
public CollectionInputDataDiffBuilder(int inputId, @Nullable Collection<Key> seq) {
super(inputId);
mySeq = seq == null ? Collections.<Key>emptySet() : seq;
}
@Override
public void differentiate(@NotNull Map<Key, Value> newData,
@NotNull KeyValueUpdateProcessor<Key, Value> addProcessor,
@NotNull KeyValueUpdateProcessor<Key, Value> updateProcessor,
@NotNull RemovedKeyProcessor<Key> removeProcessor) throws StorageException {
differentiateWithKeySeq(mySeq, newData, myInputId, addProcessor, removeProcessor);
}
public Collection<Key> getSeq() {
return mySeq;
}
static <Key, Value> void differentiateWithKeySeq(@NotNull Collection<Key> currentData,
@NotNull Map<Key, Value> newData,
int inputId,
@NotNull KeyValueUpdateProcessor<Key, Value> addProcessor,
@NotNull RemovedKeyProcessor<Key> removeProcessor) throws StorageException {
for (Key key : currentData) {
removeProcessor.process(key, inputId);
}
EmptyInputDataDiffBuilder.processKeys(newData, addProcessor, inputId);
}
}
@@ -1,62 +0,0 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import org.jetbrains.annotations.Nullable;
import java.util.Collection;
import java.util.Collections;
import java.util.Iterator;
public class CollectionInputKeyIterator<Key, Value> implements ForwardIndex.InputKeyIterator<Key, Value> {
private final Collection<Key> mySeq;
private Iterator<Key> myIt;
public CollectionInputKeyIterator(Collection<Key> seq) {
mySeq = seq;
}
@Override
public boolean isAssociatedValueEqual(@Nullable Value value) {
return false;
}
@Override
public boolean hasNext() {
init();
return myIt.hasNext();
}
@Override
public Key next() {
return myIt.next();
}
@Override
public void remove() {
throw new UnsupportedOperationException();
}
public Collection<Key> getCollection() {
return mySeq == null ? Collections.<Key>emptySet() : mySeq;
}
private void init() {
if (myIt == null) {
myIt = getCollection().iterator();
}
}
}
@@ -1,109 +0,0 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.util.ThrowableComputable;
import com.intellij.util.SystemProperties;
import com.intellij.util.ThrowableRunnable;
import com.intellij.util.indexing.ID;
import com.intellij.util.indexing.StorageException;
import gnu.trove.THashSet;
import org.jetbrains.annotations.NotNull;
import java.io.IOException;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicInteger;
public class DiffUpdateData<Key, Value> extends UpdateData<Key,Value> {
public static final boolean ourDiffUpdateEnabled = SystemProperties.getBooleanProperty("idea.disable.diff.index.update", true);
public DiffUpdateData(@NotNull Map<Key, Value> newData,
@NotNull ThrowableComputable<ForwardIndex.InputKeyIterator<Key, Value>, IOException> currentData,
@NotNull ID<Key, Value> indexId, ThrowableRunnable<IOException> forwardIndexUpdate) {
super(newData, currentData, indexId, forwardIndexUpdate);
}
@Override
public void iterateKeys(int inputId,
KeyValueUpdateProcessor<Key, Value> addProcessor,
KeyValueUpdateProcessor<Key, Value> updateProcessor,
RemovedKeyProcessor<Key> removeProcessor) throws StorageException {
final Set<Key> processedKeys = new THashSet<Key>();
int oldSize = 0; //kept for debug reasons
int addedKeys = 0;
int removedKeys = 0;
boolean newDataIsEmpty = myNewData.isEmpty();
final ForwardIndex.InputKeyIterator<Key, Value> currentData;
try {
currentData = myCurrentData.compute();
}
catch (IOException e) {
throw new StorageException(e);
}
while (currentData.hasNext()) {
oldSize++;
Key key = currentData.next();
if (!newDataIsEmpty) {
processedKeys.add(key);
}
if (newDataIsEmpty || !myNewData.containsKey(key)) {
removeProcessor.process(key, inputId);
removedKeys++;
} else {
Value newValue = myNewData.get(key);
if (!currentData.isAssociatedValueEqual(newValue)) {
updateProcessor.process(key, newValue, inputId);
removedKeys++;
addedKeys++;
}
}
}
if (!newDataIsEmpty) {
for (Map.Entry<Key, Value> entry : myNewData.entrySet()) {
if (!processedKeys.contains(entry.getKey())) {
addProcessor.process(entry.getKey(), entry.getValue(), inputId);
addedKeys++;
}
}
}
int totalRequests = requests.incrementAndGet();
totalRemovals.addAndGet(oldSize);
totalAdditions.addAndGet(myNewData.size());
incrementalAdditions.addAndGet(removedKeys);
incrementalRemovals.addAndGet(addedKeys);
if ((totalRequests & 0xFFF) == 0 && DebugAssertions.DEBUG) {
Logger.getInstance(getClass()).info("Incremental index diff update:"+requests +
", removals:" + totalRemovals + "->" + incrementalRemovals +
", additions:" +totalAdditions + "->" +incrementalAdditions);
}
}
private static final AtomicInteger requests = new AtomicInteger();
private static final AtomicInteger totalRemovals = new AtomicInteger();
private static final AtomicInteger totalAdditions = new AtomicInteger();
private static final AtomicInteger incrementalRemovals = new AtomicInteger();
private static final AtomicInteger incrementalAdditions = new AtomicInteger();
@NotNull
protected Map<Key, Value> getMap() {
return myNewData;
}
}
@@ -0,0 +1,69 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.util.indexing.StorageException;
import gnu.trove.THashMap;
import gnu.trove.TObjectObjectProcedure;
import org.jetbrains.annotations.ApiStatus;
import org.jetbrains.annotations.NotNull;
import java.util.Map;
@ApiStatus.Experimental
public class EmptyInputDataDiffBuilder<Key, Value> extends InputDataDiffBuilder<Key,Value> {
public EmptyInputDataDiffBuilder(int inputId) {
super(inputId);
}
@Override
public void differentiate(@NotNull Map<Key, Value> newData,
@NotNull final KeyValueUpdateProcessor<Key, Value> addProcessor,
@NotNull KeyValueUpdateProcessor<Key, Value> updateProcessor,
@NotNull RemovedKeyProcessor<Key> removeProcessor) throws StorageException {
processKeys(newData, addProcessor, myInputId);
}
static <Key, Value >void processKeys(@NotNull Map<Key, Value> currentData,
@NotNull final KeyValueUpdateProcessor<Key, Value> processor,
final int inputId)
throws StorageException {
if (currentData instanceof THashMap) {
final StorageException[] exception = new StorageException[]{null};
((THashMap<Key, Value>)currentData).forEachEntry(new TObjectObjectProcedure<Key, Value>() {
@Override
public boolean execute(Key k, Value v) {
try {
processor.process(k, v, inputId);
}
catch (StorageException e) {
exception[0] = e;
return false;
}
return true;
}
});
if (exception[0] != null) {
throw exception[0];
}
}
else {
for (Map.Entry<Key, Value> entry : currentData.entrySet()) {
processor.process(entry.getKey(), entry.getValue(), inputId);
}
}
}
}
@@ -1,49 +0,0 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import org.jetbrains.annotations.ApiStatus;
import org.jetbrains.annotations.Nullable;
@ApiStatus.Experimental
public class EmptyInputKeyIterator<Key, Value> implements ForwardIndex.InputKeyIterator<Key,Value> {
public static final EmptyInputKeyIterator EMPTY_INPUT_KEY_ITERATOR = new EmptyInputKeyIterator();
public static <Key, Value> ForwardIndex.InputKeyIterator<Key, Value> getInstance() {
//noinspection unchecked
return EMPTY_INPUT_KEY_ITERATOR;
}
@Override
public boolean isAssociatedValueEqual(@Nullable Value value) {
throw new UnsupportedOperationException();
}
@Override
public boolean hasNext() {
return false;
}
@Override
public Key next() {
throw new UnsupportedOperationException();
}
@Override
public void remove() {
throw new UnsupportedOperationException();
}
}
@@ -17,17 +17,24 @@ package com.intellij.util.indexing.impl;
import org.jetbrains.annotations.ApiStatus;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.io.IOException;
import java.util.Iterator;
import java.util.Map;
/**
* Represents a <a href="https://en.wikipedia.org/wiki/Search_engine_indexing#The_forward_index">forward index data structure</>:
* an index indented to hold a mappings of inputId-s to contained keys.
*/
@ApiStatus.Experimental
public interface ForwardIndex<Key, Value> {
/**
* Creates a diff builder for given inputId.
*/
@NotNull
InputKeyIterator<Key, Value> getInputKeys(int inputId) throws IOException;
InputDataDiffBuilder<Key, Value> getDiffBuilder(int inputId) throws IOException;
/**
* Update data for inputId.
*/
void putInputData(int inputId, @NotNull Map<Key, Value> data) throws IOException;
void flush();
@@ -35,9 +42,4 @@ public interface ForwardIndex<Key, Value> {
void clear() throws IOException;
void close() throws IOException;
@ApiStatus.Experimental
interface InputKeyIterator<Key, Value> extends Iterator<Key> {
boolean isAssociatedValueEqual(@Nullable Value value);
}
}
@@ -0,0 +1,39 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.util.indexing.StorageException;
import org.jetbrains.annotations.ApiStatus;
import org.jetbrains.annotations.NotNull;
import java.util.Map;
/**
* A class intended to make a diff between existing forward index data and new one.
*/
@ApiStatus.Experimental
public abstract class InputDataDiffBuilder<Key, Value> {
protected final int myInputId;
protected InputDataDiffBuilder(int id) {myInputId = id;}
/**
* produce a diff between existing data and newData and consume result to addProcessor, updateProcessor and removeProcessor.
*/
public abstract void differentiate(@NotNull Map<Key, Value> newData,
@NotNull KeyValueUpdateProcessor<Key, Value> addProcessor,
@NotNull KeyValueUpdateProcessor<Key, Value> updateProcessor,
@NotNull RemovedKeyProcessor<Key> removeProcessor) throws StorageException;
}
@@ -0,0 +1,24 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.util.indexing.StorageException;
import org.jetbrains.annotations.ApiStatus;
@ApiStatus.Experimental
public interface KeyValueUpdateProcessor<Key, Value> {
void process(Key key, Value value, int inputId) throws StorageException;
}
@@ -39,8 +39,8 @@ public abstract class MapBasedForwardIndex<Key, Value> extends AbstractForwardIn
@NotNull
@Override
public InputKeyIterator<Key, Value> getInputKeys(final int inputId) throws IOException {
return new CollectionInputKeyIterator<Key, Value>(myInputsIndex.get(inputId));
public InputDataDiffBuilder<Key, Value> getDiffBuilder(final int inputId) throws IOException {
return new CollectionInputDataDiffBuilder<Key, Value>(inputId, myInputsIndex.get(inputId));
}
@NotNull
@@ -0,0 +1,133 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.util.Comparing;
import com.intellij.util.SystemProperties;
import com.intellij.util.indexing.StorageException;
import gnu.trove.THashMap;
import gnu.trove.TObjectObjectProcedure;
import org.jetbrains.annotations.ApiStatus;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.util.Collections;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
@ApiStatus.Experimental
public class MapInputDataDiffBuilder<Key, Value> extends InputDataDiffBuilder<Key, Value> {
private static final boolean ourDiffUpdateEnabled = SystemProperties.getBooleanProperty("idea.disable.diff.index.update", true);
private final Map<Key, Value> myMap;
public MapInputDataDiffBuilder(int inputId, @Nullable Map<Key, Value> map) {
super(inputId);
myMap = map == null ? Collections.<Key, Value>emptyMap() : map;
}
@Override
public void differentiate(@NotNull Map<Key, Value> newData,
@NotNull KeyValueUpdateProcessor<Key, Value> addProcessor,
@NotNull KeyValueUpdateProcessor<Key, Value> updateProcessor,
@NotNull RemovedKeyProcessor<Key> removeProcessor) throws StorageException {
if (ourDiffUpdateEnabled) {
if (myMap.isEmpty()) {
EmptyInputDataDiffBuilder.processKeys(newData, addProcessor, myInputId);
incrementalAdditions.addAndGet(newData.size());
}
else if (newData.isEmpty()) {
processAllKeysAsDeleted(removeProcessor);
incrementalRemovals.addAndGet(myMap.size());
}
else {
int added = 0;
int removed = 0;
for (Map.Entry<Key, Value> e: myMap.entrySet()) {
final Key key = e.getKey();
final Value newValue = newData.get(key);
if (!Comparing.equal(e.getValue(), newValue) || (newValue == null && !newData.containsKey(key))) {
if (!newData.containsKey(key)) {
removeProcessor.process(key, myInputId);
removed++;
} else {
updateProcessor.process(key, newValue, myInputId);
added++;
removed++;
}
}
}
for (Map.Entry<Key, Value> e : newData.entrySet()) {
final Key key = e.getKey();
if (!myMap.containsKey(key)) {
addProcessor.process(key, e.getValue(), myInputId);
added++;
}
}
incrementalAdditions.addAndGet(added);
incrementalRemovals.addAndGet(removed);
}
int totalRequests = requests.incrementAndGet();
totalRemovals.addAndGet(myMap.size());
totalAdditions.addAndGet(newData.size());
if ((totalRequests & 0xFFF) == 0 && DebugAssertions.DEBUG) {
Logger.getInstance(getClass()).info("Incremental index diff update:" + requests +
", removals:" + totalRemovals + "->" + incrementalRemovals +
", additions:" + totalAdditions + "->" + incrementalAdditions);
}
}
else {
CollectionInputDataDiffBuilder.differentiateWithKeySeq(myMap.keySet(), newData, myInputId, addProcessor, removeProcessor);
}
}
private void processAllKeysAsDeleted(final RemovedKeyProcessor<Key> removeProcessor) throws StorageException {
if (myMap instanceof THashMap) {
final StorageException[] exception = new StorageException[]{null};
((THashMap<Key, Value>)myMap).forEachEntry(new TObjectObjectProcedure<Key, Value>() {
@Override
public boolean execute(Key k, Value v) {
try {
removeProcessor.process(k, myInputId);
}
catch (StorageException e) {
exception[0] = e;
return false;
}
return true;
}
});
if (exception[0] != null) throw exception[0];
}
else {
for (Key key : myMap.keySet()) {
removeProcessor.process(key, myInputId);
}
}
}
private static final AtomicInteger requests = new AtomicInteger();
private static final AtomicInteger totalRemovals = new AtomicInteger();
private static final AtomicInteger totalAdditions = new AtomicInteger();
private static final AtomicInteger incrementalRemovals = new AtomicInteger();
private static final AtomicInteger incrementalAdditions = new AtomicInteger();
}
@@ -1,62 +0,0 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.openapi.util.Comparing;
import org.jetbrains.annotations.Nullable;
import java.util.Collections;
import java.util.Iterator;
import java.util.Map;
public class MapInputKeyIterator<Key, Value> implements ForwardIndex.InputKeyIterator<Key,Value> {
private final Map<Key, Value> myMap;
private Iterator<Map.Entry<Key, Value>> myIterator;
private Value myCurrentValue;
public MapInputKeyIterator(Map<Key, Value> map) {
myMap = map;
}
@Override
public boolean isAssociatedValueEqual(@Nullable Value value) {
return Comparing.equal(myCurrentValue, value);
}
@Override
public boolean hasNext() {
init();
return myIterator.hasNext();
}
@Override
public Key next() {
Map.Entry<Key, Value> entry = myIterator.next();
myCurrentValue = entry.getValue();
return entry.getKey();
}
@Override
public void remove() {
throw new UnsupportedOperationException();
}
private void init() {
if (myIterator == null) {
myIterator = (myMap == null ? Collections.<Key, Value>emptyMap() : myMap).entrySet().iterator();
}
}
}
@@ -51,7 +51,6 @@ public abstract class MapReduceIndex<Key,Value, Input> implements InvertedIndex<
private final DataIndexer<Key, Value, Input> myIndexer;
protected volatile ForwardIndex<Key, Value> myForwardIndex;
private final boolean myUseDiffUpdate;
private final ReentrantReadWriteLock myLock = new ReentrantReadWriteLock();
private volatile boolean myDisposed;
@@ -85,9 +84,6 @@ public abstract class MapReduceIndex<Key,Value, Input> implements InvertedIndex<
myStorage = storage;
myValueExternalizer = extension.getValueExternalizer();
myForwardIndex = forwardIndex;
myUseDiffUpdate = DiffUpdateData.ourDiffUpdateEnabled && (forwardIndex == null ||
myForwardIndex instanceof AbstractForwardIndex &&
!((AbstractForwardIndex)myForwardIndex).hasOnlyKeysData());
}
@NotNull
@@ -222,10 +218,10 @@ public abstract class MapReduceIndex<Key,Value, Input> implements InvertedIndex<
@NotNull
protected UpdateData<Key, Value> calculateUpdateData(final int inputId, @Nullable Input content) {
final Map<Key, Value> data = mapInput(content);
return createUpdateData(data, new ThrowableComputable<ForwardIndex.InputKeyIterator<Key, Value>, IOException>() {
return createUpdateData(data, new ThrowableComputable<InputDataDiffBuilder<Key, Value>, IOException>() {
@Override
public ForwardIndex.InputKeyIterator<Key, Value> compute() throws IOException {
return readInputKeys(inputId);
public InputDataDiffBuilder<Key, Value> compute() throws IOException {
return getKeysDiffBuilder(inputId);
}
}, new ThrowableRunnable<IOException>() {
@Override
@@ -236,16 +232,15 @@ public abstract class MapReduceIndex<Key,Value, Input> implements InvertedIndex<
}
@NotNull
protected ForwardIndex.InputKeyIterator<Key, Value> readInputKeys(int inputId) throws IOException {
return myForwardIndex.getInputKeys(inputId);
protected InputDataDiffBuilder<Key, Value> getKeysDiffBuilder(int inputId) throws IOException {
return myForwardIndex.getDiffBuilder(inputId);
}
@NotNull
protected UpdateData<Key, Value> createUpdateData(Map<Key, Value> data,
ThrowableComputable<ForwardIndex.InputKeyIterator<Key, Value>, IOException> keys,
ThrowableComputable<InputDataDiffBuilder<Key, Value>, IOException> keys,
ThrowableRunnable<IOException> forwardIndexUpdate) {
return myUseDiffUpdate ? new DiffUpdateData<Key, Value>(data, keys, myIndexId, forwardIndexUpdate)
: new SimpleUpdateData<Key, Value>(data, keys, myIndexId, forwardIndexUpdate);
return new UpdateData<Key, Value>(data, keys, myIndexId, forwardIndexUpdate);
}
protected Map<Key, Value> mapInput(Input content) {
@@ -268,8 +263,8 @@ public abstract class MapReduceIndex<Key,Value, Input> implements InvertedIndex<
return myModificationStamp.get();
}
private final UpdateData.RemovedKeyProcessor<Key>
myRemovedKeyProcessor = new UpdateData.RemovedKeyProcessor<Key>() {
private final RemovedKeyProcessor<Key>
myRemovedKeyProcessor = new RemovedKeyProcessor<Key>() {
@Override
public void process(Key key, int inputId) throws StorageException {
myModificationStamp.incrementAndGet();
@@ -277,7 +272,7 @@ public abstract class MapReduceIndex<Key,Value, Input> implements InvertedIndex<
}
};
private final UpdateData.KeyValueUpdateProcessor<Key, Value> myAddedKeyProcessor = new UpdateData.KeyValueUpdateProcessor<Key, Value>() {
private final KeyValueUpdateProcessor<Key, Value> myAddedKeyProcessor = new KeyValueUpdateProcessor<Key, Value>() {
@Override
public void process(Key key, Value value, int inputId) throws StorageException {
myModificationStamp.incrementAndGet();
@@ -285,7 +280,7 @@ public abstract class MapReduceIndex<Key,Value, Input> implements InvertedIndex<
}
};
private final UpdateData.KeyValueUpdateProcessor<Key, Value> myUpdatedKeyProcessor = new UpdateData.KeyValueUpdateProcessor<Key, Value>() {
private final KeyValueUpdateProcessor<Key, Value> myUpdatedKeyProcessor = new KeyValueUpdateProcessor<Key, Value>() {
@Override
public void process(Key key, Value value, int inputId) throws StorageException {
myModificationStamp.incrementAndGet();
@@ -300,7 +295,7 @@ public abstract class MapReduceIndex<Key,Value, Input> implements InvertedIndex<
try {
try {
ValueContainerImpl.ourDebugIndexInfo.set(myIndexId);
updateData.iterateKeys(inputId, myAddedKeyProcessor, myUpdatedKeyProcessor, myRemovedKeyProcessor);
updateData.iterateKeys(myAddedKeyProcessor, myUpdatedKeyProcessor, myRemovedKeyProcessor);
updateData.updateForwardIndex();
}
catch (ProcessCanceledException e) {
@@ -0,0 +1,24 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.util.indexing.StorageException;
import org.jetbrains.annotations.ApiStatus;
@ApiStatus.Experimental
public interface RemovedKeyProcessor<Key> {
void process(Key key, int inputId) throws StorageException;
}
@@ -1,61 +0,0 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed 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 com.intellij.util.indexing.impl;
import com.intellij.openapi.util.ThrowableComputable;
import com.intellij.util.ThrowableRunnable;
import com.intellij.util.indexing.ID;
import com.intellij.util.indexing.StorageException;
import org.jetbrains.annotations.NotNull;
import java.io.IOException;
import java.util.Map;
public class SimpleUpdateData<Key, Value> extends UpdateData<Key,Value> {
public SimpleUpdateData(@NotNull Map<Key, Value> newData,
@NotNull ThrowableComputable<ForwardIndex.InputKeyIterator<Key, Value>, IOException> currentData,
@NotNull ID<Key, Value> indexId,
ThrowableRunnable<IOException> forwardIndexUpdate) {
super(newData, currentData, indexId, forwardIndexUpdate);
}
@Override
public void iterateKeys(int inputId,
KeyValueUpdateProcessor<Key, Value> addProcessor,
KeyValueUpdateProcessor<Key, Value> updateProcessor,
RemovedKeyProcessor<Key> removeProcessor) throws StorageException {
final ForwardIndex.InputKeyIterator<Key, Value> currentData;
try {
currentData = myCurrentData.compute();
}
catch (IOException e) {
throw new StorageException(e);
}
iterateKeys(inputId, addProcessor, removeProcessor, currentData);
}
protected void iterateKeys(int inputId,
KeyValueUpdateProcessor<Key, Value> addProcessor,
RemovedKeyProcessor<Key> removeProcessor, ForwardIndex.InputKeyIterator<Key, Value> currentData)
throws StorageException {
while (currentData.hasNext()) {
removeProcessor.process(currentData.next(), inputId);
}
for (Map.Entry<Key, Value> entry : myNewData.entrySet()) {
addProcessor.process(entry.getKey(), entry.getValue(), inputId);
}
}
}
@@ -27,39 +27,43 @@ import java.io.IOException;
import java.util.Map;
@ApiStatus.Experimental
public abstract class UpdateData<Key, Value> {
public class UpdateData<Key, Value> {
protected final Map<Key, Value> myNewData;
protected final ThrowableComputable<ForwardIndex.InputKeyIterator<Key, Value>, IOException> myCurrentData;
protected final ThrowableComputable<InputDataDiffBuilder<Key, Value>, IOException> myCurrentDataEvaluator;
private final ID<Key, Value> myIndexId;
private final ThrowableRunnable<IOException> myForwardIndexUpdate;
protected UpdateData(@NotNull Map<Key, Value> newData,
@NotNull ThrowableComputable<ForwardIndex.InputKeyIterator<Key, Value>, IOException> currentData,
@NotNull ID<Key, Value> indexId,
@Nullable ThrowableRunnable<IOException> forwardIndexUpdate) {
public UpdateData(@NotNull Map<Key, Value> newData,
@NotNull ThrowableComputable<InputDataDiffBuilder<Key, Value>, IOException> currentDataEvaluator,
@NotNull ID<Key, Value> indexId,
@Nullable ThrowableRunnable<IOException> forwardIndexUpdate) {
myNewData = newData;
myCurrentData = currentData;
myCurrentDataEvaluator = currentDataEvaluator;
myIndexId = indexId;
myForwardIndexUpdate = forwardIndexUpdate;
}
public abstract void iterateKeys(final int inputId,
final KeyValueUpdateProcessor<Key, Value> addProcessor,
final KeyValueUpdateProcessor<Key, Value> updateProcessor,
final RemovedKeyProcessor<Key> removeProcessor) throws StorageException;
public void iterateKeys(KeyValueUpdateProcessor<Key, Value> addProcessor,
KeyValueUpdateProcessor<Key, Value> updateProcessor,
RemovedKeyProcessor<Key> removeProcessor) throws StorageException {
final InputDataDiffBuilder<Key, Value> currentData;
try {
currentData = getCurrentDataEvaluator().compute();
}
catch (IOException e) {
throw new StorageException(e);
}
currentData.differentiate(myNewData, addProcessor, updateProcessor, removeProcessor);
}
public Map<Key, Value> getNewData() {
protected ThrowableComputable<InputDataDiffBuilder<Key, Value>, IOException> getCurrentDataEvaluator() {
return myCurrentDataEvaluator;
}
protected Map<Key, Value> getNewData() {
return myNewData;
}
public interface KeyValueUpdateProcessor<Key, Value> {
void process(Key key, Value value, int inputId) throws StorageException;
}
public interface RemovedKeyProcessor<Key> {
void process(Key key, int inputId) throws StorageException;
}
public ID<Key, Value> getIndexId() {
return myIndexId;
}
@@ -20,10 +20,7 @@ import com.intellij.openapi.progress.ProgressManager;
import com.intellij.openapi.util.Disposer;
import com.intellij.util.Consumer;
import com.intellij.util.indexing.*;
import com.intellij.util.indexing.impl.EmptyInputKeyIterator;
import com.intellij.util.indexing.impl.ForwardIndex;
import com.intellij.util.indexing.impl.MapIndexStorage;
import com.intellij.util.indexing.impl.MapReduceIndex;
import com.intellij.util.indexing.impl.*;
import com.intellij.util.io.DataExternalizer;
import com.intellij.util.io.EnumeratorIntegerDescriptor;
import com.intellij.util.io.KeyDescriptor;
@@ -194,8 +191,8 @@ public class VcsLogFullDetailsIndex<T> implements Disposable {
private static class EmptyForwardIndex<T> implements ForwardIndex<Integer, T> {
@NotNull
@Override
public InputKeyIterator<Integer, T> getInputKeys(int inputId) {
return EmptyInputKeyIterator.getInstance();
public InputDataDiffBuilder<Integer, T> getDiffBuilder(int inputId) {
return new EmptyInputDataDiffBuilder<>(inputId);
}
@Override