diff --git a/hg-store-core/src/main/java/com/baidu/hugegraph/store/PartitionInstructionProcessor.java b/hg-store-core/src/main/java/com/baidu/hugegraph/store/PartitionInstructionProcessor.java deleted file mode 100644 index b6a5080f3..000000000 --- a/hg-store-core/src/main/java/com/baidu/hugegraph/store/PartitionInstructionProcessor.java +++ /dev/null @@ -1,306 +0,0 @@ -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 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 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 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 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 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 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 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; - } -} diff --git a/hg-store-core/src/main/java/com/baidu/hugegraph/store/options/RaftRocksdbOptions.java b/hg-store-core/src/main/java/com/baidu/hugegraph/store/options/RaftRocksdbOptions.java deleted file mode 100644 index 3a5a17451..000000000 --- a/hg-store-core/src/main/java/com/baidu/hugegraph/store/options/RaftRocksdbOptions.java +++ /dev/null @@ -1,135 +0,0 @@ -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 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();} -} diff --git a/hg-store-core/src/main/java/org/apache/hugegraph/store/business/AbstractSelectIterator.java b/hg-store-core/src/main/java/org/apache/hugegraph/store/business/AbstractSelectIterator.java index 34ce98637..da9bb5f8e 100644 --- a/hg-store-core/src/main/java/org/apache/hugegraph/store/business/AbstractSelectIterator.java +++ b/hg-store-core/src/main/java/org/apache/hugegraph/store/business/AbstractSelectIterator.java @@ -25,8 +25,6 @@ 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; diff --git a/hg-store-core/src/main/java/org/apache/hugegraph/store/pd/DefaultPdProvider.java b/hg-store-core/src/main/java/org/apache/hugegraph/store/pd/DefaultPdProvider.java index f6a5b7466..35c67980b 100644 --- a/hg-store-core/src/main/java/org/apache/hugegraph/store/pd/DefaultPdProvider.java +++ b/hg-store-core/src/main/java/org/apache/hugegraph/store/pd/DefaultPdProvider.java @@ -48,7 +48,6 @@ 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; diff --git a/hg-store-rocksdb/pom.xml b/hg-store-rocksdb/pom.xml index 596ccdc99..b4dfae9e4 100644 --- a/hg-store-rocksdb/pom.xml +++ b/hg-store-rocksdb/pom.xml @@ -50,7 +50,7 @@ - com.baidu.hugegraph + org.apache.hugegraph hg-store-common diff --git a/hg-store-rocksdb/src/test/java/com/baidu/hugegraph/rocksdb/access/RocksDBFactoryTest.java b/hg-store-rocksdb/src/test/java/com/baidu/hugegraph/rocksdb/access/RocksDBFactoryTest.java deleted file mode 100644 index cd89a881b..000000000 --- a/hg-store-rocksdb/src/test/java/com/baidu/hugegraph/rocksdb/access/RocksDBFactoryTest.java +++ /dev/null @@ -1,82 +0,0 @@ -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 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(); - } -} \ No newline at end of file