GraphPlatform-2060 更改store依赖common更改为开源的org.apache.hugegraph.common包

Change-Id: Ie3997bad153358f251b242cae8de13ff8022976e
This commit is contained in:
chengxin05 2023-05-09 13:20:31 +08:00 committed by imbajin
parent e79d5f63a1
commit 33749fddaf
9 changed files with 539 additions and 15 deletions

View File

@ -0,0 +1,306 @@
package com.baidu.hugegraph.store;
import com.alipay.sofa.jraft.Status;
import com.alipay.sofa.jraft.util.Utils;
import com.baidu.hugegraph.pd.common.PDException;
import com.baidu.hugegraph.pd.grpc.MetaTask;
import com.baidu.hugegraph.pd.grpc.Metapb;
import com.baidu.hugegraph.pd.grpc.pulse.ChangeShard;
import com.baidu.hugegraph.pd.grpc.pulse.CleanPartition;
import com.baidu.hugegraph.pd.grpc.pulse.DbCompaction;
import com.baidu.hugegraph.pd.grpc.pulse.MovePartition;
import com.baidu.hugegraph.pd.grpc.pulse.PartitionKeyRange;
import com.baidu.hugegraph.pd.grpc.pulse.SplitPartition;
import com.baidu.hugegraph.pd.grpc.pulse.TransferLeader;
import com.baidu.hugegraph.store.cmd.CleanDataRequest;
import com.baidu.hugegraph.store.cmd.DbCompactionRequest;
import com.baidu.hugegraph.store.meta.MetadataKeyHelper;
import com.baidu.hugegraph.store.meta.Partition;
import com.baidu.hugegraph.store.pd.PartitionInstructionListener;
import com.baidu.hugegraph.store.raft.RaftClosure;
import com.baidu.hugegraph.store.raft.RaftOperation;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
import org.apache.hugegraph.util.Log;
import org.slf4j.Logger;
import java.io.IOException;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
/**
* PD发给Store的分区指令处理器
*/
public class PartitionInstructionProcessor implements PartitionInstructionListener {
private static final Logger LOG = Log.logger(PartitionInstructionProcessor.class);
private HgStoreEngine storeEngine;
private final ExecutorService threadPool;
public PartitionInstructionProcessor(HgStoreEngine storeEngine) {
this.storeEngine = storeEngine;
ThreadFactory namedThreadFactory = new ThreadFactoryBuilder().setNameFormat("instruct-process-pool-%d").build();
threadPool = new ThreadPoolExecutor(Runtime.getRuntime().availableProcessors(),
1000000,
180L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(1000000),
namedThreadFactory,
new ThreadPoolExecutor.AbortPolicy());
}
@Override
public void onChangeShard(long taskId, Partition partition, ChangeShard changeShard, Consumer<Integer> consumer) {
PartitionEngine engine = storeEngine.getPartitionEngine(partition.getId());
if (engine != null) {
// 清理所有的任务有失败的情况
engine.getTaskManager().deleteTask(partition.getId(), MetaTask.TaskType.Change_Shard.name());
}
if (engine != null && engine.isLeader()) {
LOG.info("Partition {}-{} Receive change shard message, {}", partition.getGraphName(),
partition.getId(), changeShard);
String graphName = partition.getGraphName();
int partitionId = partition.getId();
MetaTask.Task task = MetaTask.Task.newBuilder()
.setId(taskId)
.setPartition(partition.getProtoObj())
.setType(MetaTask.TaskType.Change_Shard)
.setState(MetaTask.TaskState.Task_Ready)
.setChangeShard(changeShard)
.build();
try {
storeEngine.addRaftTask(graphName, partitionId,
RaftOperation.create(RaftOperation.SYNC_PARTITION_TASK, task),
new RaftClosure() {
@Override
public void run(Status status) {
LOG.info("Partition {}-{} onChangeShard complete, status is {}",
graphName, partitionId, status);
consumer.accept(0);
}
});
} catch (Exception e) {
LOG.error("Partition {}-{} onSplitPartition exception {}",
graphName, partitionId, e);
}
}
}
@Override
public void onTransferLeader(long taskId, Partition partition, TransferLeader transferLeader, Consumer<Integer> consumer) {
PartitionEngine engine = storeEngine.getPartitionEngine(partition.getId());
if (engine != null && engine.isLeader()) {
consumer.accept(0);
Utils.runInThread(() -> {
LOG.info("Partition {}-{} receive TransferLeader instruction, new leader is {}"
, partition.getGraphName(), partition.getId(), transferLeader.getShard());
engine.transferLeader(partition.getGraphName(), transferLeader.getShard());
});
}
}
/**
* Leader接收到PD发送的分区分裂任务
* 添加到raft任务队列由raft进行任务分发
*/
@Override
public void onSplitPartition(long taskId, Partition partition, SplitPartition splitPartition,
Consumer<Integer> consumer) {
PartitionEngine engine = storeEngine.getPartitionEngine(partition.getId());
if (preCheckTaskId(taskId, partition.getId())) {
return;
}
if (engine != null && engine.isLeader()) {
// 先应答避免超时造成pd重复发送
consumer.accept(0);
String graphName = partition.getGraphName();
int partitionId = partition.getId();
MetaTask.Task task = MetaTask.Task.newBuilder()
.setId(taskId)
.setPartition(partition.getProtoObj())
.setType(MetaTask.TaskType.Split_Partition)
.setState(MetaTask.TaskState.Task_Ready)
.setSplitPartition(splitPartition)
.build();
try {
threadPool.submit(()->{
engine.moveData(task);
});
} catch (Exception e) {
LOG.error("Partition {}-{} onSplitPartition exception {}",
graphName, partitionId, e);
}
}
}
/**
* Leader接收到PD发送的rocksdb compaction任务
* 添加到raft任务队列由raft进行任务分发
*/
@Override
public void onDbCompaction(long taskId, Partition partition, DbCompaction dbCompaction,
Consumer<Integer> consumer) {
PartitionEngine engine = storeEngine.getPartitionEngine(partition.getId());
if (engine != null && engine.isLeader()) {
try {
DbCompactionRequest dbCompactionRequest = new DbCompactionRequest();
dbCompactionRequest.setPartitionId(partition.getId());
dbCompactionRequest.setTableName(dbCompaction.getTableName());
dbCompactionRequest.setGraphName(partition.getGraphName());
engine.addRaftTask(RaftOperation.create(RaftOperation.DB_COMPACTION,
dbCompactionRequest),
new RaftClosure() {
@Override
public void run(Status status) {
LOG.info("onRocksdbCompaction {}-{} sync partition status is {}",
partition.getGraphName(), partition.getId(), status);
}
}
);
} finally {
consumer.accept(0);
}
}
}
@Override
public void onMovePartition(long taskId, Partition partition, MovePartition movePartition,
Consumer<Integer> consumer) {
PartitionEngine engine = storeEngine.getPartitionEngine(partition.getId());
if (preCheckTaskId(taskId, partition.getId())) {
return;
}
if (engine != null && engine.isLeader()){
// 先应答避免超时造成pd重复发送
consumer.accept(0);
String graphName = partition.getGraphName();
int partitionId = partition.getId();
MetaTask.Task task = MetaTask.Task.newBuilder()
.setId(taskId)
.setPartition(partition.getProtoObj())
.setType(MetaTask.TaskType.Move_Partition)
.setState(MetaTask.TaskState.Task_Ready)
.setMovePartition(movePartition)
.build();
try {
threadPool.submit(()->{
engine.moveData(task);
});
} catch (Exception e) {
LOG.error("Partition {}-{} onMovePartition exception {}",
graphName, partitionId, e);
}
}
}
@Override
public void onCleanPartition(long taskId, Partition partition, CleanPartition cleanPartition,
Consumer<Integer> consumer) {
if (preCheckTaskId(taskId, partition.getId())) {
return;
}
PartitionEngine engine = storeEngine.getPartitionEngine(partition.getId());
if (engine != null && engine.isLeader()){
consumer.accept(0);
CleanDataRequest request = CleanDataRequest.fromCleanPartitionTask(cleanPartition, partition, taskId);
storeEngine.addRaftTask(partition.getGraphName(), partition.getId(),
RaftOperation.create(RaftOperation.IN_CLEAN_OP, request),
status -> {
LOG.info("onCleanPartition {}-{}, cleanType: {}, range:{}-{}, status:{}",
partition.getGraphName(),
partition.getId(),
cleanPartition.getCleanType(),
cleanPartition.getKeyStart(),
cleanPartition.getKeyEnd(),
status);
});
}
}
@Override
public void onPartitionKeyRangeChanged(long taskId, Partition partition, PartitionKeyRange partitionKeyRange,
Consumer<Integer> consumer) {
PartitionEngine engine = storeEngine.getPartitionEngine(partition.getId());
if (engine != null && engine.isLeader()){
consumer.accept(0);
var partitionManager = storeEngine.getPartitionManager();
var localPartition = partitionManager.getPartition(partition.getGraphName(), partition.getId());
if (localPartition == null) {
// 如果分区数据为空本地不会存储
localPartition = partitionManager.getPartitionFromPD(partition.getGraphName(), partition.getId());
LOG.info("onPartitionKeyRangeChanged, get from pd:{}-{} -> {}",
partition.getGraphName(), partition.getId(), localPartition);
if (localPartition == null){
return;
}
}
var newPartition = localPartition.getProtoObj().toBuilder()
.setStartKey(partitionKeyRange.getKeyStart())
.setEndKey(partitionKeyRange.getKeyEnd())
.setState(Metapb.PartitionState.PState_Normal)
.build();
partitionManager.updatePartition(newPartition, true);
try {
engine.addRaftTask(RaftOperation.create(RaftOperation.SYNC_PARTITION, newPartition), status -> {
LOG.info("onPartitionKeyRangeChanged, {}-{},key range: {}-{} status{}",
newPartition.getGraphName(),
newPartition.getId(),
partitionKeyRange.getKeyStart(),
partitionKeyRange.getKeyEnd(),
status);
});
LOG.info("onPartitionKeyRangeChanged: {}, update to pd", newPartition);
partitionManager.updatePartitionToPD(List.of(newPartition));
} catch (IOException e) {
LOG.error("Partition {}-{} onPartitionKeyRangeChanged exception {}",
newPartition.getGraphName(), newPartition.getId(), e);
} catch (PDException e) {
throw new RuntimeException(e);
}
}
}
/**
* is the task exists
* @param taskId task id
* @param partId partition id
* @return true if exists, false otherwise
*/
private boolean preCheckTaskId(long taskId, int partId){
if (storeEngine.getPartitionEngine(partId) == null ){
return false;
}
byte[] key = MetadataKeyHelper.getInstructionIdKey(taskId);
var wrapper = storeEngine.getPartitionManager().getWrapper();
byte[] value = wrapper.get(partId, key);
if (value != null) {
return true;
}
wrapper.put(partId, key, new byte[0]);
return false;
}
}

View File

@ -0,0 +1,135 @@
package com.baidu.hugegraph.store.options;
import com.alipay.sofa.jraft.storage.impl.RocksDBLogStorage;
import com.alipay.sofa.jraft.util.StorageOptionsFactory;
import com.baidu.hugegraph.rocksdb.access.RocksDBOptions;
import com.baidu.hugegraph.store.business.BusinessHandlerImpl;
import lombok.extern.slf4j.Slf4j;
import org.rocksdb.*;
import org.rocksdb.util.SizeUnit;
import org.apache.hugegraph.config.HugeConfig;
import java.util.Map;
@Slf4j
public class RaftRocksdbOptions {
static class RocksdbConfig{
private Env env;
private LRUCache blockCache;
private LRUCache writeCache;
private WriteBufferManager bufferManager;
private BlockBasedTableConfig tableConfig;
private long blockCacheCapacity;
private long writeCacheCapacity;
public Env getEnv(){ return env; }
public LRUCache getBlockCache() { return blockCache; }
public LRUCache getWriteCache() { return writeCache; }
public WriteBufferManager getBufferManager() { return bufferManager;}
public BlockBasedTableConfig getTableConfig(){ return tableConfig; }
public long getBlockCacheCapacity(){ return blockCacheCapacity; }
public long getWriteCacheCapacity(){ return writeCacheCapacity; }
public RocksdbConfig(HugeConfig options) {
RocksDB.loadLibrary();
this.env = Env.getDefault();
double writeBufferRatio = options.get(RocksDBOptions.WRITE_BUFFER_RATIO);
this.writeCacheCapacity = (long) (options.get(RocksDBOptions.TOTAL_MEMORY_SIZE) * writeBufferRatio);
this.blockCacheCapacity = options.get(RocksDBOptions.TOTAL_MEMORY_SIZE) - writeCacheCapacity;
this.writeCache = new LRUCache(writeCacheCapacity);
this.blockCache = new LRUCache(blockCacheCapacity);
this.bufferManager = new WriteBufferManager(writeCacheCapacity, writeCache,
options.get(RocksDBOptions.WRITE_BUFFER_ALLOW_STALL));
this.tableConfig = new BlockBasedTableConfig() //
.setIndexType(IndexType.kTwoLevelIndexSearch) //
.setPartitionFilters(true) //
.setMetadataBlockSize(8 * SizeUnit.KB) //
.setCacheIndexAndFilterBlocks(options.get(RocksDBOptions.PUT_FILTER_AND_INDEX_IN_CACHE)) //
.setCacheIndexAndFilterBlocksWithHighPriority(true) //
.setPinL0FilterAndIndexBlocksInCache(options.get(RocksDBOptions.PIN_L0_FILTER_AND_INDEX_IN_CACHE)) //
.setBlockSize(4 * SizeUnit.KB)//
.setBlockCache(blockCache);
int bitsPerKey = options.get(RocksDBOptions.BLOOM_FILTER_BITS_PER_KEY);
if (bitsPerKey >= 0) {
tableConfig.setFilterPolicy(new BloomFilter(bitsPerKey,
options.get(RocksDBOptions.BLOOM_FILTER_MODE)));
}
tableConfig.setWholeKeyFiltering(
options.get(RocksDBOptions.BLOOM_FILTER_WHOLE_KEY));
log.info("RocksdbConfig {}", options.get(RocksDBOptions.BLOOM_FILTER_BITS_PER_KEY));
}
}
private static RocksdbConfig rocksdbConfig = null;
private static RocksdbConfig getRocksdbConfig(HugeConfig options){
if ( rocksdbConfig == null){
synchronized (RocksdbConfig.class){
rocksdbConfig = new RocksdbConfig(options);
}
}
return rocksdbConfig;
}
private static void registerRaftRocksdbConfig(HugeConfig options) {
Cache blockCache = new LRUCache(1 * SizeUnit.GB);
BlockBasedTableConfig tableConfig = new BlockBasedTableConfig() //
.setIndexType(IndexType.kTwoLevelIndexSearch) //
.setPartitionFilters(true) //
.setMetadataBlockSize(8 * SizeUnit.KB) //
.setCacheIndexAndFilterBlocks(options.get(RocksDBOptions.PUT_FILTER_AND_INDEX_IN_CACHE)) //
.setCacheIndexAndFilterBlocksWithHighPriority(true) //
.setPinL0FilterAndIndexBlocksInCache(options.get(RocksDBOptions.PIN_L0_FILTER_AND_INDEX_IN_CACHE)) //
.setBlockSize(4 * SizeUnit.KB)//
.setBlockCache(blockCache);
StorageOptionsFactory.registerRocksDBTableFormatConfig(RocksDBLogStorage.class,
tableConfig);
DBOptions dbOptions = StorageOptionsFactory.getDefaultRocksDBOptions();
dbOptions.setEnv(rocksdbConfig.getEnv());
// raft rocksdb数量固定通过max_write_buffer_number可以控制
//dbOptions.setWriteBufferManager(rocksdbConfig.getBufferManager());
dbOptions.setUnorderedWrite(true);
StorageOptionsFactory.registerRocksDBOptions(RocksDBLogStorage.class,
dbOptions);
ColumnFamilyOptions cfOptions = StorageOptionsFactory.getDefaultRocksDBColumnFamilyOptions();
cfOptions.setTargetFileSizeBase(256 * SizeUnit.MB);
cfOptions.setWriteBufferSize(8 * SizeUnit.MB);
cfOptions.setNumLevels(3);
cfOptions.setMaxWriteBufferNumber(3);
cfOptions.setCompressionType(CompressionType.NO_COMPRESSION);
cfOptions.setMaxBytesForLevelBase(2048 * SizeUnit.GB);
StorageOptionsFactory.registerRocksDBColumnFamilyOptions(RocksDBLogStorage.class,
cfOptions);
}
public static void initRocksdbGlobalConfig(Map<String, Object> config){
HugeConfig hugeConfig = BusinessHandlerImpl.initRocksdb(config, null);
RocksdbConfig rocksdbConfig = getRocksdbConfig(hugeConfig);
registerRaftRocksdbConfig(hugeConfig);
config.put(RocksDBOptions.ENV, rocksdbConfig.getEnv());
config.put(RocksDBOptions.WRITE_BUFFER_MANAGER, rocksdbConfig.getBufferManager());
config.put(RocksDBOptions.BLOCK_TABLE_CONFIG, rocksdbConfig.getTableConfig());
config.put(RocksDBOptions.BLOCK_CACHE, rocksdbConfig.getBlockCache());
config.put(RocksDBOptions.WRITE_CACHE, rocksdbConfig.getWriteCache());
}
public static WriteBufferManager getWriteBufferManager(){
return rocksdbConfig.getBufferManager();
}
public static Env getEnv(){
return rocksdbConfig.getEnv();
}
public static Cache getWriteCache() { return rocksdbConfig.getWriteCache();}
public static Cache getBlockCache() { return rocksdbConfig.getBlockCache();}
public static long getWriteCacheCapacity() { return rocksdbConfig.getWriteCacheCapacity();}
public static long getBlockCacheCapacity() { return rocksdbConfig.getBlockCacheCapacity();}
}

View File

@ -25,6 +25,8 @@ import org.apache.hugegraph.rocksdb.access.ScanIterator;
import org.apache.hugegraph.structure.HugeElement;
import org.apache.hugegraph.util.Bytes;
import org.apache.tinkerpop.gremlin.structure.Edge;
import org.apache.hugegraph.iterator.CIter;
import org.apache.hugegraph.util.Bytes;
import lombok.extern.slf4j.Slf4j;

View File

@ -48,6 +48,7 @@ import org.apache.hugegraph.store.metric.HgMetricService;
import org.apache.hugegraph.store.util.Asserts;
import org.apache.hugegraph.util.Log;
import org.slf4j.Logger;
import org.apache.hugegraph.util.Log;
import lombok.extern.slf4j.Slf4j;

View File

@ -33,7 +33,7 @@
<dependency>
<groupId>org.apache.hugegraph</groupId>
<artifactId>hugegraph-common</artifactId>
<version>1.8.12</version>
<version>1.0.0</version>
<exclusions>
<exclusion>
<groupId>org.glassfish.jersey.inject</groupId>
@ -50,9 +50,8 @@
</exclusions>
</dependency>
<dependency>
<groupId>org.apache.hugegraph</groupId>
<groupId>com.baidu.hugegraph</groupId>
<artifactId>hg-store-common</artifactId>
<version>${revision}</version>
</dependency>
<dependency>
<groupId>org.rocksdb</groupId>
@ -76,14 +75,4 @@
<version>1.2.83</version>
</dependency>
</dependencies>
<distributionManagement>
<repository>
<id>Baidu_Local</id>
<url>http://maven.baidu-int.com/nexus/content/repositories/Baidu_Local</url>
</repository>
<snapshotRepository>
<id>Baidu_Local_Snapshots</id>
<url>http://maven.baidu-int.com/nexus/content/repositories/Baidu_Local_Snapshots</url>
</snapshotRepository>
</distributionManagement>
</project>

View File

@ -63,6 +63,9 @@ import org.rocksdb.Slice;
import org.rocksdb.Statistics;
import org.rocksdb.WriteBufferManager;
import org.rocksdb.WriteOptions;
import org.apache.hugegraph.config.HugeConfig;
import org.apache.hugegraph.util.Bytes;
import org.apache.hugegraph.util.E;
import lombok.extern.slf4j.Slf4j;

View File

@ -32,6 +32,7 @@ import org.rocksdb.RocksIterator;
import org.rocksdb.Slice;
import org.rocksdb.Snapshot;
import org.rocksdb.WriteBatch;
import org.apache.hugegraph.util.Bytes;
import lombok.extern.slf4j.Slf4j;

View File

@ -0,0 +1,82 @@
package com.baidu.hugegraph.rocksdb.access;
import org.apache.hugegraph.config.HugeConfig;
import org.apache.hugegraph.config.OptionSpace;
import java.util.HashMap;
import java.util.Map;
public class RocksDBFactoryTest {
// @BeforeClass
public static void init() {
OptionSpace.register("rocksdb",
"com.baidu.hugegraph.rocksdb.access.RocksDBOptions");
RocksDBOptions.instance();
Map<String, Object> configMap = new HashMap<>();
configMap.put("rocksdb.write_buffer_size", "1048576");
configMap.put("rocksdb.bloom_filter_bits_per_key", "10");
HugeConfig hConfig = new HugeConfig(configMap);
RocksDBFactory rFactory = RocksDBFactory.getInstance();
rFactory.setHugeConfig(hConfig);
}
// @Test
public void testCreateSession() throws InterruptedException {
RocksDBFactory factory = RocksDBFactory.getInstance();
try (RocksDBSession dbSession = factory.createGraphDB("./tmp", "test1")) {
SessionOperator op = dbSession.sessionOp();
op.prepare();
try {
op.put("tbl", "k1".getBytes(), "v1".getBytes());
op.commit();
}catch (Exception e){
op.rollback();
}
}
factory.destroyGraphDB("test1");
Thread.sleep(100000);
}
// @Test
public void testTotalKeys(){
RocksDBFactory dbFactory = RocksDBFactory.getInstance();
System.out.println(dbFactory.getTotalSize());
System.out.println( dbFactory.getTotalKey().entrySet()
.stream().map(e->e.getValue()).reduce(0L, Long::sum));
}
// @Test
public void releaseAllGraphDB() {
System.out.println(RocksDBFactory.class);
RocksDBFactory rFactory = RocksDBFactory.getInstance();
if(rFactory.queryGraphDB("bj01") == null) {
rFactory.createGraphDB("./tmp", "bj01");
}
if(rFactory.queryGraphDB("bj02") == null) {
rFactory.createGraphDB("./tmp", "bj02");
}
if(rFactory.queryGraphDB("bj03") == null) {
rFactory.createGraphDB("./tmp", "bj03");
}
RocksDBSession dbSession = rFactory.queryGraphDB("bj01");
dbSession.checkTable("test");
SessionOperator sessionOp = dbSession.sessionOp();
sessionOp.prepare();
sessionOp.put("test", "hi".getBytes(), "byebye".getBytes());
sessionOp.commit();
rFactory.releaseAllGraphDB();
}
}

View File

@ -23,10 +23,15 @@ import java.util.ArrayList;
import java.util.List;
import org.apache.hugegraph.store.term.HgPair;
import org.apache.commons.io.FileUtils;
import org.junit.AfterClass;
import org.junit.Assert;
import org.apache.hugegraph.config.HugeConfig;
import org.apache.hugegraph.config.OptionSpace;
import com.alibaba.fastjson.JSON;
import lombok.extern.slf4j.Slf4j;
// import org.junit.BeforeClass;
// import org.junit.Test;
@Slf4j
public class SnapshotManagerTest {