devirtualize storage for btree, forcing to use nio.ByteBuffer that allows faster getInt / putInt operations

also rearrange usages of ourLock to have more concurrency during indexing
This commit is contained in:
Maxim.Mossienko
2011-08-29 20:09:12 +04:00
parent 48346ff668
commit d2c0d64cf5
10 changed files with 395 additions and 489 deletions
@@ -1,21 +0,0 @@
package com.intellij.util.io;
import com.intellij.openapi.Forceable;
import java.io.Closeable;
interface ISimpleStorage extends Closeable, Forceable {
void put(int index, byte value);
byte get(int index);
void putInt(int index, int value);
int getInt(int index);
long length();
void put(int pos, byte[] buf, int offset, int length);
void get(int pos, byte[] buf, int offset, int length);
void putLong(int pos, long value);
long getLong(int pos);
}
@@ -1,12 +1,12 @@
package com.intellij.util.io;
import com.intellij.openapi.util.SystemInfo;
import com.intellij.openapi.util.io.FileUtil;
import gnu.trove.TIntIntHashMap;
import org.jetbrains.annotations.Nullable;
import java.io.File;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.Arrays;
/**
@@ -36,8 +36,8 @@ class IntToIntBtree {
private final byte[] buffer;
private boolean isLarge = true;
private final ISimpleStorage storage;
private final boolean offloadToSiblingsBeforeSplit = hasCachedMappings;
private final ResizeableMappedFile storage;
private final boolean offloadToSiblingsBeforeSplit = false;
private boolean indexNodeIsHashTable = true;
final int metaDataLeafPageLength;
final int hashPageCapacity;
@@ -50,8 +50,18 @@ class IntToIntBtree {
pageSize = _pageSize;
buffer = new byte[_pageSize];
if (initial) {
FileUtil.delete(file);
}
storage = new ResizeableMappedFile(file, pageSize, PersistentEnumeratorBase.ourLock, 1024 * 1024, true);
root = new BtreeIndexNodeView(this);
root.setAddress(0);
if (initial) {
nextPage(); // allocate root
root.setAddress(0);
root.setIndexLeaf(true);
}
int i = (pageSize - BtreePage.RESERVED_META_PAGE_LEN) / BtreeIndexNodeView.INTERIOR_SIZE - 1;
assert i < Short.MAX_VALUE && i % 2 == 0;
@@ -82,21 +92,6 @@ class IntToIntBtree {
assert i > 0 && i % 2 == 0;
maxLeafNodesInHash = (short) i;
if (initial) {
FileUtil.delete(file);
}
storage = new MappedFileSimpleStorage(file, pageSize, (SystemInfo.is64Bit ? 10:1) *1024 * 1024);
if (initial) {
nextPage(); // allocate root
}
((BtreePage)root).load();
if (initial) {
root.setIndexLeaf(true);
root.sync();
}
if (hasCachedMappings) {
myCachedMappings = new TIntIntHashMap(myCachedMappingsSize = 4 * maxLeafNodes);
} else {
@@ -174,18 +169,14 @@ class IntToIntBtree {
}
}
public int remove(int key) {
// TODO
//BtreeIndexNodeView currentIndexNode = new BtreeIndexNodeView(this);
//currentIndexNode.setAddress(root.address);
//int index = currentIndexNode.locate(key, false);
myAssert(BtreeIndexNodeView.haveDeleteState);
throw new UnsupportedOperationException("Remove does not work yet "+key);
}
void setRootAddress(int newRootAddress) {
root.setAddress(newRootAddress);
}
//public int remove(int key) {
// // TODO
// BtreeIndexNodeView currentIndexNode = new BtreeIndexNodeView(this);
// currentIndexNode.setAddress(root.address);
// int index = currentIndexNode.locate(key, false);
// myAssert(BtreeIndexNodeView.haveDeleteState);
// throw new UnsupportedOperationException("Remove does not work yet "+key);
//}
void dumpStatistics() {
int leafPages = height == 3 ? pagesCount - (1 + root.getChildrenCount() + 1):height == 2 ? pagesCount - 1:1;
@@ -228,8 +219,10 @@ class IntToIntBtree {
static final int RESERVED_META_PAGE_LEN = 8;
protected final IntToIntBtree btree;
protected int address;
protected int address = -1;
private short myChildrenCount;
protected int myAddressInBuffer;
protected ByteBuffer myBuffer;
public BtreePage(IntToIntBtree btree) {
this.btree = btree;
@@ -239,63 +232,80 @@ class IntToIntBtree {
void setAddress(int _address) {
if (doSanityCheck) myAssert(_address % btree.pageSize == 0);
address = _address;
syncWithStore();
}
protected void syncWithStore() {
myChildrenCount = -1;
PagedFileStorage pagedFileStorage = btree.storage.getPagedFileStorage();
myAddressInBuffer = pagedFileStorage.getOffsetInPage(address);
myBuffer = pagedFileStorage.getByteBuffer(address);
}
protected final boolean getFlag(int mask) {
return (btree.storage.get(address) & mask) == mask;
return (myBuffer.get(myAddressInBuffer) & mask) == mask;
}
protected final void setFlag(int mask, boolean flag) {
byte b = btree.storage.get(address);
byte b = myBuffer.get(myAddressInBuffer);
if (flag) b |= mask;
else b &= ~mask;
btree.storage.put(address, b);
myBuffer.put(myAddressInBuffer, b);
}
protected final short getChildrenCount() {
if (myChildrenCount == -1) {
myChildrenCount = (short)(((btree.storage.get(address + 1) & 0xFF) << 8) + (btree.storage.get(address + 2) & 0xFF));
myChildrenCount = myBuffer.getShort(myAddressInBuffer + 1);
}
return myChildrenCount;
}
protected final void setChildrenCount(short value) {
myChildrenCount = value;
btree.storage.put(address + 1, (byte)((value >> 8) & 0xFF));
btree.storage.put(address + 2, (byte)(value & 0xFF));
myBuffer.putShort(myAddressInBuffer + 1, value);
}
protected final void setNextPage(int nextPage) {
putInt(address + 3, nextPage);
putInt(3, nextPage);
}
// TODO: use it
protected final int getNextPage() {
return getInt(address + 3);
return getInt(3);
}
protected final int getInt(int address) {
return btree.storage.getInt(address);
return myBuffer.getInt(myAddressInBuffer + address);
}
protected final void putInt(int offset, int value) {
btree.storage.putInt(offset, value);
myBuffer.putInt(myAddressInBuffer + offset, value);
}
protected final void getBytes(int address, byte[] dst, int offset, int length) {
btree.storage.get(address, dst, offset, length);
protected final void getBytes(int address, byte[] dst, int length) {
myBuffer.position(address + myAddressInBuffer);
myBuffer.get(dst, 0, length);
}
private void load() {
protected final void putBytes(int address, byte[] src, int length) {
myBuffer.position(address + myAddressInBuffer);
myBuffer.put(src, 0, length);
}
protected final void putBytes(int address, byte[] src, int offset, int length) {
btree.storage.put(address, src, offset, length);
}
void sync() {
}
//protected final ByteBuffer getBytes(int address, int length) {
// ByteBuffer duplicate = myBuffer.duplicate();
//
// int newPosition = address + myAddressInBuffer;
// duplicate.position(newPosition);
// duplicate.limit(newPosition + length);
// return duplicate;
//}
//
//protected final void putBytes(int address, ByteBuffer buffer) {
// myBuffer.position(address + myAddressInBuffer);
// myBuffer.put(buffer);
//}
}
// Leaf index node
@@ -318,8 +328,8 @@ class IntToIntBtree {
}
@Override
void setAddress(int _address) {
super.setAddress(_address);
protected void syncWithStore() {
super.syncWithStore();
isIndexLeafSet = false;
isHashedLeafSet = false;
}
@@ -375,14 +385,14 @@ class IntToIntBtree {
myAssert(i < childrenCount || (!isIndexLeaf() && i == childrenCount));
metaPageLen = RESERVED_META_PAGE_LEN;
}
myAssert(offset + 4 <= address + btree.pageSize);
myAssert(offset >= address + metaPageLen);
myAssert(offset + 4 <= btree.pageSize);
myAssert(offset >= metaPageLen);
}
putInt(offset, value);
}
private final int indexToOffset(int i) {
return address + i * INTERIOR_SIZE + (isHashedLeaf() ? btree.metaDataLeafPageLength:RESERVED_META_PAGE_LEN);
return i * INTERIOR_SIZE + (isHashedLeaf() ? btree.metaDataLeafPageLength:RESERVED_META_PAGE_LEN);
}
private final int keyAt(int i) {
@@ -405,8 +415,8 @@ class IntToIntBtree {
myAssert(i < getChildrenCount());
metaPageLen = RESERVED_META_PAGE_LEN;
}
myAssert(offset + 4 <= address + btree.pageSize);
myAssert(offset >= address + metaPageLen);
myAssert(offset + 4 <= btree.pageSize);
myAssert(offset >= metaPageLen);
}
putInt(offset, value);
}
@@ -458,7 +468,7 @@ class IntToIntBtree {
int[] keys = new int[childrenCount];
if (isHashedLeaf()) {
getBytes(indexToOffset(0), btree.buffer, 0, btree.pageSize - btree.metaDataLeafPageLength);
getBytes(indexToOffset(0), btree.buffer, btree.pageSize - btree.metaDataLeafPageLength);
int keyNumber = 0;
for(int i = 0; i < btree.hashPageCapacity; ++i) {
@@ -483,7 +493,13 @@ class IntToIntBtree {
HashLeafData(BtreeIndexNodeView _nodeView, int recordCount) {
nodeView = _nodeView;
final IntToIntBtree btree = _nodeView.btree;
nodeView.getBytes(nodeView.indexToOffset(0), btree.buffer, 0, btree.pageSize - btree.metaDataLeafPageLength);
nodeView.getBytes(nodeView.indexToOffset(0), btree.buffer, btree.pageSize - btree.metaDataLeafPageLength);
// TODO: using bytebuffers should more efficient for copying data too!
//int length = btree.pageSize - btree.metaDataLeafPageLength;
//final int offset = nodeView.indexToOffset(0);
//ByteBuffer buffer = nodeView.getBytes(offset, length);
//byte[] buf = new byte[btree.pageSize];
//buffer.get(buf, 0, length);
keys = new int[recordCount];
values = new TIntIntHashMap(recordCount);
int keyNumber = 0;
@@ -491,8 +507,13 @@ class IntToIntBtree {
for(int i = 0; i < btree.hashPageCapacity; ++i) {
if (nodeView.hashGetState(i) == HASH_FULL) {
int key = Bits.getInt(btree.buffer, i * INTERIOR_SIZE + KEY_OFFSET);
//int key2 = buffer.getInt(offset + i * INTERIOR_SIZE + KEY_OFFSET);
//myAssert(key == key2);
keys[keyNumber++] = key;
values.put(key, Bits.getInt(btree.buffer, i * INTERIOR_SIZE));
int value = Bits.getInt(btree.buffer, i * INTERIOR_SIZE);
//int value2 = buffer.getInt(offset + i * INTERIOR_SIZE);
//myAssert(value == value2);
values.put(key, value);
}
}
@@ -538,6 +559,9 @@ class IntToIntBtree {
BtreeIndexNodeView newIndexNode = new BtreeIndexNodeView(btree);
newIndexNode.setAddress(btree.nextPage());
syncWithStore(); // next page can cause ByteBuffer to be invalidated!
if (parent != null) parent.syncWithStore();
btree.root.syncWithStore();
newIndexNode.setIndexLeaf(indexLeaf);
@@ -554,7 +578,7 @@ class IntToIntBtree {
boolean defaultSplit = true;
//if (keys[keys.length - 1] < newValue && btree.height <= 3) { // optimization for adding element to last block
// btree.root.setAddress(btree.root.address);
// btree.root.syncWithStore();
// if (btree.height == 2 && btree.root.search(keys[0]) == btree.root.getChildrenCount() - 1) {
// defaultSplit = false;
// } else if (btree.height == 3 &&
@@ -609,8 +633,8 @@ class IntToIntBtree {
if (btree.isLarge) {
final int bytesToMove = recordCountInNewNode * INTERIOR_SIZE;
getBytes(indexToOffset(maxIndex), btree.buffer, 0, bytesToMove);
newIndexNode.putBytes(newIndexNode.indexToOffset(0), btree.buffer, 0, bytesToMove);
getBytes(indexToOffset(maxIndex), btree.buffer, bytesToMove);
newIndexNode.putBytes(newIndexNode.indexToOffset(0), btree.buffer, bytesToMove);
} else {
for(int i = 0; i < recordCountInNewNode; ++i) {
newIndexNode.setAddressAt(i, addressAt(i + maxIndex));
@@ -653,18 +677,22 @@ class IntToIntBtree {
if (doSanityCheck) {
btree.root.dump("Splitting root:"+medianKey);
}
int newRootAddress = btree.nextPage();
newIndexNode.syncWithStore();
syncWithStore();
if (doSanityCheck) {
System.out.println("Pages:"+btree.pagesCount+", elements:"+btree.count + ", average:" + (btree.height + 1));
}
btree.setRootAddress(newRootAddress);
btree.root.setAddress(newRootAddress);
parentAddress = newRootAddress;
((BtreePage)btree.root).load();
btree.root.setChildrenCount((short)1);
btree.root.setKeyAt(0, medianKey);
btree.root.setAddressAt(0, -address);
btree.root.setAddressAt(1, -newIndexNode.address);
btree.root.sync();
if (doSanityCheck) {
btree.root.dump("New root");
@@ -673,9 +701,6 @@ class IntToIntBtree {
}
}
sync();
newIndexNode.sync();
return parentAddress;
}
@@ -725,9 +750,6 @@ class IntToIntBtree {
}
if (!isFull()) {
sync();
parent.sync();
if (doSanityCheck) {
dump("old node after split:");
parent.dump("Parent node after split");
@@ -812,8 +834,8 @@ class IntToIntBtree {
if (btree.isLarge) {
final int bytesToMove = indexOfLastChildToMove * INTERIOR_SIZE;
getBytes(indexToOffset(toMove), btree.buffer, 0, bytesToMove);
putBytes(indexToOffset(0), btree.buffer, 0, bytesToMove);
getBytes(indexToOffset(toMove), btree.buffer, bytesToMove);
putBytes(indexToOffset(0), btree.buffer, bytesToMove);
}
else {
for (int i = 0; i < indexOfLastChildToMove; ++i) {
@@ -832,9 +854,6 @@ class IntToIntBtree {
}
if (!isFull()) {
sync();
parent.sync();
if (doSanityCheck) {
dump("old node after split:");
parent.dump("Parent node after split");
@@ -937,7 +956,7 @@ class IntToIntBtree {
hashSetState(index, HASH_FULL);
setAddressAt(index, newValueId);
setChildrenCount((short)(recordCount + 1));
sync();
return;
}
}
@@ -953,8 +972,8 @@ class IntToIntBtree {
if (indexLeaf) {
if (btree.isLarge && itemsToMove > LARGE_MOVE_THRESHOLD) {
final int bytesToMove = itemsToMove * INTERIOR_SIZE;
getBytes(indexToOffset(index), btree.buffer, 0, bytesToMove);
putBytes(indexToOffset(index + 1), btree.buffer, 0, bytesToMove);
getBytes(indexToOffset(index), btree.buffer, bytesToMove);
putBytes(indexToOffset(index + 1), btree.buffer, bytesToMove);
} else {
for(int i = recordCount - 1; i >= index; --i) {
setKeyAt(i + 1, keyAt(i));
@@ -971,8 +990,8 @@ class IntToIntBtree {
int elementsAfterIndex = recordCount - index - 1;
if (elementsAfterIndex > 0) {
int bytesToMove = elementsAfterIndex * INTERIOR_SIZE;
getBytes(indexToOffset(index + 1), btree.buffer, 0, bytesToMove);
putBytes(indexToOffset(index + 2), btree.buffer, 0, bytesToMove);
getBytes(indexToOffset(index + 1), btree.buffer, bytesToMove);
putBytes(indexToOffset(index + 2), btree.buffer, bytesToMove);
}
} else {
for(int i = recordCount - 1; i > index; --i) {
@@ -991,8 +1010,6 @@ class IntToIntBtree {
if (index > 0) myAssert(keyAt(index - 1) < keyAt(index));
if (index < recordCount) myAssert(keyAt(index) < keyAt(index + 1));
}
sync();
}
private int hashIndex(int value) {
@@ -1096,7 +1113,7 @@ class IntToIntBtree {
static final int STATE_MASK_WITHOUT_DELETE = 0x1;
private final int hashGetState(int index) {
byte b = btree.storage.get(hashOccupiedStatusByteOffset(index));
byte b = myBuffer.get(hashOccupiedStatusByteOffset(index));
if (haveDeleteState) {
return ((b & 0xFF) >> hashOccupiedStatusShift(index)) & STATE_MASK;
} else {
@@ -1114,9 +1131,9 @@ class IntToIntBtree {
private final int hashOccupiedStatusByteOffset(int index) {
if (haveDeleteState) {
return address + BtreePage.RESERVED_META_PAGE_LEN + (index >> 2);
return myAddressInBuffer + BtreePage.RESERVED_META_PAGE_LEN + (index >> 2);
} else {
return address + BtreePage.RESERVED_META_PAGE_LEN + (index >> 3);
return myAddressInBuffer + BtreePage.RESERVED_META_PAGE_LEN + (index >> 3);
}
}
@@ -1129,7 +1146,7 @@ class IntToIntBtree {
}
int hashOccupiedStatusOffset = hashOccupiedStatusByteOffset(index);
byte b = btree.storage.get(hashOccupiedStatusOffset);
byte b = myBuffer.get(hashOccupiedStatusOffset);
int shift = hashOccupiedStatusShift(index);
if (haveDeleteState) {
@@ -1137,7 +1154,7 @@ class IntToIntBtree {
} else {
b = (byte)(((b & 0xFF) & ~(STATE_MASK_WITHOUT_DELETE << shift)) | (value << shift));
}
btree.storage.put(hashOccupiedStatusOffset, b);
myBuffer.put(hashOccupiedStatusOffset, b);
if (doSanityCheck) myAssert(hashGetState(index) == value);
}
@@ -1,80 +0,0 @@
package com.intellij.util.io;
import com.intellij.openapi.util.io.FileUtil;
import java.io.File;
import java.io.IOException;
final class MappedFileSimpleStorage implements ISimpleStorage {
private final ResizeableMappedFile storage;
public MappedFileSimpleStorage(File file, int initialSize) throws IOException {
FileUtil.createIfDoesntExist(file);
storage = new ResizeableMappedFile(file, initialSize, PersistentEnumeratorBase.ourLock);
}
public MappedFileSimpleStorage(File file, int initialSize, int pageSize) throws IOException {
FileUtil.createIfDoesntExist(file);
storage = new ResizeableMappedFile(file, initialSize, PersistentEnumeratorBase.ourLock, pageSize);
}
@Override
public void put(int index, byte value) {
storage.put(index, value);
}
@Override
public byte get(int index) {
return storage.get(index);
}
@Override
public void putInt(int index, int value) {
storage.putInt(index, value);
}
@Override
public int getInt(int index) {
return storage.getInt(index);
}
@Override
public long length() {
return storage.length();
}
@Override
public void put(int pos, byte[] buf, int offset, int length) {
storage.put(pos, buf, offset, length);
}
@Override
public void get(int pos, byte[] buf, int offset, int length) {
storage.get(pos, buf, offset, length);
}
@Override
public void putLong(int pos, long value) {
storage.putLong(pos, value);
}
@Override
public long getLong(int pos) {
return storage.getLong(pos);
}
@Override
public void close() throws IOException {
storage.close();
}
@Override
public boolean isDirty() {
return storage.isDirty();
}
@Override
public void force() {
storage.force();
}
}
@@ -29,6 +29,7 @@ import java.io.RandomAccessFile;
import java.nio.ByteBuffer;
import java.nio.MappedByteBuffer;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
/**
* @author max
@@ -40,9 +41,9 @@ public class PagedFileStorage implements Forceable {
final static int DEFAULT_BUFFER_SIZE;
private final static int UPPER_LIMIT;
public static final int LOWER_LIMIT_IN_MEGABYTES = 100;
private static final int LOWER_LIMIT_IN_MEGABYTES = 100;
private final static int LOWER_LIMIT = LOWER_LIMIT_IN_MEGABYTES * MEGABYTE;
public static final int UNKNOWN_PAGE = -1;
private static final int UNKNOWN_PAGE = -1;
static {
String maxPagedStorageCacheProperty = System.getProperty("idea.max.paged.storage.cache");
@@ -62,8 +63,13 @@ public class PagedFileStorage implements Forceable {
private MappedBufferWrapper myLastBuffer2;
private int myLastChangeCount;
private int myLastChangeCount2;
private final int myStorageIndex;
private static final int MAX_PAGES_COUNT = 0xFFFF;
public static class StorageLock {
private static final int FILE_INDEX_MASK = 0xFFFF0000;
private static final int FILE_INDEX_SHIFT = 16;
private final boolean checkThreadAccess;
public StorageLock() {
@@ -74,28 +80,88 @@ public class PagedFileStorage implements Forceable {
this.checkThreadAccess = checkThreadAccess;
}
final BuffersCache myBuffersCache = new BuffersCache();
private final BuffersCache myBuffersCache = new BuffersCache();
private final ConcurrentHashMap<Integer, PagedFileStorage> myIndex2Storage = new ConcurrentHashMap<Integer, PagedFileStorage>();
private int registerPagedFileStorage(PagedFileStorage storage) {
int value;
while(myIndex2Storage.putIfAbsent(value = (myIndex2Storage.size() << FILE_INDEX_SHIFT), storage) != null) ;
return value;
}
private PagedFileStorage getRegisteredPagedFileStorageByIndex(int index) {
return myIndex2Storage.get(index);
}
private class BuffersCache extends MyCache {
private class BuffersCache {
private int changeCount;
private final LinkedHashMap<Integer, MappedBufferWrapper> myMap;
private long mySizeLimit;
private long mySize;
public BuffersCache() {
super(UPPER_LIMIT);
private BuffersCache() {
mySizeLimit = UPPER_LIMIT;
myMap = new LinkedHashMap<Integer, MappedBufferWrapper>(10) {
@Override
protected boolean removeEldestEntry(Map.Entry<Integer, MappedBufferWrapper> eldest) {
return mySize > mySizeLimit;
}
@Nullable
@Override
public MappedBufferWrapper remove(Object key) {
// this method can be called after removeEldestEntry
MappedBufferWrapper wrapper = super.remove(key);
if (wrapper != null) {
mySize -= wrapper.myLength;
wrapper.dispose();
}
return wrapper;
}
};
}
private MappedBufferWrapper get(Integer key) {
MappedBufferWrapper wrapper = myMap.get(key);
if (wrapper != null) {
return wrapper;
}
long started = IOStatistics.DEBUG ? System.currentTimeMillis() : 0;
wrapper = createValue(key);
mySize += wrapper.myLength;
if (IOStatistics.DEBUG) {
long finished = System.currentTimeMillis();
if (finished - started > IOStatistics.MIN_IO_TIME_TO_REPORT) {
IOStatistics.dump(
"Mapping " + wrapper.myLength + " from " + wrapper.myPosition + " file:" + wrapper.myFile + " for " + (finished - started));
}
}
myMap.put(key, wrapper);
ensureSize(mySizeLimit);
return wrapper;
}
private void ensureSize(long sizeLimit) {
while (mySize > sizeLimit) {
// we still have to drop something
myMap.doRemoveEldestEntry();
}
}
@NotNull
public MappedBufferWrapper createValue(PageKey key) {
if (checkThreadAccess && !Thread.holdsLock(StorageLock.this)) {
throw new IllegalStateException("Must hold StorageLock lock to access PagedFileStorage");
}
int off = key.page * key.owner.myPageSize;
if (off > key.owner.length()) {
throw new IndexOutOfBoundsException("off=" + off + " key.owner.length()=" + key.owner.length());
private MappedBufferWrapper createValue(Integer key) {
checkThreadAccess();
PagedFileStorage owner = getRegisteredPagedFileStorageByIndex(key & FILE_INDEX_MASK);
int off = (key & MAX_PAGES_COUNT) * owner.myPageSize;
if (off > owner.length()) {
throw new IndexOutOfBoundsException("off=" + off + " key.owner.length()=" + owner.length());
}
++changeCount;
ReadWriteMappedBufferWrapper wrapper =
new ReadWriteMappedBufferWrapper(key.owner.myFile, off, Math.min((int)(key.owner.length() - off), key.owner.myPageSize));
new ReadWriteMappedBufferWrapper(owner.myFile, off, Math.min((int)(owner.length() - off), owner.myPageSize));
IOException oome = null;
while (true) {
try {
@@ -112,17 +178,18 @@ public class PagedFileStorage implements Forceable {
if (e.getCause() instanceof OutOfMemoryError) {
oome = e;
if (mySizeLimit > LOWER_LIMIT) {
mySizeLimit -= key.owner.myPageSize;
mySizeLimit -= owner.myPageSize;
}
long newSize = getSize() - key.owner.myPageSize;
long newSize = mySize - owner.myPageSize;
if (newSize >= 0) {
ensureSize(newSize);
continue; // next try
}
else {
throw new MappingFailedException("Cannot recover from OOME in memory mapping: -Xmx=" + Runtime.getRuntime().maxMemory() / MEGABYTE + "MB " +
"new size limit: " + mySizeLimit / MEGABYTE + "MB " +
"trying to allocate " + wrapper.myLength + " block", e);
throw new MappingFailedException(
"Cannot recover from OOME in memory mapping: -Xmx=" + Runtime.getRuntime().maxMemory() / MEGABYTE + "MB " +
"new size limit: " + mySizeLimit / MEGABYTE + "MB " +
"trying to allocate " + wrapper.myLength + " block", e);
}
}
throw new MappingFailedException("Cannot map buffer", e);
@@ -130,56 +197,72 @@ public class PagedFileStorage implements Forceable {
}
}
public void onDropFromCache(PageKey key, MappedBufferWrapper buf) {
buf.dispose();
private void checkThreadAccess() {
if (checkThreadAccess && !Thread.holdsLock(StorageLock.this)) {
throw new IllegalStateException("Must hold StorageLock lock to access PagedFileStorage");
}
}
private @Nullable Map<Integer, MappedBufferWrapper> getBuffersOrderedForOwner(int index) {
checkThreadAccess();
Map<Integer, MappedBufferWrapper> mineBuffers = null;
for (Map.Entry<Integer, MappedBufferWrapper> entry : myMap.entrySet()) {
if ((entry.getKey() & FILE_INDEX_MASK) == index) {
if (mineBuffers == null) {
mineBuffers = new TreeMap<Integer, MappedBufferWrapper>(new Comparator<Integer>() {
@Override
public int compare(Integer o1, Integer o2) {
return o1 - o2;
}
});
}
mineBuffers.put(entry.getKey(), entry.getValue());
}
}
return mineBuffers;
}
private void unmapBuffersForOwner(int index) {
final Map<Integer, MappedBufferWrapper> buffers = getBuffersOrderedForOwner(index);
if (buffers != null) {
for (Integer key : buffers.keySet()) {
myMap.remove(key);
}
}
}
private void flushBuffersForOwner(int index) {
Map<Integer, MappedBufferWrapper> buffers = getBuffersOrderedForOwner(index);
if (buffers != null) {
for(MappedBufferWrapper buffer:buffers.values()) {
buffer.flush();
}
}
}
}
}
private static class PageKey {
private final PagedFileStorage owner;
private final int page;
public PageKey(PagedFileStorage owner, int page) {
this.owner = owner;
this.page = page;
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof PageKey)) return false;
PageKey pageKey = (PageKey)o;
if (!owner.equals(pageKey.owner)) return false;
if (page != pageKey.page) return false;
return true;
}
@Override
public int hashCode() {
return 31 * owner.hashCode() + page;
}
}
private final byte[] myTypedIOBuffer = new byte[8];
private final byte[] myTypedIOBuffer;
private boolean isDirty = false;
private final File myFile;
protected long mySize = -1;
protected final int myPageSize;
protected final boolean myValuesAreBufferAligned;
@NonNls private static final String RW = "rw";
public PagedFileStorage(File file, StorageLock lock, int pageSize) throws IOException {
public PagedFileStorage(File file, StorageLock lock, int pageSize, boolean valuesAreBufferAligned) throws IOException {
myFile = file;
myLock = lock;
myPageSize = Math.max(pageSize, Page.PAGE_SIZE);
myValuesAreBufferAligned = valuesAreBufferAligned;
myStorageIndex = lock.registerPagedFileStorage(this);
myTypedIOBuffer = valuesAreBufferAligned ? null:new byte[8];
}
public PagedFileStorage(File file, StorageLock lock) throws IOException {
this(file, lock, DEFAULT_BUFFER_SIZE);
this(file, lock, DEFAULT_BUFFER_SIZE, false);
}
public File getFile() {
@@ -187,34 +270,89 @@ public class PagedFileStorage implements Forceable {
}
public void putInt(int addr, int value) {
Bits.putInt(myTypedIOBuffer, 0, value);
put(addr, myTypedIOBuffer, 0, 4);
if (myValuesAreBufferAligned) {
isDirty = true;
int page = addr / myPageSize;
int page_offset = addr % myPageSize;
getBuffer(page).putInt(page_offset, value);
} else {
Bits.putInt(myTypedIOBuffer, 0, value);
put(addr, myTypedIOBuffer, 0, 4);
}
}
public int getInt(int addr) {
get(addr, myTypedIOBuffer, 0, 4);
return Bits.getInt(myTypedIOBuffer, 0);
if (myValuesAreBufferAligned) {
int page = addr / myPageSize;
int page_offset = addr % myPageSize;
return getBuffer(page).getInt(page_offset);
} else {
get(addr, myTypedIOBuffer, 0, 4);
return Bits.getInt(myTypedIOBuffer, 0);
}
}
public final void putShort(int addr, short value) {
if (myValuesAreBufferAligned) {
isDirty = true;
int page = addr / myPageSize;
int page_offset = addr % myPageSize;
getBuffer(page).putShort(page_offset, value);
} else {
Bits.putShort(myTypedIOBuffer, 0, value);
put(addr, myTypedIOBuffer, 0, 2);
}
}
int getOffsetInPage(int addr) {
return addr % myPageSize;
}
ByteBuffer getByteBuffer(int address) {
return getBuffer(address / myPageSize);
}
public final short getShort(int addr) {
if (myValuesAreBufferAligned) {
int page = addr / myPageSize;
int page_offset = addr % myPageSize;
return getBuffer(page).getShort(page_offset);
} else {
get(addr, myTypedIOBuffer, 0, 2);
return Bits.getShort(myTypedIOBuffer, 0);
}
}
public void putLong(int addr, long value) {
Bits.putLong(myTypedIOBuffer, 0, value);
put(addr, myTypedIOBuffer, 0, 8);
if (myValuesAreBufferAligned) {
isDirty = true;
int page = addr / myPageSize;
int page_offset = addr % myPageSize;
getBuffer(page).putLong(page_offset, value);
} else {
Bits.putLong(myTypedIOBuffer, 0, value);
put(addr, myTypedIOBuffer, 0, 8);
}
}
@SuppressWarnings({"UnusedDeclaration"})
public void putByte(final int addr, final byte b) {
myTypedIOBuffer[0] = b;
put(addr, myTypedIOBuffer, 0, 1);
put(addr, b);
}
public byte getByte(int addr) {
get(addr, myTypedIOBuffer, 0, 1);
return myTypedIOBuffer[0];
return get(addr);
}
public long getLong(int addr) {
get(addr, myTypedIOBuffer, 0, 8);
return Bits.getLong(myTypedIOBuffer, 0);
if (myValuesAreBufferAligned) {
int page = addr / myPageSize;
int page_offset = addr % myPageSize;
return getBuffer(page).getLong(page_offset);
} else {
get(addr, myTypedIOBuffer, 0, 8);
return Bits.getLong(myTypedIOBuffer, 0);
}
}
public byte get(int index) {
@@ -297,13 +435,7 @@ public class PagedFileStorage implements Forceable {
}
private void unmapAll() {
final Map<PageKey, MappedBufferWrapper> mineBuffers = getMineBuffersOrdered();
if (mineBuffers != null) {
for (PageKey key : mineBuffers.keySet()) {
myLock.myBuffersCache.remove(key);
}
}
myLock.myBuffersCache.unmapBuffersForOwner(myStorageIndex);
myLastPage = UNKNOWN_PAGE;
myLastPage2 = UNKNOWN_PAGE;
@@ -378,7 +510,8 @@ public class PagedFileStorage implements Forceable {
}
try {
MappedBufferWrapper mappedBufferWrapper = myLock.myBuffersCache.get(new PageKey(this, page));
assert page <= MAX_PAGES_COUNT;
MappedBufferWrapper mappedBufferWrapper = myLock.myBuffersCache.get(myStorageIndex | page);
MappedByteBuffer buf = mappedBufferWrapper.buf();
if (myLastPage != page) {
@@ -402,13 +535,8 @@ public class PagedFileStorage implements Forceable {
public void force() {
long started = IOStatistics.DEBUG ? System.currentTimeMillis():0;
Map<PageKey, MappedBufferWrapper> mineBuffers = getMineBuffersOrdered();
if (mineBuffers != null) {
for(MappedBufferWrapper buffer:mineBuffers.values()) {
buffer.flush();
}
}
myLock.myBuffersCache.flushBuffersForOwner(myStorageIndex);
isDirty = false;
if (IOStatistics.DEBUG) {
long finished = System.currentTimeMillis();
@@ -418,98 +546,7 @@ public class PagedFileStorage implements Forceable {
}
}
private @Nullable
Map<PageKey, MappedBufferWrapper> getMineBuffersOrdered() {
Map<PageKey, MappedBufferWrapper> mineBuffers = null;
for (Map.Entry<PageKey,MappedBufferWrapper> entry : myLock.myBuffersCache.entrySet()) {
if (entry.getKey().owner == this) {
if (mineBuffers == null) mineBuffers = new TreeMap<PageKey, MappedBufferWrapper>(new Comparator<PageKey>() {
@Override
public int compare(PageKey o1, PageKey o2) {
return o1.page - o2.page;
}
});
mineBuffers.put(entry.getKey(), entry.getValue());
}
}
return mineBuffers;
}
public boolean isDirty() {
return isDirty;
}
private static abstract class MyCache {
private final LinkedHashMap<PageKey, MappedBufferWrapper> myMap;
protected long mySizeLimit;
private long mySize;
protected MyCache(long sizeLimit) {
mySizeLimit = sizeLimit;
myMap = new LinkedHashMap<PageKey, MappedBufferWrapper>(10) {
@Override
protected boolean removeEldestEntry(Map.Entry<PageKey, MappedBufferWrapper> eldest) {
return mySize > mySizeLimit;
}
@Nullable
@Override
public MappedBufferWrapper remove(Object key) {
// this method can be called after removeEldestEntry
MappedBufferWrapper wrapper = super.remove(key);
if (wrapper != null) {
mySize -= wrapper.myLength;
onDropFromCache((PageKey)key, wrapper);
}
return wrapper;
}
};
}
public MappedBufferWrapper get(PageKey key) {
MappedBufferWrapper wrapper = myMap.get(key);
if (wrapper != null) {
return wrapper;
}
long started = IOStatistics.DEBUG ? System.currentTimeMillis() : 0;
wrapper = createValue(key);
mySize += wrapper.myLength;
if (IOStatistics.DEBUG) {
long finished = System.currentTimeMillis();
if (finished - started > IOStatistics.MIN_IO_TIME_TO_REPORT) {
IOStatistics.dump("Mapping " + wrapper.myLength + " from " + wrapper.myPosition + " file:"+wrapper.myFile + " for "+(finished - started));
}
}
myMap.put(key, wrapper);
ensureSize(mySizeLimit);
return wrapper;
}
protected void ensureSize(long sizeLimit) {
while (mySize > sizeLimit) {
// we still have to drop something
myMap.doRemoveEldestEntry();
}
}
public long getSize() {
return mySize;
}
public Set<Map.Entry<PageKey, MappedBufferWrapper>> entrySet() {
return myMap.entrySet();
}
public void remove(PageKey key) {
myMap.remove(key);
}
protected abstract MappedBufferWrapper createValue(PageKey key);
protected abstract void onDropFromCache(PageKey key, MappedBufferWrapper wrapper);
}
}
@@ -42,6 +42,10 @@ public class PersistentBTreeEnumerator<Data> extends PersistentEnumeratorBase<Da
}
private static final int RECORD_SIZE = 4;
private static final int VALUE_PAGE_SIZE = 1024 * 1024;
static {
assert VALUE_PAGE_SIZE % PAGE_SIZE == 0:"Page size should be divisor of " + VALUE_PAGE_SIZE;
}
private int myLogicalFileLength;
private int myDataPageStart;
@@ -65,7 +69,7 @@ public class PersistentBTreeEnumerator<Data> extends PersistentEnumeratorBase<Da
private static final int KEY_SHIFT = 1;
public PersistentBTreeEnumerator(File file, KeyDescriptor<Data> dataDescriptor, int initialSize) throws IOException {
super(file, new MappedFileSimpleStorage(file, initialSize, 1024 * 1024), dataDescriptor, initialSize,
super(file, new ResizeableMappedFile(file, initialSize, ourLock, VALUE_PAGE_SIZE, true), dataDescriptor, initialSize,
ourVersion, new RecordBufferHandler(), false);
myInlineKeysNoMapping = myDataDescriptor instanceof InlineKeyDescriptor && !wantKeyMapping();
@@ -169,7 +173,8 @@ public class PersistentBTreeEnumerator<Data> extends PersistentEnumeratorBase<Da
synchronized (ourLock) {
List<IntToIntBtree.BtreeIndexNodeView> leafPages = new ArrayList<IntToIntBtree.BtreeIndexNodeView> ();
btree.root.setAddress(btree.root.address);
btree.doFlush();
btree.root.syncWithStore();
collectLeafPages(btree.root, leafPages);
Collections.sort(leafPages, new Comparator<IntToIntBtree.BtreeIndexNodeView>() {
@Override
@@ -229,7 +234,15 @@ public class PersistentBTreeEnumerator<Data> extends PersistentEnumeratorBase<Da
@Override
protected int setupValueId(int hashCode, int dataOff) {
if (myExternalKeysNoMapping) return dataOff + KEY_SHIFT;
return super.setupValueId(hashCode, dataOff);
final PersistentEnumeratorBase.RecordBufferHandler<PersistentEnumeratorBase> recordHandler = getRecordHandler();
final byte[] buf = recordHandler.getRecordBuffer(this);
// optimization for using putInt / getInt on aligned empty storage (our page always contains on storage ByteBuffer)
final int pos = recordHandler.recordWriteOffset(this, buf);
myStorage.ensureSize(pos + buf.length);
if (!myInlineKeysNoMapping) myStorage.putInt(pos, dataOff);
return pos;
}
@Override
@@ -49,7 +49,7 @@ public class PersistentEnumerator<Data> extends PersistentEnumeratorBase<Data> {
private static final Version ourVersion = new Version(CORRECTLY_CLOSED_MAGIC, DIRTY_MAGIC);
public PersistentEnumerator(File file, KeyDescriptor<Data> dataDescriptor, int initialSize) throws IOException {
super(file, new MappedFileSimpleStorage(file, initialSize), dataDescriptor, initialSize, ourVersion,
super(file, new ResizeableMappedFile(file, initialSize, ourLock), dataDescriptor, initialSize, ourVersion,
new RecordBufferHandler(), true);
}
@@ -44,7 +44,7 @@ abstract class PersistentEnumeratorBase<Data> implements Forceable, Closeable {
private static final int META_DATA_OFFSET = 4;
protected static final int DATA_START = META_DATA_OFFSET + 16;
protected final ISimpleStorage myStorage;
protected final ResizeableMappedFile myStorage;
private final ResizeableMappedFile myKeyStorage;
private boolean myClosed = false;
@@ -137,7 +137,7 @@ abstract class PersistentEnumeratorBase<Data> implements Forceable, Closeable {
}
}
public PersistentEnumeratorBase(File file, ISimpleStorage storage, KeyDescriptor<Data> dataDescriptor, int initialSize,
public PersistentEnumeratorBase(File file, ResizeableMappedFile storage, KeyDescriptor<Data> dataDescriptor, int initialSize,
Version version, RecordBufferHandler<? extends PersistentEnumeratorBase> recordBufferHandler,
boolean doCaching) throws IOException {
myDataDescriptor = dataDescriptor;
@@ -357,28 +357,30 @@ abstract class PersistentEnumeratorBase<Data> implements Forceable, Closeable {
}
protected boolean iterateData(final Processor<Data> processor) throws IOException {
if (myKeyStorage == null) {
throw new UnsupportedOperationException("Iteration over InlineIntegerKeyDescriptors is not supported");
}
synchronized (ourLock) {
if (myKeyStorage == null) {
throw new UnsupportedOperationException("Iteration over InlineIntegerKeyDescriptors is not supported");
}
myKeyStorage.force();
myKeyStorage.force();
DataInputStream keysStream = new DataInputStream(new BufferedInputStream(new LimitedInputStream(new FileInputStream(keystreamFile()),
(int)myKeyStorage.length())));
try {
DataInputStream keysStream = new DataInputStream(new BufferedInputStream(new LimitedInputStream(new FileInputStream(keystreamFile()),
(int)myKeyStorage.length())));
try {
while (true) {
Data key = myDataDescriptor.read(keysStream);
if (!processor.process(key)) return false;
try {
while (true) {
Data key = myDataDescriptor.read(keysStream);
if (!processor.process(key)) return false;
}
}
catch (EOFException e) {
// Done
}
return true;
}
catch (EOFException e) {
// Done
finally {
keysStream.close();
}
return true;
}
finally {
keysStream.close();
}
}
@@ -386,7 +388,7 @@ abstract class PersistentEnumeratorBase<Data> implements Forceable, Closeable {
return new File(myFile.getPath() + ".keystream");
}
public synchronized Data valueOf(int idx) throws IOException {
public Data valueOf(int idx) throws IOException {
synchronized (ourLock) {
try {
int addr = indexToAddr(idx);
@@ -108,22 +108,24 @@ public class PersistentHashMap<Key, Value> extends PersistentEnumeratorDelegate<
}
protected void onDropFromCache(final Key key, final AppendStream value) {
try {
final ByteSequence bytes = value.getInternalBuffer();
final int id = enumerate(key);
HeaderRecord oldHeaderRecord = readValueId(id);
synchronized (PersistentEnumerator.ourLock) {
try {
final ByteSequence bytes = value.getInternalBuffer();
final int id = enumerate(key);
HeaderRecord oldHeaderRecord = readValueId(id);
HeaderRecord headerRecord = new HeaderRecord(
myValueStorage.appendBytes(bytes, oldHeaderRecord.address)
);
HeaderRecord headerRecord = new HeaderRecord(
myValueStorage.appendBytes(bytes, oldHeaderRecord.address)
);
updateValueId(id, headerRecord, oldHeaderRecord, key, 0);
if (oldHeaderRecord == HeaderRecord.EMPTY) myLiveAndGarbageKeysCounter += LIVE_KEY_MASK;
updateValueId(id, headerRecord, oldHeaderRecord, key, 0);
if (oldHeaderRecord == HeaderRecord.EMPTY) myLiveAndGarbageKeysCounter += LIVE_KEY_MASK;
myStreamPool.recycle(value);
}
catch (IOException e) {
throw new RuntimeException(e);
myStreamPool.recycle(value);
}
catch (IOException e) {
throw new RuntimeException(e);
}
}
}
};
@@ -280,12 +282,10 @@ public class PersistentHashMap<Key, Value> extends PersistentEnumeratorDelegate<
}
public synchronized void appendData(Key key, ValueDataAppender appender) throws IOException {
synchronized (PersistentEnumerator.ourLock) {
myEnumerator.markDirty(true);
final AppendStream stream = myAppendCache.get(key);
appender.append(stream);
}
myEnumerator.markDirty(true);
final AppendStream stream = myAppendCache.get(key);
appender.append(stream);
}
/**
@@ -293,10 +293,8 @@ public class PersistentHashMap<Key, Value> extends PersistentEnumeratorDelegate<
* {@link #processKeysWithExistingMapping(com.intellij.util.Processor)} to process only keys with existing mappings
*/
public synchronized boolean processKeys(Processor<Key> processor) throws IOException {
synchronized (PersistentEnumerator.ourLock) {
myAppendCache.clear();
return myEnumerator.iterateData(processor);
}
myAppendCache.clear();
return myEnumerator.iterateData(processor);
}
public Collection<Key> getAllKeysWithExistingMapping() throws IOException {
@@ -1,75 +0,0 @@
package com.intellij.util.io;
import com.intellij.openapi.util.io.FileUtil;
import java.io.File;
import java.io.IOException;
class RandomAccessFileSimpleStorage implements ISimpleStorage {
private final RandomAccessDataFile storage;
public RandomAccessFileSimpleStorage(File file, PagePool pool) throws IOException {
FileUtil.createIfDoesntExist(file);
storage = new RandomAccessDataFile(file, pool);
}
@Override
public void put(int index, byte value) {
storage.putByte(index, value);
}
@Override
public byte get(int index) {
return storage.getByte(index);
}
@Override
public void putInt(int index, int value) {
storage.putInt(index, value);
}
@Override
public int getInt(int index) {
return storage.getInt(index);
}
@Override
public void putLong(int pos, long value) {
storage.putLong(pos, value);
}
@Override
public long getLong(int pos) {
return storage.getLong(pos);
}
@Override
public long length() {
return storage.length();
}
@Override
public void put(int pos, byte[] buf, int offset, int length) {
storage.put(pos, buf, offset, length);
}
@Override
public void get(int pos, byte[] buf, int offset, int length) {
storage.get(pos, buf, offset, length);
}
@Override
public void close() throws IOException {
storage.close();
}
@Override
public boolean isDirty() {
return storage.isDirty();
}
@Override
public void force() {
storage.force();
}
}
@@ -21,6 +21,7 @@ package com.intellij.util.io;
import com.intellij.openapi.Forceable;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.util.io.FileUtil;
import java.io.*;
@@ -30,9 +31,11 @@ public class ResizeableMappedFile implements Forceable {
private long myLogicalSize;
private final PagedFileStorage myStorage;
public ResizeableMappedFile(final File file, int initialSize, PagedFileStorage.StorageLock lock, int pageSize) throws IOException {
myStorage = new PagedFileStorage(file, lock, pageSize);
if (!file.exists() || file.length() == 0) {
public ResizeableMappedFile(final File file, int initialSize, PagedFileStorage.StorageLock lock, int pageSize, boolean valuesAreBufferAligned) throws IOException {
myStorage = new PagedFileStorage(file, lock, pageSize, valuesAreBufferAligned);
boolean exists = file.exists();
if (!exists || file.length() == 0) {
if (!exists) FileUtil.createParentDirs(file);
writeLength(0);
}
@@ -45,7 +48,7 @@ public class ResizeableMappedFile implements Forceable {
}
public ResizeableMappedFile(final File file, int initialSize, PagedFileStorage.StorageLock lock) throws IOException {
this(file, initialSize, lock, PagedFileStorage.DEFAULT_BUFFER_SIZE);
this(file, initialSize, lock, PagedFileStorage.DEFAULT_BUFFER_SIZE, false);
}
public long length() {
@@ -65,7 +68,7 @@ public class ResizeableMappedFile implements Forceable {
}
}
private void ensureSize(final long pos) {
void ensureSize(final long pos) {
if (pos + 16 > Integer.MAX_VALUE) throw new RuntimeException("FATAL ERROR: Can't get over 2^32 address space");
myLogicalSize = Math.max(pos, myLogicalSize);
while (pos >= realSize()) {
@@ -150,6 +153,15 @@ public class ResizeableMappedFile implements Forceable {
myStorage.putInt(index, value);
}
public short getShort(int index) {
return myStorage.getShort(index);
}
public void putShort(int index, short value) {
ensureSize(index + 2);
myStorage.putShort(index, value);
}
public long getLong(int index) {
return myStorage.getLong(index);
}
@@ -186,4 +198,7 @@ public class ResizeableMappedFile implements Forceable {
}
}
public PagedFileStorage getPagedFileStorage() {
return myStorage;
}
}