Skip to content
Merged
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 @@ -152,6 +152,7 @@ boolean cleanPartition(String graph, int partId, long startKey, long endKey,

default void doBatch(String graph, int partId, List<BatchEntry> entryList) {
BusinessHandler.TxBuilder builder = txBuilder(graph, partId);
BusinessHandler.Tx transaction = builder.build();
try {
for (BatchEntry b : entryList) {
Key start = b.getStartKey();
Expand Down Expand Up @@ -185,12 +186,16 @@ default void doBatch(String graph, int partId, List<BatchEntry> entryList) {
}
}
}
builder.build().commit();
transaction.commit();
} catch (Throwable e) {
String msg =
String.format("graph data %s-%s do batch insert with error:", graph, partId);
log.error(msg, e);
builder.build().rollback();
try {
transaction.rollback();
} catch (Throwable rollbackError) {
e.addSuppressed(rollbackError);
}
throw e;
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,9 @@
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.BiFunction;
import java.util.function.Consumer;
import java.util.function.Function;
Expand Down Expand Up @@ -121,6 +124,8 @@ public class BusinessHandlerImpl implements BusinessHandler {

private static final Map<String, HugeGraphSupplier> GRAPH_SUPPLIER_CACHE =
new ConcurrentHashMap<>();
private static final int GRAPH_LOCK_STRIPES = 1024;
private static final ReadWriteLock[] GRAPH_LOCKS = createGraphLocks();
private static final int batchSize = 10000;
private static Long indexDataSize = 50 * 1024L;
private static final RocksDBFactory factory = RocksDBFactory.getInstance();
Expand Down Expand Up @@ -172,6 +177,20 @@ public void onDBSessionReleased(RocksDBSession dbSession) {
});
}

private static ReadWriteLock[] createGraphLocks() {
ReadWriteLock[] locks = new ReadWriteLock[GRAPH_LOCK_STRIPES];
for (int i = 0; i < locks.length; i++) {
locks[i] = new ReentrantReadWriteLock(true);
}
return locks;
}

private static ReadWriteLock graphLock(String graph, int partId) {
int hash = 31 * partId + graph.hashCode();
hash ^= hash >>> 16;
return GRAPH_LOCKS[hash & (GRAPH_LOCK_STRIPES - 1)];
}

public static HugeConfig initRocksdb(Map<String, Object> rocksdbConfig,
RocksdbChangedListener listener) {
// Register rocksdb configuration
Expand Down Expand Up @@ -1014,11 +1033,17 @@ public void batchGet(String graph, String table, Supplier<HgPair<Integer, byte[]
public void truncate(String graphName, int partId) throws HgStoreException {
// Each partition corresponds to a rocksdb instance, so the rocksdb instance name is
// rocksdb + partId
try (RocksDBSession dbSession = getSession(graphName, partId)) {
dbSession.sessionOp().deleteRange(keyCreator.getStartKey(partId, graphName),
keyCreator.getEndKey(partId, graphName));
// Release map ID
keyCreator.delGraphId(partId, graphName);
Lock lifecycleLock = graphLock(graphName, partId).writeLock();
lifecycleLock.lock();
try {
try (RocksDBSession dbSession = getSession(graphName, partId)) {
dbSession.sessionOp().deleteRange(keyCreator.getStartKey(partId, graphName),
keyCreator.getEndKey(partId, graphName));
// Release map ID
keyCreator.delGraphId(partId, graphName);
}
} finally {
lifecycleLock.unlock();
}
}

Expand Down Expand Up @@ -1287,7 +1312,7 @@ private void deleteGraphDatabase(String graph, int partId) throws IOException {

@Override
public TxBuilder txBuilder(String graph, int partId) throws HgStoreException {
return new TxBuilderImpl(graph, partId, getSession(graph, partId));
return new TxBuilderImpl(graph, partId);
}

@Override
Expand Down Expand Up @@ -1564,20 +1589,42 @@ private class TxBuilderImpl implements TxBuilder {
private final int partId;
private final RocksDBSession dbSession;
private final SessionOperator op;
private final Lock lifecycleLock;
private boolean completed;

private TxBuilderImpl(String graph, int partId, RocksDBSession dbSession) {
private TxBuilderImpl(String graph, int partId) {
this.graph = graph;
this.partId = partId;
this.dbSession = dbSession;
this.op = this.dbSession.sessionOp();
this.op.prepare();
this.lifecycleLock = graphLock(graph, partId).readLock();
this.lifecycleLock.lock();

RocksDBSession session = null;
SessionOperator operator = null;
try {
session = getSession(graph, partId);
operator = session.sessionOp();
operator.prepare();
} catch (RuntimeException | Error e) {
try {
if (session != null) {
session.close();
}
} catch (Throwable closeError) {
e.addSuppressed(closeError);
} finally {
this.lifecycleLock.unlock();
}
throw e;
}
this.dbSession = session;
this.op = operator;
}

@Override
public TxBuilder put(int code, String table, byte[] key, byte[] value) throws
HgStoreException {
try {
byte[] targetKey = keyCreator.getKey(this.partId, graph, code, key);
byte[] targetKey = keyCreator.getKeyOrCreate(this.partId, graph, code, key);
Comment thread
imbajin marked this conversation as resolved.
Comment thread
contrueCT marked this conversation as resolved.
this.op.put(table, targetKey, value);
} catch (DBStoreException e) {
throw new HgStoreException(HgStoreException.EC_RKDB_DOPUT_FAIL, e.toString());
Expand Down Expand Up @@ -1642,7 +1689,7 @@ public TxBuilder merge(int code, String table, byte[] key, byte[] value) throws
HgStoreException {

try {
byte[] targetKey = keyCreator.getKey(this.partId, graph, code, key);
byte[] targetKey = keyCreator.getKeyOrCreate(this.partId, graph, code, key);
Comment thread
contrueCT marked this conversation as resolved.
op.merge(table, targetKey, value);
} catch (DBStoreException e) {
throw new HgStoreException(HgStoreException.EC_RKDB_DOMERGE_FAIL, e.toString());
Expand All @@ -1655,21 +1702,37 @@ public Tx build() {
return new Tx() {
@Override
public void commit() throws HgStoreException {
if (completed) {
return;
}
op.commit(); // After an exception occurs in commit, rollback must be
// called, otherwise it will cause the lock not to be released.
dbSession.close();
completed = true;
release();
}

@Override
public void rollback() throws HgStoreException {
if (completed) {
return;
}
try {
op.rollback();
} finally {
dbSession.close();
completed = true;
release();
}
}
};
}

private void release() {
try {
this.dbSession.close();
} finally {
this.lifecycleLock.unlock();
}
}
}

public static void clearCache() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -222,12 +222,22 @@ public void cleanData(Metapb.Partition partition) {

@Override
public void write(BatchPutRequest request) {
BusinessHandler.TxBuilder tx =
BusinessHandler.TxBuilder builder =
businessHandler.txBuilder(request.getGraphName(), request.getPartitionId());
for (BatchPutRequest.KV kv : request.getEntries()) {
tx.put(kv.getCode(), kv.getTable(), kv.getKey(), kv.getValue());
BusinessHandler.Tx transaction = builder.build();
try {
for (BatchPutRequest.KV kv : request.getEntries()) {
builder.put(kv.getCode(), kv.getTable(), kv.getKey(), kv.getValue());
}
transaction.commit();
} catch (Throwable e) {
try {
transaction.rollback();
} catch (Throwable rollbackError) {
e.addSuppressed(rollbackError);
}
throw e;
}
tx.build().commit();
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -221,12 +221,22 @@ public void cleanData(Metapb.Partition partition) {

@Override
public void doWriteData(BatchPutRequest request) {
BusinessHandler.TxBuilder tx =
BusinessHandler.TxBuilder builder =
businessHandler.txBuilder(request.getGraphName(), request.getPartitionId());
for (BatchPutRequest.KV kv : request.getEntries()) {
tx.put(kv.getCode(), kv.getTable(), kv.getKey(), kv.getValue());
BusinessHandler.Tx transaction = builder.build();
try {
for (BatchPutRequest.KV kv : request.getEntries()) {
builder.put(kv.getCode(), kv.getTable(), kv.getKey(), kv.getValue());
}
transaction.commit();
} catch (Throwable e) {
try {
transaction.rollback();
} catch (Throwable rollbackError) {
e.addSuppressed(rollbackError);
}
throw e;
}
tx.build().commit();
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@

package org.apache.hugegraph.store.meta;

import static org.apache.hugegraph.store.constant.HugeServerTables.VERTEX_TABLE;

import java.nio.ByteBuffer;
import java.util.Arrays;
import java.util.List;
Expand Down Expand Up @@ -128,8 +130,13 @@ public long releaseGraphId(String graphName) {
private boolean checkCount(long l) {
var start = new byte[2];
Bits.putShort(start, 0, (short) l);
try (var itr = sessionBuilder.getSession(partitionId).sessionOp().scan("g+v", start)) {
return itr == null || !itr.hasNext();
try (var session = sessionBuilder.getSession(partitionId)) {
if (!session.tableIsExist(VERTEX_TABLE)) {
return true;
}
try (var itr = session.sessionOp().scan(VERTEX_TABLE, start)) {
return itr == null || !itr.hasNext();
}
}
}

Expand Down
1 change: 1 addition & 0 deletions hugegraph-store/hg-store-test/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,7 @@
</testClassesDirectory>
<includes>
<include>**/CoreSuiteTest.java</include>
<include>**/BatchGraphIsolationTest.java</include>
</includes>
</configuration>
</execution>
Expand Down
Loading
Loading