diff --git a/hg-store-client/src/main/java/org/apache/hugegraph/store/HgKvStore.java b/hg-store-client/src/main/java/org/apache/hugegraph/store/HgKvStore.java index fd1e48367..ae79ed828 100644 --- a/hg-store-client/src/main/java/org/apache/hugegraph/store/HgKvStore.java +++ b/hg-store-client/src/main/java/org/apache/hugegraph/store/HgKvStore.java @@ -97,6 +97,8 @@ public interface HgKvStore { HgKvIterator scanIterator(ScanStreamReq.Builder scanReqBuilder); + long count(String table); + boolean truncate(); default boolean existsTable(String table) { diff --git a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/NodeTxSessionProxy.java b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/NodeTxSessionProxy.java index dc850c37d..3379d6653 100644 --- a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/NodeTxSessionProxy.java +++ b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/NodeTxSessionProxy.java @@ -505,6 +505,17 @@ class NodeTxSessionProxy implements HgStoreSession { return this.toHgKvIteratorProxy(iterators, scanReqBuilder.getLimit()); } + @Override + public long count(String table) { + return this.toNodeTkvList(table) + .parallelStream() + .map( + e -> this.getStoreNode(e.getNodeId()).openSession(this.graphName) + .count(e.getTable()) + ) + .collect(Collectors.summingLong(l -> l)); + } + @Override public List> scanBatch(HgScanQuery scanQuery) { HgAssert.isArgumentNotNull(scanQuery, "scanQuery"); diff --git a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/AbstractGrpcClient.java b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/AbstractGrpcClient.java index 5a9df45a0..b29b5c837 100644 --- a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/AbstractGrpcClient.java +++ b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/AbstractGrpcClient.java @@ -19,10 +19,13 @@ package org.apache.hugegraph.store.client.grpc; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; import java.util.stream.IntStream; +import org.apache.hugegraph.store.client.util.ExecutorPool; import org.apache.hugegraph.store.client.util.HgStoreClientConfig; import org.apache.hugegraph.store.term.HgPair; @@ -30,22 +33,27 @@ import io.grpc.ManagedChannel; import io.grpc.ManagedChannelBuilder; import io.grpc.stub.AbstractAsyncStub; import io.grpc.stub.AbstractBlockingStub; +import io.grpc.stub.AbstractStub; /** * @date 2023/3/28 **/ public abstract class AbstractGrpcClient { - private static final Map channels = new ConcurrentHashMap<>(); - private static final int n = 5; - private static final int concurrency = 1 << n; - private static final AtomicLong counter = new AtomicLong(0); - private static final long limit = Long.MAX_VALUE >> 1; - private static final HgStoreClientConfig config = HgStoreClientConfig.of(); - private final Map[]> blockingStubs = - new ConcurrentHashMap<>(); - private final Map[]> asyncStubs = + private static Map channels = new ConcurrentHashMap<>(); + private static int n = 5; + private static int concurrency = 1 << n; + private static AtomicLong counter = new AtomicLong(0); + private static long limit = Long.MAX_VALUE >> 1; + private Map[]> blockingStubs = new ConcurrentHashMap<>(); + private Map[]> asyncStubs = new ConcurrentHashMap<>(); + private static HgStoreClientConfig config = HgStoreClientConfig.of(); + private ThreadPoolExecutor executor; + + { + executor = ExecutorPool.createExecutor("common", 60, concurrency, concurrency); + } public AbstractGrpcClient() { @@ -56,10 +64,26 @@ public abstract class AbstractGrpcClient { if ((tc = channels.get(target)) == null) { synchronized (channels) { if ((tc = channels.get(target)) == null) { - ManagedChannel[] value = new ManagedChannel[concurrency]; - IntStream.range(0, concurrency).parallel() - .forEach(i -> value[i] = getManagedChannel(target)); - channels.put(target, tc = value); + try { + ManagedChannel[] value = new ManagedChannel[concurrency]; + CountDownLatch latch = new CountDownLatch(concurrency); + for (int i = 0; i < concurrency; i++) { + int fi = i; + executor.execute(() -> { + try{ + value[fi] = getManagedChannel(target); + } catch (Exception e) { + throw new RuntimeException(e); + } finally { + latch.countDown(); + } + }); + } + latch.await(); + channels.put(target, tc = value); + } catch (Exception e) { + throw new RuntimeException(e); + } } } } @@ -81,25 +105,27 @@ public abstract class AbstractGrpcClient { pairs = blockingStubs.get(target); if (pairs == null) { HgPair[] value = new HgPair[concurrency]; - IntStream.range(0, concurrency).parallel().forEach(i -> { + IntStream.range(0, concurrency).forEach(i -> { ManagedChannel channel = channels[index]; AbstractBlockingStub stub = getBlockingStub(channel); - stub.withMaxInboundMessageSize(config.getGrpcMaxInboundMessageSize()) - .withMaxOutboundMessageSize(config.getGrpcMaxOutboundMessageSize()); value[i] = new HgPair<>(channel, stub); // log.info("create channel for {}",target); }); blockingStubs.put(target, value); AbstractBlockingStub stub = value[index].getValue(); - return (AbstractBlockingStub) stub.withDeadlineAfter( - config.getGrpcTimeoutSeconds(), - TimeUnit.SECONDS); + return (AbstractBlockingStub) setBlockingStubOption(stub); } } } - return (AbstractBlockingStub) pairs[index].getValue() - .withDeadlineAfter(config.getGrpcTimeoutSeconds(), - TimeUnit.SECONDS); + return (AbstractBlockingStub) setBlockingStubOption(pairs[index].getValue()); + } + + private AbstractStub setBlockingStubOption(AbstractBlockingStub stub) { + return stub.withDeadlineAfter(config.getGrpcTimeoutSeconds(),TimeUnit.SECONDS) + .withMaxInboundMessageSize( + config.getGrpcMaxInboundMessageSize()) + .withMaxOutboundMessageSize( + config.getGrpcMaxOutboundMessageSize()); } public AbstractAsyncStub getAsyncStub(ManagedChannel channel) { @@ -122,21 +148,28 @@ public abstract class AbstractGrpcClient { IntStream.range(0, concurrency).parallel().forEach(i -> { ManagedChannel channel = channels[index]; AbstractAsyncStub stub = getAsyncStub(channel); - stub.withMaxInboundMessageSize(config.getGrpcMaxInboundMessageSize()) - .withMaxOutboundMessageSize(config.getGrpcMaxOutboundMessageSize()); + // stub.withMaxInboundMessageSize(config.getGrpcMaxInboundMessageSize()) + // .withMaxOutboundMessageSize(config.getGrpcMaxOutboundMessageSize()); value[i] = new HgPair<>(channel, stub); // log.info("create channel for {}",target); }); asyncStubs.put(target, value); - AbstractAsyncStub stub = value[index].getValue(); + AbstractAsyncStub stub = + (AbstractAsyncStub) setStubOption(value[index].getValue()); return stub; } } } - return pairs[index].getValue(); + return (AbstractAsyncStub) setStubOption(pairs[index].getValue()); } + private AbstractStub setStubOption(AbstractStub value) { + return value.withMaxInboundMessageSize( + config.getGrpcMaxInboundMessageSize()) + .withMaxOutboundMessageSize( + config.getGrpcMaxOutboundMessageSize()); + } private ManagedChannel getManagedChannel(String target) { return ManagedChannelBuilder.forTarget(target).usePlaintext().build(); diff --git a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/GrpcStoreNodeSessionImpl.java b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/GrpcStoreNodeSessionImpl.java index a0eee04cb..944742ace 100644 --- a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/GrpcStoreNodeSessionImpl.java +++ b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/GrpcStoreNodeSessionImpl.java @@ -368,6 +368,11 @@ class GrpcStoreNodeSessionImpl implements HgStoreNodeSession { return GrpcKvIteratorImpl.of(this, scanner); } + @Override + public long count(String table) { + return this.storeSessionClient.count(this, table); + } + @Override public HgKvIterator scanIterator(String table, byte[] query) { diff --git a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/GrpcStoreSessionClient.java b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/GrpcStoreSessionClient.java index 246c80eb8..af69add87 100644 --- a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/GrpcStoreSessionClient.java +++ b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/grpc/GrpcStoreSessionClient.java @@ -18,6 +18,7 @@ package org.apache.hugegraph.store.client.grpc; import java.util.List; +import java.util.concurrent.TimeUnit; import javax.annotation.concurrent.ThreadSafe; @@ -37,6 +38,7 @@ import org.apache.hugegraph.store.grpc.session.HgStoreSessionGrpc; import org.apache.hugegraph.store.grpc.session.HgStoreSessionGrpc.HgStoreSessionBlockingStub; import org.apache.hugegraph.store.grpc.session.TableReq; +import io.grpc.Deadline; import io.grpc.ManagedChannel; import lombok.extern.slf4j.Slf4j; @@ -137,6 +139,17 @@ class GrpcStoreSessionClient extends AbstractGrpcClient { .build() ); } + + public long count(HgStoreNodeSession nodeSession, String table) { + Agg agg = this.getBlockingStub(nodeSession).withDeadline(Deadline.after(24, TimeUnit.HOURS)) + .count(ScanStreamReq.newBuilder() + .setHeader(getHeader(nodeSession)) + .setTable(table) + .setMethod(ScanMethod.ALL) + .build() + ); + return agg.getCount(); + } } diff --git a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/util/ExecutorPool.java b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/util/ExecutorPool.java index be28245f8..f0aa85f48 100644 --- a/hg-store-client/src/main/java/org/apache/hugegraph/store/client/util/ExecutorPool.java +++ b/hg-store-client/src/main/java/org/apache/hugegraph/store/client/util/ExecutorPool.java @@ -25,20 +25,12 @@ import java.util.concurrent.atomic.AtomicInteger; import lombok.extern.slf4j.Slf4j; -/** - * 2021/11/22 - */ @Slf4j public final class ExecutorPool { - public static ThreadFactory newThreadFactory(String namePrefix, int priority) { - HgAssert.isArgumentNotNull(namePrefix, "namePrefix"); - return new HgThreadFactory(namePrefix, priority); - } - public static ThreadFactory newThreadFactory(String namePrefix) { HgAssert.isArgumentNotNull(namePrefix, "namePrefix"); - return new HgDefaultThreadFactory(namePrefix); + return new DefaultThreadFactory(namePrefix); } public static ThreadPoolExecutor createExecutor(String name, long keepAliveTime, @@ -50,97 +42,20 @@ public final class ExecutorPool { ); } - // - // private final static ExecutorService executor = new ThreadPoolExecutor(0, Integer.MAX_VALUE, - // 10L, TimeUnit.SECONDS, - // new - // LinkedBlockingQueue<>(), - // new - // HgDefaultThreadFactory - // ("store-common")); - // private final static ExecutorService grpcExecutor = createExecutor("store-grpc", 60L, 600, - // 10240); - //// public final static ExecutorService scannerExecutor = createExecutor("scanner", 10l, - // 200, 1024); - // - // static { - // Thread hook = new Thread(() -> { - // try { - // executor.shutdown(); - // executor.awaitTermination(30, TimeUnit.SECONDS); - // return; - // } catch (InterruptedException e) { - // log.error("failed to await executorService.shutdown()", e); - // } - // executor.shutdownNow().forEach(r -> log.error("not run task:" + r.toString())); - // try { - // executor.awaitTermination(30, TimeUnit.SECONDS); - // } catch (InterruptedException e) { - // log.error("failed to await executorService.shutdownNow()", e); - // throw HgStoreClientException.of(e); - // } - // } - // ); - // Runtime.getRuntime().addShutdownHook(hook); - //} - // - //// public static void execute(Runnable command) { - //// isArgumentNotNull(command, "command"); - //// executor.execute(command); - //// } - //// - // public static ExecutorService getGrpcCommonExecutor() { - // return grpcExecutor; - // } - // + public static class DefaultThreadFactory implements ThreadFactory { - /** - * The default thread factory - */ - static class HgThreadFactory implements ThreadFactory { - private final AtomicInteger threadNumber = new AtomicInteger(1); - private final String namePrefix; - private final int priority; - - HgThreadFactory(String namePrefix, int priority) { - this.namePrefix = namePrefix; - this.priority = priority; - } - - @Override - public Thread newThread(Runnable r) { - Thread t = new Thread(null, r, namePrefix + "-" + threadNumber.getAndIncrement(), 0); - if (t.isDaemon()) { - t.setDaemon(false); - } - if (t.getPriority() != priority) { - t.setPriority(priority); - } - return t; - } - } - - /** - * The default thread factory, which added threadNamePrefix in construction method. - */ - static class HgDefaultThreadFactory implements ThreadFactory { - private static final AtomicInteger poolNumber = new AtomicInteger(1); private final AtomicInteger threadNumber = new AtomicInteger(1); private final String namePrefix; - HgDefaultThreadFactory(String threadNamePrefix) { - this.namePrefix = threadNamePrefix + "-" + poolNumber.getAndIncrement() + "-thread-"; + public DefaultThreadFactory(String threadNamePrefix) { + this.namePrefix = threadNamePrefix + "-"; } @Override public Thread newThread(Runnable r) { Thread t = new Thread(null, r, namePrefix + threadNumber.getAndIncrement(), 0); - if (t.isDaemon()) { - t.setDaemon(false); - } - if (t.getPriority() != Thread.NORM_PRIORITY) { - t.setPriority(Thread.NORM_PRIORITY); - } + t.setDaemon(true); + t.setPriority(Thread.NORM_PRIORITY); return t; } } diff --git a/hg-store-core/src/main/java/org/apache/hugegraph/store/HgStoreEngine.java b/hg-store-core/src/main/java/org/apache/hugegraph/store/HgStoreEngine.java index e8755d532..ff1a3e9ad 100644 --- a/hg-store-core/src/main/java/org/apache/hugegraph/store/HgStoreEngine.java +++ b/hg-store-core/src/main/java/org/apache/hugegraph/store/HgStoreEngine.java @@ -80,9 +80,7 @@ public class HgStoreEngine implements Lifecycle, HgStoreSt private HgMetricService metricService; private DataMover dataMover; - private HgStoreEngine() { - - } + private static ConcurrentHashMap engineLocks = new ConcurrentHashMap<>(); public static HgStoreEngine getInstance() { return instance; @@ -290,7 +288,8 @@ public class HgStoreEngine implements Lifecycle, HgStoreSt Configuration conf) { PartitionEngine engine; if ((engine = partitionEngines.get(groupId)) == null) { - synchronized (this) { + engineLocks.computeIfAbsent(groupId, k -> new Object()); + synchronized (engineLocks.get(groupId)) { // 分区分裂时特殊情况(集群中图分区数量不一样),会导致分裂的分区,可能不在本机器上. if (conf != null) { var list = conf.listPeers(); @@ -566,7 +565,8 @@ public class HgStoreEngine implements Lifecycle, HgStoreSt RaftClosure closure) { PartitionEngine engine = getPartitionEngine(graphName, partId); if (engine == null) { - synchronized (this) { + engineLocks.computeIfAbsent(partId, k -> new Object()); + synchronized (engineLocks.get(partId)) { engine = getPartitionEngine(graphName, partId); if (engine == null) { Partition partition = partitionManager.findPartition(graphName, partId); diff --git a/hg-store-core/src/main/java/org/apache/hugegraph/store/business/BusinessHandler.java b/hg-store-core/src/main/java/org/apache/hugegraph/store/business/BusinessHandler.java index 2ef0c8c2e..be608f2a9 100644 --- a/hg-store-core/src/main/java/org/apache/hugegraph/store/business/BusinessHandler.java +++ b/hg-store-core/src/main/java/org/apache/hugegraph/store/business/BusinessHandler.java @@ -22,8 +22,6 @@ import java.util.Map; import java.util.function.Consumer; import java.util.function.Supplier; -import javax.annotation.concurrent.NotThreadSafe; - import org.apache.hugegraph.pd.grpc.pulse.CleanType; import org.apache.hugegraph.rocksdb.access.ScanIterator; import org.apache.hugegraph.store.grpc.Graphpb; @@ -198,28 +196,5 @@ public interface BusinessHandler extends DBSessionBuilder { void destroyGraphDB(String graphName, int partId) throws HgStoreException; - @NotThreadSafe - interface TxBuilder { - TxBuilder put(int code, String table, byte[] key, byte[] value) throws HgStoreException; - - TxBuilder del(int code, String table, byte[] key) throws HgStoreException; - - TxBuilder delSingle(int code, String table, byte[] key) throws HgStoreException; - - TxBuilder delPrefix(int code, String table, byte[] prefix) throws HgStoreException; - - TxBuilder delRange(int code, String table, byte[] start, byte[] end) throws - HgStoreException; - - TxBuilder merge(int code, String table, byte[] key, byte[] value) throws HgStoreException; - - Tx build(); - } - - interface Tx { - void commit() throws HgStoreException; - - void rollback() throws HgStoreException; - } - + long count(String graphName, String table); } diff --git a/hg-store-core/src/main/java/org/apache/hugegraph/store/business/BusinessHandlerImpl.java b/hg-store-core/src/main/java/org/apache/hugegraph/store/business/BusinessHandlerImpl.java index 8172171e5..8430b62e2 100644 --- a/hg-store-core/src/main/java/org/apache/hugegraph/store/business/BusinessHandlerImpl.java +++ b/hg-store-core/src/main/java/org/apache/hugegraph/store/business/BusinessHandlerImpl.java @@ -33,8 +33,6 @@ import java.util.function.Function; import java.util.function.Supplier; import java.util.stream.Collectors; -import javax.annotation.concurrent.NotThreadSafe; - import org.apache.commons.lang.ArrayUtils; import org.apache.commons.lang.StringUtils; import org.apache.hugegraph.config.HugeConfig; @@ -802,117 +800,35 @@ public class BusinessHandlerImpl implements BusinessHandler { keyCreator.clearCache(partId); } - @NotThreadSafe - private class TxBuilderImpl implements TxBuilder { - private final String graph; - private final int partId; - private final RocksDBSession dbSession; - private final SessionOperator op; - - private TxBuilderImpl(String graph, int partId, RocksDBSession dbSession) { - this.graph = graph; - this.partId = partId; - this.dbSession = dbSession; - this.op = this.dbSession.sessionOp(); - this.op.prepare(); - } - - @Override - public TxBuilder put(int code, String table, byte[] key, byte[] value) throws - HgStoreException { - try { - byte[] targetKey = keyCreator.getKey(this.partId, graph, code, key); - this.op.put(table, targetKey, value); - } catch (DBStoreException e) { - throw new HgStoreException(HgStoreException.EC_RKDB_DOPUT_FAIL, e.toString()); - } - return this; - } - - @Override - public TxBuilder del(int code, String table, byte[] key) throws HgStoreException { - try { - byte[] targetKey = keyCreator.getKey(this.partId, graph, code, key); - this.op.delete(table, targetKey); - } catch (DBStoreException e) { - throw new HgStoreException(HgStoreException.EC_RKDB_DODEL_FAIL, e.toString()); - } - - return this; - } - - @Override - public TxBuilder delSingle(int code, String table, byte[] key) throws HgStoreException { - - try { - byte[] targetKey = keyCreator.getKey(this.partId, graph, code, key); - op.deleteSingle(table, targetKey); - } catch (DBStoreException e) { - throw new HgStoreException(HgStoreException.EC_RDKDB_DOSINGLEDEL_FAIL, - e.toString()); - } - - return this; - } - - @Override - public TxBuilder delPrefix(int code, String table, byte[] prefix) throws HgStoreException { - - try { - this.op.deletePrefix(table, keyCreator.getPrefixKey(this.partId, graph, prefix)); - } catch (DBStoreException e) { - throw new HgStoreException(HgStoreException.EC_RKDB_DODELPREFIX_FAIL, e.toString()); - } - - return this; - } - - @Override - public TxBuilder delRange(int code, String table, byte[] start, byte[] end) throws - HgStoreException { - - try { - this.op.deleteRange(table, keyCreator.getStartKey(this.partId, graph, start), - keyCreator.getEndKey(this.partId, graph, end)); - } catch (DBStoreException e) { - throw new HgStoreException(HgStoreException.EC_RKDB_DODELRANGE_FAIL, e.toString()); - } - - return this; - } - - @Override - public TxBuilder merge(int code, String table, byte[] key, byte[] value) throws - HgStoreException { - - try { - byte[] targetKey = keyCreator.getKey(this.partId, graph, code, key); - op.merge(table, targetKey, value); - } catch (DBStoreException e) { - throw new HgStoreException(HgStoreException.EC_RKDB_DOMERGE_FAIL, e.toString()); - } - - return this; - } - - @Override - public Tx build() { - return new Tx() { - @Override - public void commit() throws HgStoreException { - op.commit(); // commit发生异常后,必须调用rollback,否则造成锁未释放 - dbSession.close(); + @Override + public long count(String graph, String table) { + List ids = this.getLeaderPartitionIds(graph); + Long all = ids.parallelStream().map((id) -> { + InnerKeyFilter it = null; + try (RocksDBSession dbSession = getSession(graph, table, id)) { + long count = 0; + SessionOperator op = dbSession.sessionOp(); + it = new InnerKeyFilter(op.scan(table, + keyCreator.getStartKey(id, graph), + keyCreator.getEndKey(id, graph), + ScanIterator.Trait.SCAN_LT_END)); + while (it.hasNext()) { + it.next(); + count++; } - - @Override - public void rollback() throws HgStoreException { + return count; + } catch (Exception e) { + throw e; + } finally { + if (it != null) { try { - op.rollback(); - } finally { - dbSession.close(); + it.close(); + } catch (Exception e) { + } } - }; - } + } + }).collect(Collectors.summingLong(l -> l)); + return all; } } 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 35c67980b..b4d8c4c43 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 @@ -21,8 +21,6 @@ package org.apache.hugegraph.store.pd; import java.util.ArrayList; import java.util.Collections; import java.util.List; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; import java.util.function.Consumer; import org.apache.hugegraph.pd.client.PDClient; @@ -54,12 +52,14 @@ import lombok.extern.slf4j.Slf4j; @Slf4j public class DefaultPdProvider implements PdProvider { private static final Logger LOG = Log.logger(DefaultPdProvider.class); - private static final ConcurrentMap quotas = - new ConcurrentHashMap<>(); private final PDClient pdClient; private final String pdServerAddress; + + private Consumer hbOnError = null; private List partitionCommandListeners = - Collections.synchronizedList(new ArrayList()); + Collections.synchronizedList(new ArrayList<>()); + private final PDPulse pulseClient; + private PDPulse.Notifier pdPulse; private GraphManager graphManager = null; PDClient.PDEventListener listener = new PDClient.PDEventListener() { @@ -70,6 +70,11 @@ public class DefaultPdProvider implements PdProvider { log.info("store raft group changed!, {}", event); pdClient.invalidStoreCache(event.getNodeId()); HgStoreEngine.getInstance().rebuildRaftGroup(event.getNodeId()); + } else if (event.getEventType() == NodeEvent.EventType.NODE_PD_LEADER_CHANGE) { + log.info("pd leader changed!, {}. restart heart beat", event); + if (pulseClient.resetStub(event.getGraph(), pdPulse)) { + startHeartbeatStream(hbOnError); + } } } @@ -94,6 +99,8 @@ public class DefaultPdProvider implements PdProvider { this.pdClient.addEventListener(listener); this.pdServerAddress = pdAddress; partitionCommandListeners = Collections.synchronizedList(new ArrayList()); + log.info("pulse client connect to {}", pdClient.getLeaderIp()); + this.pulseClient = new PDPulseImpl(pdClient.getLeaderIp()); } @Override @@ -244,61 +251,77 @@ public class DefaultPdProvider implements PdProvider { */ @Override public boolean startHeartbeatStream(Consumer onError) { - PDPulse pulse = pdClient.getPulseClient(); - pdPulse = pulse.connectPartition(new PDPulse.Listener() { + this.hbOnError = onError; + pdPulse = pulseClient.connectPartition(new PDPulse.Listener<>() { @Override - public void onNotice(PulseServerNotice response) { - PartitionHeartbeatResponse instruction = response.getContent(); - LOG.debug("Partition heartbeat receive instruction: {}", instruction); - Partition partition = new Partition(instruction.getPartition()); + public void onNotice(PulseServerNotice response) { + PulseResponse content = response.getContent(); // 消息消费应答,能够正确消费消息,调用accept返回状态码,否则不要调用accept - Consumer consumer = new Consumer<>() { - @Override - public void accept(Integer integer) { - LOG.debug("Partition heartbeat accept instruction: {}", instruction); - // LOG.info("accept notice id : {}, ts:{}", response.getNoticeId(), - // System.currentTimeMillis()); - // http2 并发问题,需要加锁 - // synchronized (pdPulse) { - response.ack(); - // } - } + Consumer consumer = integer -> { + LOG.debug("Partition heartbeat accept instruction: {}", content); + // LOG.info("accept notice id : {}, ts:{}", response.getNoticeId(), System + // .currentTimeMillis()); + // http2 并发问题,需要加锁 + // synchronized (pdPulse) { + response.ack(); + // } }; + if (content.hasInstructionResponse()) { + var pdInstruction = content.getInstructionResponse(); + consumer.accept(0); + // 当前的链接变成了follower,重新链接 + if (pdInstruction.getInstructionType() == + PdInstructionType.CHANGE_TO_FOLLOWER) { + onCompleted(); + log.info("got pulse instruction, change leader to {}", + pdInstruction.getLeaderIp()); + if (pulseClient.resetStub(pdInstruction.getLeaderIp(), pdPulse)) { + startHeartbeatStream(hbOnError); + } + } + return; + } + + PartitionHeartbeatResponse instruct = content.getPartitionHeartbeatResponse(); + LOG.debug("Partition heartbeat receive instruction: {}", instruct); + + Partition partition = new Partition(instruct.getPartition()); + for (PartitionInstructionListener event : partitionCommandListeners) { - if (instruction.hasChangeShard()) { - event.onChangeShard(instruction.getId(), partition, - instruction.getChangeShard(), consumer); + if (instruct.hasChangeShard()) { + event.onChangeShard(instruct.getId(), partition, instruct.getChangeShard(), + consumer); } - if (instruction.hasSplitPartition()) { - event.onSplitPartition(instruction.getId(), partition, - instruction.getSplitPartition(), consumer); + if (instruct.hasSplitPartition()) { + event.onSplitPartition(instruct.getId(), partition, + instruct.getSplitPartition(), consumer); } - if (instruction.hasTransferLeader()) { - event.onTransferLeader(instruction.getId(), partition, - instruction.getTransferLeader(), consumer); + if (instruct.hasTransferLeader()) { + event.onTransferLeader(instruct.getId(), partition, + instruct.getTransferLeader(), consumer); } - if (instruction.hasDbCompaction()) { - event.onDbCompaction(instruction.getId(), partition, - instruction.getDbCompaction(), consumer); + if (instruct.hasDbCompaction()) { + event.onDbCompaction(instruct.getId(), partition, + instruct.getDbCompaction(), consumer); } - if (instruction.hasMovePartition()) { - event.onMovePartition(instruction.getId(), partition, - instruction.getMovePartition(), consumer); + if (instruct.hasMovePartition()) { + event.onMovePartition(instruct.getId(), partition, + instruct.getMovePartition(), consumer); } - if (instruction.hasCleanPartition()) { - event.onCleanPartition(instruction.getId(), partition, - instruction.getCleanPartition(), + if (instruct.hasCleanPartition()) { + event.onCleanPartition(instruct.getId(), partition, + instruct.getCleanPartition(), consumer); } - if (instruction.hasKeyRange()) { - event.onPartitionKeyRangeChanged(instruction.getId(), partition, - instruction.getKeyRange(), + if (instruct.hasKeyRange()) { + event.onPartitionKeyRangeChanged(instruct.getId(), partition, + instruct.getKeyRange(), consumer); } } @@ -307,6 +330,7 @@ public class DefaultPdProvider implements PdProvider { @Override public void onError(Throwable throwable) { LOG.error("Partition heartbeat stream error. {}", throwable); + pulseClient.resetStub(pdClient.getLeaderIp(), pdPulse); onError.accept(throwable); } diff --git a/hg-store-core/src/test/resources/version.txt b/hg-store-core/src/test/resources/version.txt index 77a069e39..b55f10804 100644 --- a/hg-store-core/src/test/resources/version.txt +++ b/hg-store-core/src/test/resources/version.txt @@ -1 +1 @@ -3.6.2 \ No newline at end of file +3.6.5 \ No newline at end of file diff --git a/hg-store-grpc/src/main/proto/store_session.proto b/hg-store-grpc/src/main/proto/store_session.proto index 62dcdd2ad..3110a9b4b 100644 --- a/hg-store-grpc/src/main/proto/store_session.proto +++ b/hg-store-grpc/src/main/proto/store_session.proto @@ -5,6 +5,7 @@ option java_package = "org.apache.hugegraph.store.grpc.session"; option java_outer_classname = "HgStoreSessionProto"; import "store_common.proto"; +import "store_stream_meta.proto"; service HgStoreSession { rpc Get2(GetReq) returns (FeedbackRes) {} @@ -13,6 +14,7 @@ service HgStoreSession { rpc Table(TableReq) returns (FeedbackRes){}; rpc Graph(GraphReq) returns (FeedbackRes){}; rpc Clean(CleanReq) returns (FeedbackRes) {} + rpc Count(ScanStreamReq) returns (Agg) {} } message TableReq{ @@ -112,5 +114,8 @@ enum PartitionFaultType{ PARTITION_FAULT_TYPE_NOT_LOCAL = 3; } - +message Agg { + Header header = 1; + int64 count = 2; +} diff --git a/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/HgStoreSessionImpl.java b/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/HgStoreSessionImpl.java index f8dd483fe..634ccb1cf 100644 --- a/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/HgStoreSessionImpl.java +++ b/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/HgStoreSessionImpl.java @@ -26,6 +26,8 @@ import java.util.concurrent.atomic.AtomicInteger; import org.apache.hugegraph.pd.common.PDException; import org.apache.hugegraph.pd.grpc.Metapb; import org.apache.hugegraph.pd.grpc.Metapb.GraphMode; +import org.apache.hugegraph.rocksdb.access.ScanIterator; +import org.apache.hugegraph.store.business.BusinessHandler; import org.apache.hugegraph.store.grpc.common.Key; import org.apache.hugegraph.store.grpc.common.Kv; import org.apache.hugegraph.store.grpc.common.ResCode; @@ -523,4 +525,25 @@ public class HgStoreSessionImpl extends HgStoreSessionGrpc.HgStoreSessionImplBas } GrpcClosure.setResult(response, builder.build()); } + + @Override + public void count(ScanStreamReq request, StreamObserver observer) { + ScanIterator it = null; + try { + BusinessHandler handler = storeService.getStoreEngine().getBusinessHandler(); + long count = handler.count(request.getHeader().getGraph(), request.getTable()); + observer.onNext(Agg.newBuilder().setCount(count).build()); + observer.onCompleted(); + } catch (Exception e) { + observer.onError(e); + } finally { + if (it != null) { + try { + it.close(); + } catch (Exception e) { + + } + } + } + } } diff --git a/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/ScanBatchResponse.java b/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/ScanBatchResponse.java index 4e03cbe29..220931e63 100644 --- a/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/ScanBatchResponse.java +++ b/hg-store-node/src/main/java/org/apache/hugegraph/store/node/grpc/ScanBatchResponse.java @@ -20,7 +20,6 @@ package org.apache.hugegraph.store.node.grpc; import java.util.List; import java.util.concurrent.ThreadPoolExecutor; -import java.util.concurrent.locks.ReentrantLock; import org.apache.hugegraph.rocksdb.access.ScanIterator; import org.apache.hugegraph.store.buffer.ByteBufferAllocator; @@ -44,24 +43,23 @@ import lombok.extern.slf4j.Slf4j; public class ScanBatchResponse implements StreamObserver { static ByteBufferAllocator bfAllocator = new ByteBufferAllocator(ParallelScanIterator.maxBodySize * 3 / 2, 1000); + static ByteBufferAllocator alloc = + new ByteBufferAllocator(ParallelScanIterator.maxBodySize * 3 / 2, 1000); private final int maxInFlightCount = PropertyUtil.getInt("app.scan.stream.inflight", 16); private final int activeTimeout = PropertyUtil.getInt("app.scan.stream.timeout", 60); //单位秒 - private final Object stateLock = new Object(); - private final ReentrantLock iteratorLock = new ReentrantLock(); private final StreamObserver sender; private final HgStoreWrapperEx wrapper; private final ThreadPoolExecutor executor; - private final long logId; // 当前正在遍历的迭代器 private ScanIterator iterator; // 下一次发送的序号 - private volatile int nextSeqNo; + private volatile int seqNo; // Client已消费的序号 private volatile int clientSeqNo; // 已经发送的条目数 - private volatile long entriesCounter; + private volatile long count; // 客户端要求返回的最大条目数 - private volatile long clientLimit; // 客户端要求的最大条目数 + private volatile long limit; private ScanQueryRequest query; // 上次读取数据时间 private long activeTime; @@ -73,10 +71,9 @@ public class ScanBatchResponse implements StreamObserver { this.wrapper = wrapper; this.executor = executor; this.iterator = null; - this.nextSeqNo = 1; + this.seqNo = 1; this.state = State.IDLE; this.activeTime = System.currentTimeMillis(); - this.logId = response.hashCode(); } /** @@ -95,7 +92,7 @@ public class ScanBatchResponse implements StreamObserver { break; case RECEIPT_REQUEST: // 消息异步应答 this.clientSeqNo = request.getReceiptRequest().getTimes(); - if (nextSeqNo - clientSeqNo < maxInFlightCount) { + if (seqNo - clientSeqNo < maxInFlightCount) { synchronized (stateLock) { if (state == State.IDLE) { state = State.DOING; @@ -136,16 +133,9 @@ public class ScanBatchResponse implements StreamObserver { */ private void startQuery(String graphName, ScanQueryRequest request) { this.query = request; - - // log.info("Stream {} startQuery graphName is {}, query degree/keylimit/limit is " + - // "{}/{}/{}, scanType is {}, orderType is {}", - // this.logId, graphName, query.getSkipDegree(), query.getPerKeyLimit(), query - // .getLimit(), - // query.getScanType(), query.getOrderType()); - - this.clientLimit = request.getLimit(); - this.entriesCounter = 0; - this.iterator = ScanUtil.getParallelIterator(graphName, request, this.wrapper, executor); + this.limit = request.getLimit(); + this.count = 0; + this.iterator = getParallelIterator(graphName, request, this.wrapper, executor); synchronized (stateLock) { if (state == State.IDLE) { state = State.DOING; @@ -161,63 +151,75 @@ public class ScanBatchResponse implements StreamObserver { */ private void closeQuery() { setStateDone(); - iteratorLock.lock(); + try { + closeIter(); + this.sender.onCompleted(); + } catch (Exception e) { + log.error("exception ", e); + } + int active = ScanBatchResponseFactory.getInstance().removeStreamObserver(this); + log.info("ScanBatchResponse closeQuery, active count is {}", active); + } + + private void closeIter() { try { if (this.iterator != null) { this.iterator.close(); this.iterator = null; } - this.sender.onCompleted(); } catch (Exception e) { - log.error("exception ", e); - } finally { - iteratorLock.unlock(); + } - int active = ScanBatchResponseFactory.getInstance().removeStreamObserver(this); - log.info("ScanBatchResponse closeQuery, active count is {}", active); } /** * 发送数据 */ private void sendEntries() { + if (state == State.DONE || iterator == null) { + setStateIdle(); + return; + } iteratorLock.lock(); try { if (state == State.DONE || iterator == null) { setStateIdle(); return; } - KvStream.Builder dataBuilder = KvStream.newBuilder() - .setVersion(1); - while (iterator.hasNext() - && (nextSeqNo - clientSeqNo < maxInFlightCount) - && this.entriesCounter < clientLimit - && state != State.DONE) { - KVByteBuffer buffer = new KVByteBuffer(bfAllocator.get()); - List dataList = iterator.next(); + KvStream.Builder dataBuilder = KvStream.newBuilder().setVersion(1); + while (state != State.DONE && iterator.hasNext() + && (seqNo - clientSeqNo < maxInFlightCount) + && this.count < limit) { + KVByteBuffer buffer = new KVByteBuffer(alloc.get()); + List dataList = iterator.next(); dataList.forEach(kv -> { kv.write(buffer); - this.entriesCounter++; + this.count++; }); dataBuilder.setStream(buffer.flip().getBuffer()); - dataBuilder.setSeqNo(nextSeqNo++); - dataBuilder.complete(e -> { - bfAllocator.release(buffer.getBuffer()); - }); + dataBuilder.setSeqNo(seqNo++); + dataBuilder.complete(e -> alloc.release(buffer.getBuffer())); this.sender.onNext(dataBuilder.build()); this.activeTime = System.currentTimeMillis(); } - if (!iterator.hasNext() || this.entriesCounter >= clientLimit || state == State.DONE) { + if (!iterator.hasNext() || this.count >= limit || state == State.DONE) { + closeIter(); this.sender.onNext(KvStream.newBuilder().setOver(true).build()); setStateDone(); } else { setStateIdle(); } } catch (Throwable e) { - log.error("exception ", e); - setStateIdle(); - if (this.sender != null) { - this.sender.onError(e); + if (this.state != State.DONE) { + log.error(" send data exception: ", e); + setStateIdle(); + if (this.sender != null) { + try { + this.sender.onError(e); + } catch (Exception ex) { + + } + } } } finally { iteratorLock.unlock(); diff --git a/hg-store-node/src/main/resources/version.txt b/hg-store-node/src/main/resources/version.txt index 1ac53bb4b..b55f10804 100644 --- a/hg-store-node/src/main/resources/version.txt +++ b/hg-store-node/src/main/resources/version.txt @@ -1 +1 @@ -3.6.3 \ No newline at end of file +3.6.5 \ No newline at end of file diff --git a/hg-store-rocksdb/src/main/java/org/apache/hugegraph/rocksdb/access/RocksDBSession.java b/hg-store-rocksdb/src/main/java/org/apache/hugegraph/rocksdb/access/RocksDBSession.java index 40ad796c8..0edadffb6 100644 --- a/hg-store-rocksdb/src/main/java/org/apache/hugegraph/rocksdb/access/RocksDBSession.java +++ b/hg-store-rocksdb/src/main/java/org/apache/hugegraph/rocksdb/access/RocksDBSession.java @@ -63,9 +63,6 @@ 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; diff --git a/hg-store-rocksdb/src/main/java/org/apache/hugegraph/rocksdb/access/SessionOperatorImpl.java b/hg-store-rocksdb/src/main/java/org/apache/hugegraph/rocksdb/access/SessionOperatorImpl.java index c9056260a..cca5c6b09 100644 --- a/hg-store-rocksdb/src/main/java/org/apache/hugegraph/rocksdb/access/SessionOperatorImpl.java +++ b/hg-store-rocksdb/src/main/java/org/apache/hugegraph/rocksdb/access/SessionOperatorImpl.java @@ -32,7 +32,6 @@ 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; diff --git a/hg-store-rocksdb/src/test/java/org/apache/hugegraph/rocksdb/access/SnapshotManagerTest.java b/hg-store-rocksdb/src/test/java/org/apache/hugegraph/rocksdb/access/SnapshotManagerTest.java index 789020dfe..77eb2d554 100644 --- a/hg-store-rocksdb/src/test/java/org/apache/hugegraph/rocksdb/access/SnapshotManagerTest.java +++ b/hg-store-rocksdb/src/test/java/org/apache/hugegraph/rocksdb/access/SnapshotManagerTest.java @@ -23,13 +23,6 @@ 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; diff --git a/hg-store-test/src/main/resources/version.txt b/hg-store-test/src/main/resources/version.txt index 084e244ce..b55f10804 100644 --- a/hg-store-test/src/main/resources/version.txt +++ b/hg-store-test/src/main/resources/version.txt @@ -1 +1 @@ -3.6.0 \ No newline at end of file +3.6.5 \ No newline at end of file