Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -16,25 +16,26 @@

package org.dizitart.no2.mvstore;

import java.io.File;
import java.util.HashSet;
import java.util.Set;

import org.dizitart.no2.store.events.StoreEventListener;
import org.h2.mvstore.FileStore;

import lombok.AccessLevel;
import lombok.Getter;
import lombok.Setter;
import lombok.experimental.Accessors;
import org.dizitart.no2.store.events.StoreEventListener;
import org.h2.mvstore.FileStore;

import java.io.File;
import java.util.HashSet;
import java.util.Set;

/**
* The MVStoreModuleBuilder class is responsible for building an instance of
* {@link MVStoreModule}. It provides methods to set various configuration
* options for the MVStore database.
*
* @since 4.0
* @see MVStoreModule
*
* @author Anindya Chatterjee
* @see MVStoreModule
* @since 4.0
*/
@Getter
@Setter
Expand Down Expand Up @@ -81,6 +82,13 @@ public class MVStoreModuleBuilder {
*/
private boolean autoCommit = true;

/**
* Flag to enable/disable auto-compact mode. If set to true, fragmented
* chunks or chunks that are sufficiently below the target fill rate of 90%
* are reclaimed. This will typically shrink the file.
*/
private boolean autoCompact = true;

/**
* Indicates whether the MVStore should be opened in recovery mode or not.
*/
Expand Down Expand Up @@ -174,6 +182,7 @@ public MVStoreModule build() {
dbConfig.compress(compress());
dbConfig.compressHigh(compressHigh());
dbConfig.autoCommit(autoCommit());
dbConfig.autoCompact(autoCompact());
dbConfig.recoveryMode(recoveryMode());
dbConfig.cacheSize(cacheSize());
dbConfig.cacheConcurrency(cacheConcurrency());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,17 +16,18 @@

package org.dizitart.no2.mvstore;

import lombok.extern.slf4j.Slf4j;
import static org.dizitart.no2.common.util.StringUtils.isNullOrEmpty;

import java.io.File;
import java.util.Map;

import org.dizitart.no2.exceptions.InvalidOperationException;
import org.dizitart.no2.exceptions.NitriteIOException;
import org.dizitart.no2.mvstore.compat.v1.UpgradeUtil;
import org.h2.mvstore.MVStore;
import org.h2.mvstore.MVStoreException;

import java.io.File;
import java.util.Map;

import static org.dizitart.no2.common.util.StringUtils.isNullOrEmpty;
import lombok.extern.slf4j.Slf4j;

/**
* @author Anindya Chatterjee.
Expand Down Expand Up @@ -55,7 +56,7 @@ static MVStore openOrCreate(MVStoreConfig storeConfig) {
}
} catch (MVStoreException me) {
if (me.getMessage().contains("file is locked")) {
throw new NitriteIOException("Database is already opened in other process");
throw new NitriteIOException("Database is already opened in other process", me);
}

if (dbFile != null) {
Expand Down Expand Up @@ -128,8 +129,10 @@ private static MVStore.Builder createBuilder(MVStoreConfig mvStoreConfig) {
builder = builder.autoCommitBufferSize(mvStoreConfig.autoCommitBufferSize());
}

// auto compact disabled github issue #41
builder.autoCompactFillRate(0);
if (!mvStoreConfig.autoCompact()) {
// disables background compaction
builder.autoCompactFillRate(0);
}

if (mvStoreConfig.encryptionKey() != null) {
builder = builder.encryptionKey(mvStoreConfig.encryptionKey());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,45 +16,56 @@

package org.dizitart.no2.mvstore;

import static org.dizitart.no2.common.util.ValidationUtils.notNull;

import java.lang.ref.Cleaner;
import java.util.Iterator;
import java.util.Map;
import java.util.NoSuchElementException;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Supplier;

import org.dizitart.no2.common.RecordStream;
import org.dizitart.no2.common.tuples.Pair;
import org.dizitart.no2.exceptions.NitriteIOException;
import org.dizitart.no2.store.NitriteMap;
import org.dizitart.no2.store.NitriteStore;
import org.h2.mvstore.MVMap;
import org.h2.mvstore.MVStore;

import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;

import static org.dizitart.no2.common.util.ValidationUtils.notNull;

/**
* @since 1.0
* @author Anindya Chatterjee
* @since 1.0
*/
class NitriteMVMap<Key, Value> implements NitriteMap<Key, Value> {

private static final Cleaner CLEANER = Cleaner.create();

private final MVMap<Key, Value> mvMap;
private final NitriteStore<?> nitriteStore;
private final MVStore mvStore;
private final AtomicBoolean droppedFlag;
private final AtomicBoolean closedFlag;
private final Set<VersionUsage> versionUsages;

NitriteMVMap(MVMap<Key, Value> mvMap, NitriteStore<?> nitriteStore) {
NitriteMVMap(final MVMap<Key, Value> mvMap, final NitriteStore<?> nitriteStore) {
this.mvMap = mvMap;
this.nitriteStore = nitriteStore;
this.mvStore = mvMap.getStore();
this.closedFlag = new AtomicBoolean(false);
this.droppedFlag = new AtomicBoolean(false);
this.versionUsages = ConcurrentHashMap.newKeySet();
}

@Override
public boolean containsKey(Key key) {
public boolean containsKey(final Key key) {
return mvMap.containsKey(key);
}

@Override
public Value get(Key key) {
public Value get(final Key key) {
return mvMap.get(key);
}

Expand All @@ -65,7 +76,7 @@ public NitriteStore<?> getStore() {

@Override
public void clear() {
MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
final MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
try {
mvMap.clear();
updateLastModifiedTime();
Expand All @@ -81,14 +92,14 @@ public String getName() {

@Override
public RecordStream<Value> values() {
return RecordStream.fromIterable(mvMap.values());
return () -> versionedIterator(() -> mvMap.values().iterator());
}

@Override
public Value remove(Key key) {
MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
public Value remove(final Key key) {
final MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
try {
Value value = mvMap.remove(key);
final Value value = mvMap.remove(key);
updateLastModifiedTime();
return value;
} finally {
Expand All @@ -98,13 +109,13 @@ public Value remove(Key key) {

@Override
public RecordStream<Key> keys() {
return RecordStream.fromIterable(mvMap.keySet());
return () -> versionedIterator(() -> mvMap.keySet().iterator());
}

@Override
public void put(Key key, Value value) {
public void put(final Key key, final Value value) {
notNull(value, "value cannot be null");
MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
final MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
try {
mvMap.put(key, value);
updateLastModifiedTime();
Expand All @@ -119,11 +130,11 @@ public long size() {
}

@Override
public Value putIfAbsent(Key key, Value value) {
public Value putIfAbsent(final Key key, final Value value) {
notNull(value, "value cannot be null");
MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
final MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
try {
Value v = mvMap.putIfAbsent(key, value);
final Value v = mvMap.putIfAbsent(key, value);
updateLastModifiedTime();
return v;
} finally {
Expand All @@ -133,8 +144,8 @@ public Value putIfAbsent(Key key, Value value) {

@Override
public RecordStream<Pair<Key, Value>> entries() {
return () -> new Iterator<>() {
final Iterator<Map.Entry<Key, Value>> entryIterator = mvMap.entrySet().iterator();
return () -> versionedIterator(() -> new Iterator<>() {
private final Iterator<Map.Entry<Key, Value>> entryIterator = mvMap.entrySet().iterator();

@Override
public boolean hasNext() {
Expand All @@ -143,15 +154,19 @@ public boolean hasNext() {

@Override
public Pair<Key, Value> next() {
Map.Entry<Key, Value> entry = entryIterator.next();
final Map.Entry<Key, Value> entry = entryIterator.next();
return new Pair<>(entry.getKey(), entry.getValue());
}
};
});
}

@Override
public RecordStream<Pair<Key, Value>> reversedEntries() {
return () -> new ReverseIterator<>(mvMap);
return () -> versionedIterator(() -> new ReverseIterator<>(mvMap));
}

private <Element> Iterator<Element> versionedIterator(final Supplier<Iterator<Element>> iteratorSupplier) {
return new VersionedIterator<>(mvStore, iteratorSupplier, versionUsages);
}

@Override
Expand All @@ -165,22 +180,22 @@ public Key lastKey() {
}

@Override
public Key higherKey(Key key) {
public Key higherKey(final Key key) {
return mvMap.higherKey(key);
}

@Override
public Key ceilingKey(Key key) {
public Key ceilingKey(final Key key) {
return mvMap.ceilingKey(key);
}

@Override
public Key lowerKey(Key key) {
public Key lowerKey(final Key key) {
return mvMap.lowerKey(key);
}

@Override
public Key floorKey(Key key) {
public Key floorKey(final Key key) {
return mvMap.floorKey(key);
}

Expand All @@ -194,11 +209,12 @@ public void drop() {
if (!droppedFlag.get()) {
droppedFlag.compareAndSet(false, true);
closedFlag.compareAndSet(false, true);
releaseVersionUsages();

MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
final MVStore.TxCounter txCounter = mvStore.registerVersionUsage();
try {
nitriteStore.closeMap(getName());
nitriteStore.removeMap(getName());
nitriteStore.closeMap(mvMap.getName());
nitriteStore.removeMap(mvMap.getName());
} finally {
mvStore.deregisterVersionUsage(txCounter);
}
Expand All @@ -214,12 +230,82 @@ public boolean isDropped() {
public void close() {
if (!closedFlag.get() && !droppedFlag.get()) {
closedFlag.compareAndSet(false, true);
nitriteStore.closeMap(getName());
releaseVersionUsages();
nitriteStore.closeMap(mvMap.getName());
}
}

@Override
public boolean isClosed() {
return closedFlag.get();
}

private void releaseVersionUsages() {
for (final VersionUsage versionUsage : versionUsages) {
versionUsage.release();
}
}

private static class VersionedIterator<Element> implements Iterator<Element> {

private final Iterator<Element> iterator;
private final Cleaner.Cleanable cleanable;
private final VersionUsage versionUsage;
private boolean exhausted;

private VersionedIterator(final MVStore mvStore,
final Supplier<Iterator<Element>> iteratorSupplier,
final Set<VersionUsage> versionUsages) {

versionUsage = new VersionUsage(mvStore, mvStore.registerVersionUsage(), versionUsages);
versionUsages.add(versionUsage);

try {
this.iterator = iteratorSupplier.get();
this.cleanable = CLEANER.register(this, versionUsage::release);
} catch (final RuntimeException | Error e) {
versionUsage.release();
throw e;
}
}

@Override
public boolean hasNext() {
if (exhausted) {
return false;
}
ensureOpen();
try {
final boolean hasNext = iterator.hasNext();
if (!hasNext) {
exhausted = true;
cleanable.clean();
}
return hasNext;
} catch (final RuntimeException | Error e) {
cleanable.clean();
throw e;
}
}

@Override
public Element next() {
if (exhausted) {
throw new NoSuchElementException();
}
ensureOpen();
try {
return iterator.next();
} catch (final RuntimeException | Error e) {
cleanable.clean();
throw e;
}
}

private void ensureOpen() {
if (versionUsage.isReleased()) {
throw new NitriteIOException("MVStore is closed");
}
}
}
}
Loading
Loading