This is an automated email from the ASF dual-hosted git repository.
albumenj pushed a commit to branch 3.0
in repository https://gitbox.apache.org/repos/asf/dubbo.git
The following commit(s) were added to refs/heads/3.0 by this push:
new 3531a7d [3.0] Feature/state router (#6844)
3531a7d is described below
commit 3531a7d681da5574f31609beeb2ed288af7aaebc
Author: panxiaojun233 <[email protected]>
AuthorDate: Tue May 18 17:13:50 2021 +0800
[3.0] Feature/state router (#6844)
* state router
* fix router chain
* fix router chain
* fix comment
* fix comment
* format files
* deal dependecy
* fix comment
* fix comment
* fix
* fix ut
* fix ut
* fix
* fix comment
* add code comments
* add code comments
* add code comments
* add code comments
---
dubbo-cluster/pom.xml | 4 +
.../org/apache/dubbo/rpc/cluster/RouterChain.java | 202 +++++++++++++++-
.../rpc/cluster/directory/StaticDirectory.java | 1 +
.../cluster/router/state/AbstractStateRouter.java | 119 +++++++++
.../dubbo/rpc/cluster/router/state/AddrCache.java | 52 ++++
.../dubbo/rpc/cluster/router/state/BitList.java | 232 ++++++++++++++++++
.../router/state/CacheableStateRouterFactory.java | 38 +++
.../rpc/cluster/router/state/RouterCache.java | 49 ++++
.../rpc/cluster/router/state/StateRouter.java | 100 ++++++++
.../cluster/router/state/StateRouterFactory.java | 35 +++
.../cluster/router/tag/TagDynamicStateRouter.java | 267 +++++++++++++++++++++
.../router/tag/TagDynamicStateRouterFactory.java | 37 +++
.../cluster/router/tag/TagStaticStateRouter.java | 155 ++++++++++++
.../router/tag/TagStaticStateRouterFactory.java | 37 +++
...bbo.rpc.cluster.router.state.StateRouterFactory | 2 +
.../manager/DefaultExecutorRepository.java | 40 ++-
.../threadpool/manager/ExecutorRepository.java | 5 +-
dubbo-dependencies-bom/pom.xml | 6 +
dubbo-distribution/dubbo-all/pom.xml | 4 +
.../java/org/apache/dubbo/rpc/RpcException.java | 1 +
20 files changed, 1367 insertions(+), 19 deletions(-)
diff --git a/dubbo-cluster/pom.xml b/dubbo-cluster/pom.xml
index bb58849..3055454 100644
--- a/dubbo-cluster/pom.xml
+++ b/dubbo-cluster/pom.xml
@@ -32,6 +32,10 @@
</properties>
<dependencies>
<dependency>
+ <groupId>org.roaringbitmap</groupId>
+ <artifactId>RoaringBitmap</artifactId>
+ </dependency>
+ <dependency>
<groupId>org.apache.dubbo</groupId>
<artifactId>dubbo-rpc-api</artifactId>
<version>${project.parent.version}</version>
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/RouterChain.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/RouterChain.java
index a16b2aa..eb92945 100644
--- a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/RouterChain.java
+++ b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/RouterChain.java
@@ -16,18 +16,31 @@
*/
package org.apache.dubbo.rpc.cluster;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Semaphore;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.stream.Collectors;
+
import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.Version;
import org.apache.dubbo.common.extension.ExtensionLoader;
import org.apache.dubbo.common.logger.Logger;
import org.apache.dubbo.common.logger.LoggerFactory;
+import org.apache.dubbo.common.threadpool.manager.ExecutorRepository;
import org.apache.dubbo.common.utils.CollectionUtils;
+import org.apache.dubbo.common.utils.NetUtils;
import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
-
-import java.util.ArrayList;
-import java.util.Collections;
-import java.util.List;
-import java.util.stream.Collectors;
+import org.apache.dubbo.rpc.RpcException;
+import org.apache.dubbo.rpc.cluster.directory.StaticDirectory;
+import org.apache.dubbo.rpc.cluster.router.state.AddrCache;
+import org.apache.dubbo.rpc.cluster.router.state.BitList;
+import org.apache.dubbo.rpc.cluster.router.state.RouterCache;
+import org.apache.dubbo.rpc.cluster.router.state.StateRouter;
+import org.apache.dubbo.rpc.cluster.router.state.StateRouterFactory;
/**
* Router chain
@@ -45,19 +58,50 @@ public class RouterChain<T> {
// instance will never delete or recreate.
private List<Router> builtinRouters = Collections.emptyList();
+ private List<StateRouter> builtinStateRouters = Collections.emptyList();
+ private List<StateRouter> stateRouters = Collections.emptyList();
+ private final ExecutorRepository executorRepository =
ExtensionLoader.getExtensionLoader(ExecutorRepository.class)
+ .getDefaultExtension();
+
+ protected URL url;
+
+ protected AtomicReference<AddrCache<T>> cache = new AtomicReference<>();
+
+ private final Semaphore loopPermit = new Semaphore(1);
+ private final Semaphore loopPermitNotify = new Semaphore(1);
+
+ private final ExecutorService loopPool;
+
+ private boolean firstBuildCache = true;
+
+ private static final Logger logger =
LoggerFactory.getLogger(StaticDirectory.class);
+
public static <T> RouterChain<T> buildChain(URL url) {
return new RouterChain<>(url);
}
private RouterChain(URL url) {
+ loopPool = executorRepository.nextExecutorExecutor();
List<RouterFactory> extensionFactories =
ExtensionLoader.getExtensionLoader(RouterFactory.class)
- .getActivateExtension(url, "router");
+ .getActivateExtension(url, "router");
List<Router> routers = extensionFactories.stream()
- .map(factory -> factory.getRouter(url))
- .collect(Collectors.toList());
+ .map(factory -> factory.getRouter(url))
+ .collect(Collectors.toList());
initWithRouters(routers);
+
+ List<StateRouterFactory> extensionStateRouterFactories =
ExtensionLoader.getExtensionLoader(
+ StateRouterFactory.class)
+ .getActivateExtension(url, "stateRouter");
+
+ List<StateRouter> stateRouters = extensionStateRouterFactories.stream()
+ .map(factory -> factory.getRouter(url, this))
+ .sorted(StateRouter::compareTo)
+ .collect(Collectors.toList());
+
+ // init state routers
+ initWithStateRouters(stateRouters);
}
/**
@@ -70,6 +114,11 @@ public class RouterChain<T> {
this.sort();
}
+ public void initWithStateRouters(List<StateRouter> builtinRouters) {
+ this.builtinStateRouters = builtinRouters;
+ this.stateRouters = new ArrayList<>(builtinRouters);
+ }
+
/**
* If we use route:// protocol in version before 2.7.0, each URL will
generate a Router instance, so we should
* keep the routers up to date, that is, each time router URLs changes, we
should update the routers list, only
@@ -86,6 +135,14 @@ public class RouterChain<T> {
this.routers = newRouters;
}
+ public void addStateRouters(List<StateRouter> stateRouters) {
+ List<StateRouter> newStateRouters = new ArrayList<>();
+ newStateRouters.addAll(builtinStateRouters);
+ newStateRouters.addAll(stateRouters);
+ CollectionUtils.sort(newStateRouters);
+ this.stateRouters = newStateRouters;
+ }
+
public List<Router> getRouters() {
return routers;
}
@@ -95,13 +152,35 @@ public class RouterChain<T> {
}
/**
- *
* @param url
* @param invocation
* @return
*/
public List<Invoker<T>> route(URL url, Invocation invocation) {
- List<Invoker<T>> finalInvokers = invokers;
+
+ AddrCache<T> cache = this.cache.get();
+ if (cache == null) {
+ throw new RpcException(RpcException.ROUTER_CACHE_NOT_BUILD,
"Failed to invoke the method "
+ + invocation.getMethodName() + " in the service " +
url.getServiceInterface()
+ + ". address cache not build "
+ + " on the consumer " + NetUtils.getLocalHost()
+ + " using the dubbo version " + Version.getVersion()
+ + ".");
+ }
+ BitList<Invoker<T>> finalBitListInvokers = new BitList<>(invokers,
false);
+ for (StateRouter stateRouter : stateRouters) {
+ if (stateRouter.isEnable()) {
+ RouterCache<T> routerCache =
cache.getCache().get(stateRouter.getName());
+ finalBitListInvokers = stateRouter.route(finalBitListInvokers,
routerCache, url, invocation);
+ }
+ }
+
+ List<Invoker<T>> finalInvokers = new
ArrayList<>(finalBitListInvokers.size());
+
+ for(Invoker<T> invoker: finalBitListInvokers) {
+ finalInvokers.add(invoker);
+ }
+
for (Router router : routers) {
finalInvokers = router.route(finalInvokers, url, invocation);
}
@@ -114,7 +193,99 @@ public class RouterChain<T> {
*/
public void setInvokers(List<Invoker<T>> invokers) {
this.invokers = (invokers == null ? Collections.emptyList() :
invokers);
+ stateRouters.forEach(router -> router.notify(this.invokers));
routers.forEach(router -> router.notify(this.invokers));
+ loop(true);
+ }
+
+ /**
+ * Build the asynchronous address cache for stateRouter.
+ * @param notify Whether the addresses in registry has changed.
+ */
+ private void buildCache(boolean notify) {
+ if (invokers == null || invokers.size() <= 0) {
+ return;
+ }
+ AddrCache<T> origin = cache.get();
+ List<Invoker<T>> copyInvokers = new ArrayList<>(this.invokers);
+ AddrCache<T> newCache = new AddrCache<T>();
+ newCache.setInvokers(invokers);
+ for (StateRouter stateRouter : stateRouters) {
+ RouterCache routerCache;
+ try {
+ routerCache = poolRouter(stateRouter, origin, copyInvokers,
notify);
+ //file cache
+ newCache.getCache().put(stateRouter.getName(), routerCache);
+ } catch (Throwable t) {
+ logger.error("Failed to pool router: " + stateRouter.getUrl()
+ ", cause: " + t.getMessage(), t);
+ return;
+ }
+ }
+
+ this.cache.set(newCache);
+ }
+
+ /**
+ * Cache the address list for each StateRouter.
+ * @param router router
+ * @param orign The original address cache
+ * @param invokers The full address list
+ * @param notify Whether the addresses in registry has changed.
+ * @return
+ */
+ private RouterCache poolRouter(StateRouter router, AddrCache<T> orign,
List<Invoker<T>> invokers, boolean notify) {
+ String routerName = router.getName();
+ RouterCache routerCache;
+ if (isCacheMiss(orign, routerName) || router.shouldRePool() || notify)
{
+ return router.pool(invokers);
+ } else {
+ routerCache = orign.getCache().get(routerName);
+ }
+ if (routerCache == null) {
+ return new RouterCache();
+ }
+ return routerCache;
+ }
+
+ private boolean isCacheMiss(AddrCache<T> cache, String routerName) {
+ return cache == null || cache.getCache() == null ||
cache.getInvokers() == null || cache.getCache().get(
+ routerName)
+ == null;
+ }
+
+ /***
+ * Build the asynchronous address cache for stateRouter.
+ * @param notify Whether the addresses in registry has changed.
+ */
+ public void loop(boolean notify) {
+ if (firstBuildCache) {
+ firstBuildCache = false;
+ buildCache(notify);
+ }
+ if (notify) {
+ if (loopPermitNotify.tryAcquire()) {
+ loopPool.submit(new NotifyLoopRunnable(true));
+ }
+ } else {
+ if (loopPermit.tryAcquire()) {
+ loopPool.submit(new NotifyLoopRunnable(false));
+ }
+ }
+ }
+
+ class NotifyLoopRunnable implements Runnable {
+
+ private final boolean notify;
+
+ public NotifyLoopRunnable(boolean notify) {
+ this.notify = notify;
+ }
+
+ @Override
+ public void run() {
+ loopPermitNotify.release();
+ buildCache(notify);
+ }
}
public void destroy() {
@@ -128,5 +299,16 @@ public class RouterChain<T> {
}
routers = Collections.emptyList();
builtinRouters = Collections.emptyList();
+
+ for (StateRouter router : stateRouters) {
+ try {
+ router.stop();
+ } catch (Exception e) {
+ LOGGER.error("Error trying to stop stateRouter " +
router.getClass(), e);
+ }
+ }
+ stateRouters = Collections.emptyList();
+ builtinStateRouters = Collections.emptyList();
}
+
}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/directory/StaticDirectory.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/directory/StaticDirectory.java
index db1b202..50110aa 100644
---
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/directory/StaticDirectory.java
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/directory/StaticDirectory.java
@@ -94,6 +94,7 @@ public class StaticDirectory<T> extends AbstractDirectory<T> {
public void buildRouterChain() {
RouterChain<T> routerChain = RouterChain.buildChain(getUrl());
routerChain.setInvokers(invokers);
+ routerChain.loop(true);
this.setRouterChain(routerChain);
}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/AbstractStateRouter.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/AbstractStateRouter.java
new file mode 100644
index 0000000..fc5e362
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/AbstractStateRouter.java
@@ -0,0 +1,119 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.state;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.extension.ExtensionLoader;
+import org.apache.dubbo.rpc.Invocation;
+import org.apache.dubbo.rpc.Invoker;
+import org.apache.dubbo.rpc.RpcException;
+import org.apache.dubbo.rpc.cluster.RouterChain;
+import org.apache.dubbo.rpc.cluster.governance.GovernanceRuleRepository;
+
+/***
+ * The abstract class of StateRoute.
+ * @since 3.0
+ */
+public abstract class AbstractStateRouter implements StateRouter {
+ final protected RouterChain chain;
+ protected int priority = DEFAULT_PRIORITY;
+ protected boolean force = false;
+ protected URL url;
+ protected List<Invoker> invokers;
+ protected AtomicReference<AddrCache> cache;
+ protected GovernanceRuleRepository ruleRepository;
+
+ public AbstractStateRouter(URL url, RouterChain chain) {
+ this.ruleRepository =
ExtensionLoader.getExtensionLoader(GovernanceRuleRepository.class).getDefaultExtension();
+ this.chain = chain;
+ this.url = url;
+ }
+
+ @Override
+ public <T> void notify(List<Invoker<T>> invokers) {
+ this.invokers = (List)invokers;
+ }
+
+ @Override
+ public URL getUrl() {
+ return url;
+ }
+
+ public void setUrl(URL url) {
+ this.url = url;
+ }
+
+ @Override
+ public boolean isRuntime() {
+ return true;
+ }
+
+ @Override
+ public boolean isForce() {
+ return force;
+ }
+
+ public void setForce(boolean force) {
+ this.force = force;
+ }
+
+ @Override
+ public int getPriority() {
+ return priority;
+ }
+
+ public void setPriority(int priority) {
+ this.priority = priority;
+ }
+
+ @Override
+ public <T> BitList<Invoker<T>> route(BitList<Invoker<T>> invokers,
RouterCache<T> cache, URL url,
+ Invocation invocation) throws RpcException {
+
+ List<String> tags = getTags(url, invocation);
+
+ if (tags == null) {
+ return invokers;
+ }
+ for (String tag : tags) {
+ BitList<Invoker<T>> tagInvokers = cache.getAddrPool().get(tag);
+ if (tagMatchFail(tagInvokers)) {
+ continue;
+ }
+ return tagInvokers.intersect(invokers,
invokers.getUnmodifiableList());
+ }
+
+ return invokers;
+ }
+
+ protected List<String> getTags(URL url, Invocation invocation) {
+ return new ArrayList<String>();
+ }
+
+ public <T> Boolean tagMatchFail(BitList<Invoker<T>> invokers) {
+ return invokers == null || invokers.isEmpty();
+ }
+
+ @Override
+ public void pool() {
+ chain.loop(false);
+ }
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/AddrCache.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/AddrCache.java
new file mode 100644
index 0000000..9d1e9b3
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/AddrCache.java
@@ -0,0 +1,52 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.state;
+
+import java.util.List;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+
+import org.apache.dubbo.rpc.Invoker;
+
+/***
+ * Address cache,
+ * used to cache the results of the StaterRouter's asynchronous address list
calculations.
+ * @param <T>
+ * @since 3.0
+ */
+public class AddrCache<T> {
+ private final static ConcurrentHashMap EMPTY_MAP = new
ConcurrentHashMap<>();
+
+ protected List<Invoker<T>> invokers;
+ protected ConcurrentMap<String, RouterCache<T>> cache = EMPTY_MAP;
+
+ public List<Invoker<T>> getInvokers() {
+ return invokers;
+ }
+
+ public void setInvokers(List<Invoker<T>> invokers) {
+ this.invokers = invokers;
+ }
+
+ public ConcurrentMap<String, RouterCache<T>> getCache() {
+ return cache;
+ }
+
+ public void setCache(ConcurrentHashMap<String, RouterCache<T>> cache) {
+ this.cache = cache;
+ }
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/BitList.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/BitList.java
new file mode 100644
index 0000000..71456cb
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/BitList.java
@@ -0,0 +1,232 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.state;
+
+import java.util.Arrays;
+import java.util.Collection;
+import java.util.Iterator;
+import java.util.List;
+import java.util.ListIterator;
+
+import org.roaringbitmap.IntIterator;
+import org.roaringbitmap.RoaringBitmap;
+
+/**
+ * BitList based on BitMap implementation.
+ * @param <E>
+ * @since 3.0
+ */
+public class BitList<E> implements List<E> {
+ private final RoaringBitmap rootMap;
+ private final List<E> unmodifiableList;
+
+ public BitList(List<E> unmodifiableList, boolean empty) {
+ this.unmodifiableList = unmodifiableList;
+ this.rootMap = new RoaringBitmap();
+ if (!empty) {
+ this.rootMap.add(0L, unmodifiableList.size());
+ }
+ }
+
+ public BitList(List<E> unmodifiableList, RoaringBitmap rootMap) {
+ this.unmodifiableList = unmodifiableList;
+ this.rootMap = rootMap;
+ }
+
+ public BitList(List<E> unmodifiableList) {
+ this(unmodifiableList, false);
+ }
+
+ public List<E> getUnmodifiableList() {
+ return unmodifiableList;
+ }
+
+ public void addIndex(int index) {
+ this.rootMap.add(index);
+ }
+
+ public BitList<E> intersect(List<E> b, List<E> totalList) {
+ RoaringBitmap resultMap = rootMap.clone();
+ resultMap.and(((BitList)b).rootMap);
+ return new BitList<>(totalList, resultMap);
+ }
+
+ @Override
+ public int size() {
+ return rootMap.getCardinality();
+ }
+
+ @Override
+ public boolean isEmpty() {
+ return rootMap.isEmpty();
+ }
+
+ @Override
+ public boolean contains(Object o) {
+ int idx = unmodifiableList.indexOf(o);
+ return idx >= 0 && rootMap.contains(idx);
+ }
+
+ @Override
+ public Iterator<E> iterator() {
+ return new Iterator<E>() {
+ private int prev = -1;
+
+ @Override
+ public boolean hasNext() {
+ return -1 != rootMap.nextValue(prev + 1);
+ }
+
+ @Override
+ public E next() {
+ prev = (int)rootMap.nextValue(prev + 1);
+ return unmodifiableList.get(prev);
+ }
+
+ @Override
+ public void remove() {
+ rootMap.remove(prev);
+ }
+ };
+ }
+
+ @Override
+ public Object[] toArray() {
+ int size = size();
+ Object[] obj = new Object[size];
+ for (int i = 0; i < size; i++) {
+ obj[i] = unmodifiableList.get(rootMap.select(i));
+ }
+ return obj;
+ }
+
+ @Override
+ public <T> T[] toArray(T[] a) {
+ int size = size();
+ Object[] arr = toArray();
+ if (a.length < size)
+ // Make a new array of a's runtime type, but my contents:
+ { return (T[])Arrays.copyOf(arr, size, a.getClass()); }
+ System.arraycopy(arr, 0, a, 0, size);
+ if (a.length > size) { a[size] = null; }
+ return null;
+ }
+
+ @Override
+ public boolean add(E e) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public boolean remove(Object o) {
+ int idx = unmodifiableList.indexOf(o);
+ if (idx > -1) {
+ rootMap.remove(idx);
+ }
+ return true;
+ }
+
+ @Override
+ public boolean containsAll(Collection<?> c) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public boolean addAll(Collection<? extends E> c) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public boolean addAll(int index, Collection<? extends E> c) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public boolean removeAll(Collection<?> c) {
+ for (Object o : c) {
+ remove(o);
+ }
+ return true;
+ }
+
+ @Override
+ public boolean retainAll(Collection<?> c) {
+ return false;
+ }
+
+ @Override
+ public void clear() {
+ rootMap.clear();
+ }
+
+ @Override
+ public E get(int index) {
+ int real = rootMap.select(index);
+ return unmodifiableList.get(real);
+ }
+
+ @Override
+ public E set(int index, E element) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public void add(int index, E element) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public E remove(int index) {
+ rootMap.remove(index);
+ return null;
+ }
+
+ @Override
+ public int indexOf(Object o) {
+ IntIterator intIterator = rootMap.getIntIterator();
+ int st = 0;
+ while (intIterator.hasNext()) {
+ int idxInMap = intIterator.next();
+ if (unmodifiableList.get(idxInMap).equals(o)) {
+ return st;
+ }
+ st++;
+ }
+ return -1;
+ }
+
+ @Override
+ public int lastIndexOf(Object o) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public ListIterator<E> listIterator() {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public ListIterator<E> listIterator(int index) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public List<E> subList(int fromIndex, int toIndex) {
+ throw new UnsupportedOperationException();
+ }
+
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/CacheableStateRouterFactory.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/CacheableStateRouterFactory.java
new file mode 100644
index 0000000..5f38ccb
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/CacheableStateRouterFactory.java
@@ -0,0 +1,38 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.state;
+
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.rpc.cluster.RouterChain;
+
+/**
+ * If you want to provide a router implementation based on design of v2.7.0,
please extend from this abstract class.
+ * For 2.6.x style router, please implement and use RouterFactory directly.
+ */
+public abstract class CacheableStateRouterFactory implements
StateRouterFactory {
+ private final ConcurrentMap<String, StateRouter> routerMap = new
ConcurrentHashMap<>();
+
+ @Override
+ public StateRouter getRouter(URL url, RouterChain chain) {
+ return routerMap.computeIfAbsent(url.getServiceKey(), k ->
createRouter(url, chain));
+ }
+
+ protected abstract StateRouter createRouter(URL url, RouterChain chain);
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/RouterCache.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/RouterCache.java
new file mode 100644
index 0000000..0b7b499
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/RouterCache.java
@@ -0,0 +1,49 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.state;
+
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+
+import org.apache.dubbo.rpc.Invoker;
+
+/***
+ * Cache the address list for each Router.
+ * @param <T>
+ * @since 3.0
+ */
+public class RouterCache<T> {
+ private final static ConcurrentHashMap EMPTY_MAP = new
ConcurrentHashMap<>();
+ protected ConcurrentMap<String, BitList<Invoker<T>>> addrPool = EMPTY_MAP;
+ protected Object addrMetadata;
+
+ public ConcurrentMap<String, BitList<Invoker<T>>> getAddrPool() {
+ return addrPool;
+ }
+
+ public void setAddrPool(ConcurrentHashMap<String, BitList<Invoker<T>>>
addrPool) {
+ this.addrPool = addrPool;
+ }
+
+ public Object getAddrMetadata() {
+ return addrMetadata;
+ }
+
+ public void setAddrMetadata(Object addrMetadata) {
+ this.addrMetadata = addrMetadata;
+ }
+}
\ No newline at end of file
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/StateRouter.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/StateRouter.java
new file mode 100644
index 0000000..ee28459
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/StateRouter.java
@@ -0,0 +1,100 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.state;
+
+import java.util.List;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.rpc.Invocation;
+import org.apache.dubbo.rpc.Invoker;
+import org.apache.dubbo.rpc.RpcException;
+import org.apache.dubbo.rpc.cluster.Directory;
+
+/**
+ * State Router. (SPI, Prototype, ThreadSafe)
+ * <p>
+ * <a href="http://en.wikipedia.org/wiki/Routing">Routing</a>
+ *
+ * @see org.apache.dubbo.rpc.cluster.Cluster#join(Directory)
+ * @see Directory#list(Invocation)
+ * @since 3.0
+ */
+public interface StateRouter extends Comparable<StateRouter> {
+
+ int DEFAULT_PRIORITY = Integer.MAX_VALUE;
+
+ /**
+ * Get the router url.
+ *
+ * @return url
+ */
+ URL getUrl();
+
+ /***
+ * Filter invokers with current routing rule and only return the invokers
that comply with the rule.
+ * Caching address lists in BitMap mode improves routing performance.
+ * @param invokers invoker bit list
+ * @param cache router address cache
+ * @param url refer url
+ * @param invocation invocation
+ * @param <T>
+ * @return routed invokers
+ * @throws RpcException
+ * @Since 3.0
+ */
+ <T> BitList<Invoker<T>> route(BitList<Invoker<T>> invokers, RouterCache<T>
cache, URL url, Invocation invocation)
+ throws
+ RpcException;
+
+ default <T> void notify(List<Invoker<T>> invokers) {
+
+ }
+
+ /**
+ * To decide whether this router need to execute every time an RPC comes
or should only execute when addresses or
+ * rule change.
+ *
+ * @return true if the router need to execute every time.
+ */
+ boolean isRuntime();
+
+ boolean isEnable();
+
+ boolean isForce();
+
+ int getPriority();
+
+ @Override
+ default int compareTo(StateRouter o) {
+ if (o == null) {
+ throw new IllegalArgumentException();
+ }
+ return Integer.compare(this.getPriority(), o.getPriority());
+ }
+
+ String getName();
+
+ boolean shouldRePool();
+
+ <T> RouterCache<T> pool(List<Invoker<T>> invokers);
+
+ void pool();
+
+ default void stop() {
+ //do nothing by default
+ }
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/StateRouterFactory.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/StateRouterFactory.java
new file mode 100644
index 0000000..9780c36
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/state/StateRouterFactory.java
@@ -0,0 +1,35 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.state;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.extension.Adaptive;
+import org.apache.dubbo.common.extension.SPI;
+import org.apache.dubbo.rpc.cluster.RouterChain;
+
+@SPI
+public interface StateRouterFactory {
+ /**
+ * Create state router.
+ *
+ * @param url url
+ * @return router instance
+ * @since 3.0
+ */
+ @Adaptive("protocol")
+ <T> StateRouter getRouter(URL url, RouterChain<T> chain);
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagDynamicStateRouter.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagDynamicStateRouter.java
new file mode 100644
index 0000000..231aa2c
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagDynamicStateRouter.java
@@ -0,0 +1,267 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.tag;
+
+import java.net.UnknownHostException;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.config.configcenter.ConfigChangeType;
+import org.apache.dubbo.common.config.configcenter.ConfigChangedEvent;
+import org.apache.dubbo.common.config.configcenter.ConfigurationListener;
+import org.apache.dubbo.common.config.configcenter.DynamicConfiguration;
+import org.apache.dubbo.common.logger.Logger;
+import org.apache.dubbo.common.logger.LoggerFactory;
+import org.apache.dubbo.common.utils.CollectionUtils;
+import org.apache.dubbo.common.utils.NetUtils;
+import org.apache.dubbo.common.utils.StringUtils;
+import org.apache.dubbo.rpc.Invocation;
+import org.apache.dubbo.rpc.Invoker;
+import org.apache.dubbo.rpc.RpcException;
+import org.apache.dubbo.rpc.cluster.RouterChain;
+import org.apache.dubbo.rpc.cluster.router.state.AbstractStateRouter;
+import org.apache.dubbo.rpc.cluster.router.state.BitList;
+import org.apache.dubbo.rpc.cluster.router.state.RouterCache;
+import org.apache.dubbo.rpc.cluster.router.tag.model.TagRouterRule;
+import org.apache.dubbo.rpc.cluster.router.tag.model.TagRuleParser;
+
+import static org.apache.dubbo.common.constants.CommonConstants.ANYHOST_VALUE;
+import static org.apache.dubbo.common.constants.CommonConstants.TAG_KEY;
+import static org.apache.dubbo.rpc.Constants.FORCE_USE_TAG;
+
+/**
+ * TagDynamicStateRouter, "application.tag-router"
+ */
+public class TagDynamicStateRouter extends AbstractStateRouter implements
ConfigurationListener {
+ public static final String NAME = "TAG_ROUTER";
+ private static final int TAG_ROUTER_DEFAULT_PRIORITY = 100;
+ private static final Logger logger =
LoggerFactory.getLogger(TagDynamicStateRouter.class);
+ private static final String RULE_SUFFIX = ".tag-router";
+ private static final String NO_TAG = "noTag";
+
+ private TagRouterRule tagRouterRule;
+ private String application;
+
+ public TagDynamicStateRouter(URL url, RouterChain chain) {
+ super(url, chain);
+ this.priority = TAG_ROUTER_DEFAULT_PRIORITY;
+ }
+
+ @Override
+ public synchronized void process(ConfigChangedEvent event) {
+ if (logger.isDebugEnabled()) {
+ logger.debug("Notification of tag rule, change type is: " +
event.getChangeType() + ", raw rule is:\n " +
+ event.getContent());
+ }
+
+ try {
+ if (event.getChangeType().equals(ConfigChangeType.DELETED)) {
+ this.tagRouterRule = null;
+ } else {
+ this.tagRouterRule = TagRuleParser.parse(event.getContent());
+ }
+ } catch (Exception e) {
+ logger.error("Failed to parse the raw tag router rule and it will
not take effect, please check if the " +
+ "rule matches with the template, the raw rule is:\n ", e);
+ }
+ }
+
+ @Override
+ public URL getUrl() {
+ return url;
+ }
+
+ @Override
+ public <T> BitList<Invoker<T>> route(BitList<Invoker<T>> invokers,
RouterCache<T> cache, URL url,
+ Invocation invocation) throws RpcException {
+
+
+ final TagRouterRule tagRouterRuleCopy =
(TagRouterRule)cache.getAddrMetadata();
+
+ String tag = StringUtils.isEmpty(invocation.getAttachment(TAG_KEY)) ?
url.getParameter(TAG_KEY) :
+ invocation.getAttachment(TAG_KEY);
+
+ ConcurrentMap<String, BitList<Invoker<T>>> addrPool =
cache.getAddrPool();
+
+ if (StringUtils.isEmpty(tag)) {
+ return invokers.intersect(addrPool.get(NO_TAG),
invokers.getUnmodifiableList());
+ } else {
+ BitList<Invoker<T>> result = addrPool.get(tag);
+
+ if (CollectionUtils.isNotEmpty(result) || (tagRouterRuleCopy !=
null && tagRouterRuleCopy.isForce())
+ || isForceUseTag(invocation)) {
+ return invokers.intersect(result,
invokers.getUnmodifiableList());
+ } else {
+ invocation.setAttachment(TAG_KEY, NO_TAG);
+ return invokers;
+ }
+ }
+ }
+
+ private boolean isForceUseTag(Invocation invocation) {
+ return Boolean.valueOf(invocation.getAttachment(FORCE_USE_TAG,
url.getParameter(FORCE_USE_TAG, "false")));
+ }
+
+ @Override
+ public boolean isRuntime() {
+ return tagRouterRule != null && tagRouterRule.isRuntime();
+ }
+
+ @Override
+ public boolean isEnable() {
+ return true;
+ }
+
+ @Override
+ public boolean isForce() {
+ return tagRouterRule != null && tagRouterRule.isForce();
+ }
+
+ @Override
+ public String getName() {
+ return "TagDynamic";
+ }
+
+ @Override
+ public boolean shouldRePool() {
+ return false;
+ }
+
+ @Override
+ public <T> RouterCache<T> pool(List<Invoker<T>> invokers) {
+
+ RouterCache<T> routerCache = new RouterCache<>();
+ ConcurrentHashMap<String, BitList<Invoker<T>>> addrPool = new
ConcurrentHashMap<>();
+
+ final TagRouterRule tagRouterRuleCopy = tagRouterRule;
+
+
+ if (tagRouterRuleCopy == null || !tagRouterRuleCopy.isValid() ||
!tagRouterRuleCopy.isEnabled()) {
+ BitList<Invoker<T>> noTagList = new BitList<>(invokers, true);
+
+ for (int index = 0; index < invokers.size(); index++) {
+ noTagList.addIndex(index);
+ }
+ addrPool.put(NO_TAG, noTagList);
+ routerCache.setAddrPool(addrPool);
+ return routerCache;
+ }
+
+ List<String> tagNames = tagRouterRuleCopy.getTagNames();
+ Map<String, List<String>> tagnameToAddresses =
tagRouterRuleCopy.getTagnameToAddresses();
+
+ for (String tag : tagNames) {
+ List<String> addresses = tagnameToAddresses.get(tag);
+ BitList<Invoker<T>> list = new BitList<>(invokers, true);
+
+ if (CollectionUtils.isEmpty(addresses)) {
+ list.addAll(invokers);
+ } else {
+ for (int index = 0; index < invokers.size(); index++) {
+ Invoker<T> invoker = invokers.get(index);
+ if (addressMatches(invoker.getUrl(), addresses)) {
+ list.addIndex(index);
+ }
+ }
+ }
+
+ addrPool.put(tag, list);
+ }
+
+ List<String> addresses = tagRouterRuleCopy.getAddresses();
+ BitList<Invoker<T>> noTagList = new BitList<>(invokers, true);
+
+ for (int index = 0; index < invokers.size(); index++) {
+ Invoker<T> invoker = invokers.get(index);
+ if (addressNotMatches(invoker.getUrl(), addresses)) {
+ noTagList.addIndex(index);
+ }
+ }
+ addrPool.put(NO_TAG, noTagList);
+ routerCache.setAddrPool(addrPool);
+ routerCache.setAddrMetadata(tagRouterRuleCopy);
+
+ return routerCache;
+ }
+
+ private boolean addressMatches(URL url, List<String> addresses) {
+ return addresses != null && checkAddressMatch(addresses,
url.getHost(), url.getPort());
+ }
+
+ private boolean addressNotMatches(URL url, List<String> addresses) {
+ return addresses == null || !checkAddressMatch(addresses,
url.getHost(), url.getPort());
+ }
+
+ private boolean checkAddressMatch(List<String> addresses, String host, int
port) {
+ for (String address : addresses) {
+ try {
+ if (NetUtils.matchIpExpression(address, host, port)) {
+ return true;
+ }
+ if ((ANYHOST_VALUE + ":" + port).equals(address)) {
+ return true;
+ }
+ } catch (UnknownHostException e) {
+ logger.error("The format of ip address is invalid in tag
route. Address :" + address, e);
+ } catch (Exception e) {
+ logger.error("The format of ip address is invalid in tag
route. Address :" + address, e);
+ }
+ }
+ return false;
+ }
+
+ public void setApplication(String app) {
+ this.application = app;
+ }
+
+ @Override
+ public <T> void notify(List<Invoker<T>> invokers) {
+ if (CollectionUtils.isEmpty(invokers)) {
+ return;
+ }
+
+ Invoker<T> invoker = invokers.get(0);
+ URL url = invoker.getUrl();
+ String providerApplication = url.getRemoteApplication();
+
+ if (StringUtils.isEmpty(providerApplication)) {
+ logger.error("TagRouter must getConfig from or subscribe to a
specific application, but the application " +
+ "in this TagRouter is not specified.");
+ return;
+ }
+
+ synchronized (this) {
+ if (!providerApplication.equals(application)) {
+ if (!StringUtils.isEmpty(application)) {
+ ruleRepository.removeListener(application + RULE_SUFFIX,
this);
+ }
+ String key = providerApplication + RULE_SUFFIX;
+ ruleRepository.addListener(key, this);
+ application = providerApplication;
+ String rawRule = ruleRepository.getRule(key,
DynamicConfiguration.DEFAULT_GROUP);
+ if (StringUtils.isNotEmpty(rawRule)) {
+ this.process(new ConfigChangedEvent(key,
DynamicConfiguration.DEFAULT_GROUP, rawRule));
+ }
+ }
+ }
+ pool(invokers);
+ }
+
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagDynamicStateRouterFactory.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagDynamicStateRouterFactory.java
new file mode 100644
index 0000000..cdcdb84
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagDynamicStateRouterFactory.java
@@ -0,0 +1,37 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.tag;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.extension.Activate;
+import org.apache.dubbo.rpc.cluster.RouterChain;
+import org.apache.dubbo.rpc.cluster.router.state.CacheableStateRouterFactory;
+import org.apache.dubbo.rpc.cluster.router.state.StateRouter;
+
+/**
+ * Tag router factory
+ */
+@Activate(order = 99)
+public class TagDynamicStateRouterFactory extends CacheableStateRouterFactory {
+
+ public static final String NAME = "tag-dynamic";
+
+ @Override
+ protected StateRouter createRouter(URL url, RouterChain chain) {
+ return new TagDynamicStateRouter(url, chain);
+ }
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagStaticStateRouter.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagStaticStateRouter.java
new file mode 100644
index 0000000..d15745f
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagStaticStateRouter.java
@@ -0,0 +1,155 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.tag;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.logger.Logger;
+import org.apache.dubbo.common.logger.LoggerFactory;
+import org.apache.dubbo.common.utils.CollectionUtils;
+import org.apache.dubbo.common.utils.StringUtils;
+import org.apache.dubbo.rpc.Invocation;
+import org.apache.dubbo.rpc.Invoker;
+import org.apache.dubbo.rpc.RpcException;
+import org.apache.dubbo.rpc.cluster.RouterChain;
+import org.apache.dubbo.rpc.cluster.router.state.AbstractStateRouter;
+import org.apache.dubbo.rpc.cluster.router.state.BitList;
+import org.apache.dubbo.rpc.cluster.router.state.RouterCache;
+import org.apache.dubbo.rpc.cluster.router.tag.model.TagRouterRule;
+import static org.apache.dubbo.common.constants.CommonConstants.TAG_KEY;
+
+/**
+ * TagStaticStateRouter, "application.tag-router"
+ */
+public class TagStaticStateRouter extends AbstractStateRouter {
+ public static final String NAME = "TAG_ROUTER";
+ private static final int TAG_ROUTER_DEFAULT_PRIORITY = 100;
+ private static final Logger logger =
LoggerFactory.getLogger(TagStaticStateRouter.class);
+ private static final String NO_TAG = "noTag";
+
+ private TagRouterRule tagRouterRule;
+
+ public TagStaticStateRouter(URL url, RouterChain chain) {
+ super(url, chain);
+ this.priority = TAG_ROUTER_DEFAULT_PRIORITY;
+ }
+
+ @Override
+ public URL getUrl() {
+ return url;
+ }
+
+ @Override
+ public <T> BitList<Invoker<T>> route(BitList<Invoker<T>> invokers,
RouterCache<T> routerCache, URL url, Invocation invocation)
+ throws RpcException {
+
+ String tag = StringUtils.isEmpty(invocation.getAttachment(TAG_KEY)) ?
url.getParameter(TAG_KEY) :
+ invocation.getAttachment(TAG_KEY);
+ if (StringUtils.isEmpty(tag)) {
+ tag = NO_TAG;
+ }
+
+ ConcurrentMap<String, BitList<Invoker<T>>> pool =
routerCache.getAddrPool();
+ BitList<Invoker<T>> res = pool.get(tag);
+ if (res == null) {
+ return invokers;
+ }
+ return invokers.intersect(res, invokers.getUnmodifiableList());
+ }
+
+ @Override
+ protected List<String> getTags(URL url, Invocation invocation) {
+ List<String> tags = new ArrayList<>();
+ String tag = StringUtils.isEmpty(invocation.getAttachment(TAG_KEY)) ?
url.getParameter(TAG_KEY) :
+ invocation.getAttachment(TAG_KEY);
+ if (StringUtils.isEmpty(tag)) {
+ tag = NO_TAG;
+ }
+ tags.add(tag);
+ return tags;
+ }
+
+ @Override
+ public boolean isRuntime() {
+ return tagRouterRule != null && tagRouterRule.isRuntime();
+ }
+
+ @Override
+ public boolean isEnable() {
+ return true;
+ }
+
+ @Override
+ public boolean isForce() {
+ // FIXME
+ return false;
+ }
+
+ @Override
+ public String getName() {
+ return "TagStatic";
+ }
+
+ @Override
+ public boolean shouldRePool() {
+ return false;
+ }
+
+ @Override
+ public <T> RouterCache<T> pool(List<Invoker<T>> invokers) {
+
+ RouterCache<T> routerCache = new RouterCache<>();
+ ConcurrentHashMap<String, BitList<Invoker<T>>> addrPool = new
ConcurrentHashMap<>();
+
+ for (int index = 0; index < invokers.size(); index++) {
+ Invoker<T> invoker = invokers.get(index);
+ String tag = invoker.getUrl().getParameter(TAG_KEY);
+ if (StringUtils.isEmpty(tag)) {
+ BitList<Invoker<T>> noTagList = addrPool.putIfAbsent(NO_TAG,
new BitList<>(invokers, true));
+ if (noTagList == null) {
+ noTagList = addrPool.get(NO_TAG);
+ }
+ noTagList.addIndex(index);
+ } else {
+ BitList<Invoker<T>> list = addrPool.putIfAbsent(tag, new
BitList<>(invokers, true));
+ if (list == null) {
+ list = addrPool.get(tag);
+ }
+ list.addIndex(index);
+ }
+ }
+
+ routerCache.setAddrPool(addrPool);
+
+ return routerCache;
+ }
+
+
+ @Override
+ public <T> void notify(List<Invoker<T>> invokers) {
+ if (CollectionUtils.isEmpty(invokers)) {
+ return;
+ }
+
+ pool(invokers);
+ }
+
+}
diff --git
a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagStaticStateRouterFactory.java
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagStaticStateRouterFactory.java
new file mode 100644
index 0000000..21fe348
--- /dev/null
+++
b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/router/tag/TagStaticStateRouterFactory.java
@@ -0,0 +1,37 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.dubbo.rpc.cluster.router.tag;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.extension.Activate;
+import org.apache.dubbo.rpc.cluster.RouterChain;
+import org.apache.dubbo.rpc.cluster.router.state.CacheableStateRouterFactory;
+import org.apache.dubbo.rpc.cluster.router.state.StateRouter;
+
+/**
+ * Tag router factory
+ */
+@Activate(order = 100)
+public class TagStaticStateRouterFactory extends CacheableStateRouterFactory {
+
+ public static final String NAME = "tag-static";
+
+ @Override
+ protected StateRouter createRouter(URL url, RouterChain chain) {
+ return new TagStaticStateRouter(url, chain);
+ }
+}
diff --git
a/dubbo-cluster/src/main/resources/META-INF/dubbo/internal/org.apache.dubbo.rpc.cluster.router.state.StateRouterFactory
b/dubbo-cluster/src/main/resources/META-INF/dubbo/internal/org.apache.dubbo.rpc.cluster.router.state.StateRouterFactory
new file mode 100644
index 0000000..c69b3ab
--- /dev/null
+++
b/dubbo-cluster/src/main/resources/META-INF/dubbo/internal/org.apache.dubbo.rpc.cluster.router.state.StateRouterFactory
@@ -0,0 +1,2 @@
+tag-dynamic=org.apache.dubbo.rpc.cluster.router.tag.TagDynamicStateRouterFactory
+tag-static=org.apache.dubbo.rpc.cluster.router.tag.TagStaticStateRouterFactory
\ No newline at end of file
diff --git
a/dubbo-common/src/main/java/org/apache/dubbo/common/threadpool/manager/DefaultExecutorRepository.java
b/dubbo-common/src/main/java/org/apache/dubbo/common/threadpool/manager/DefaultExecutorRepository.java
index b8eb41e..02bd6c9 100644
---
a/dubbo-common/src/main/java/org/apache/dubbo/common/threadpool/manager/DefaultExecutorRepository.java
+++
b/dubbo-common/src/main/java/org/apache/dubbo/common/threadpool/manager/DefaultExecutorRepository.java
@@ -16,20 +16,23 @@
*/
package org.apache.dubbo.common.threadpool.manager;
-import org.apache.dubbo.common.URL;
-import org.apache.dubbo.common.extension.ExtensionLoader;
-import org.apache.dubbo.common.logger.Logger;
-import org.apache.dubbo.common.logger.LoggerFactory;
-import org.apache.dubbo.common.threadpool.ThreadPool;
-import org.apache.dubbo.common.utils.NamedThreadFactory;
-
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
+import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.dubbo.common.URL;
+import org.apache.dubbo.common.extension.ExtensionLoader;
+import org.apache.dubbo.common.logger.Logger;
+import org.apache.dubbo.common.logger.LoggerFactory;
+import org.apache.dubbo.common.threadlocal.NamedInternalThreadFactory;
+import org.apache.dubbo.common.threadpool.ThreadPool;
+import org.apache.dubbo.common.utils.NamedThreadFactory;
import static org.apache.dubbo.common.constants.CommonConstants.CONSUMER_SIDE;
import static
org.apache.dubbo.common.constants.CommonConstants.EXECUTOR_SERVICE_COMPONENT_KEY;
@@ -59,13 +62,25 @@ public class DefaultExecutorRepository implements
ExecutorRepository {
private ConcurrentMap<String, ConcurrentMap<Integer, ExecutorService>>
data = new ConcurrentHashMap<>();
+ private ExecutorService poolRouterExecutor;
+
+ private static Ring<ExecutorService> executorServiceRing = new
Ring<ExecutorService>();
+
public DefaultExecutorRepository() {
for (int i = 0; i < DEFAULT_SCHEDULER_SIZE; i++) {
- ScheduledExecutorService scheduler =
Executors.newSingleThreadScheduledExecutor(new
NamedThreadFactory("Dubbo-framework-scheduler"));
+ ScheduledExecutorService scheduler =
Executors.newSingleThreadScheduledExecutor(
+ new NamedThreadFactory("Dubbo-framework-scheduler"));
scheduledExecutors.addItem(scheduler);
+
+ executorServiceRing.addItem(new ThreadPoolExecutor(1, 1,
+ 0L, TimeUnit.MILLISECONDS,
+ new LinkedBlockingQueue<Runnable>(1024), new
NamedInternalThreadFactory("Dubbo-state-router-loop", true)
+ , new ThreadPoolExecutor.AbortPolicy()));
}
//
// reconnectScheduledExecutor =
Executors.newSingleThreadScheduledExecutor(new
NamedThreadFactory("Dubbo-reconnect-scheduler"));
+ poolRouterExecutor = new ThreadPoolExecutor(1, 10, 0L,
TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>(1024),
+ new NamedInternalThreadFactory("Dubbo-state-router-pool-router",
true), new ThreadPoolExecutor.AbortPolicy());
serviceExporterExecutor = Executors.newScheduledThreadPool(1, new
NamedThreadFactory("Dubbo-exporter-scheduler"));
serviceDiscveryAddressNotificationExecutor =
Executors.newSingleThreadScheduledExecutor(new
NamedThreadFactory("Dubbo-SD-address-refresh"));
registryNotificationExecutor =
Executors.newSingleThreadScheduledExecutor(new
NamedThreadFactory("Dubbo-registry-notification"));
@@ -158,6 +173,11 @@ public class DefaultExecutorRepository implements
ExecutorRepository {
}
@Override
+ public ExecutorService nextExecutorExecutor() {
+ return executorServiceRing.pollItem();
+ }
+
+ @Override
public ScheduledExecutorService getServiceExporterExecutor() {
return serviceExporterExecutor;
}
@@ -185,4 +205,8 @@ public class DefaultExecutorRepository implements
ExecutorRepository {
return (ExecutorService)
ExtensionLoader.getExtensionLoader(ThreadPool.class).getAdaptiveExtension().getExecutor(url);
}
+ @Override
+ public ExecutorService getPoolRouterExecutor() {
+ return poolRouterExecutor;
+ }
}
diff --git
a/dubbo-common/src/main/java/org/apache/dubbo/common/threadpool/manager/ExecutorRepository.java
b/dubbo-common/src/main/java/org/apache/dubbo/common/threadpool/manager/ExecutorRepository.java
index d00a7bb..0a8b145 100644
---
a/dubbo-common/src/main/java/org/apache/dubbo/common/threadpool/manager/ExecutorRepository.java
+++
b/dubbo-common/src/main/java/org/apache/dubbo/common/threadpool/manager/ExecutorRepository.java
@@ -55,6 +55,8 @@ public interface ExecutorRepository {
*/
ScheduledExecutorService nextScheduledExecutor();
+ ExecutorService nextExecutorExecutor();
+
ScheduledExecutorService getServiceExporterExecutor();
ScheduledExecutorService getServiceDiscoveryAddressNotificationExecutor();
@@ -68,7 +70,6 @@ public interface ExecutorRepository {
*/
ScheduledExecutorService getRegistryNotificationExecutor();
-
/**
* Get the default shared threadpool.
*
@@ -76,4 +77,6 @@ public interface ExecutorRepository {
*/
ExecutorService getSharedExecutor();
+ ExecutorService getPoolRouterExecutor();
+
}
diff --git a/dubbo-dependencies-bom/pom.xml b/dubbo-dependencies-bom/pom.xml
index 675b7c7..5d59d8d 100644
--- a/dubbo-dependencies-bom/pom.xml
+++ b/dubbo-dependencies-bom/pom.xml
@@ -160,6 +160,7 @@
<hessian_lite_version>3.2.8</hessian_lite_version>
<swagger_version>1.5.19</swagger_version>
+ <roaringbitmap_version>0.9.0</roaringbitmap_version>
<metrics_version>2.0.1</metrics_version>
<sofa_registry_version>5.2.0</sofa_registry_version>
<gson_version>2.8.5</gson_version>
@@ -338,6 +339,11 @@
<artifactId>protobuf-java-util</artifactId>
<version>${protobuf-java_version}</version>
</dependency>
+ <dependency>
+ <groupId>org.roaringbitmap</groupId>
+ <artifactId>RoaringBitmap</artifactId>
+ <version>${roaringbitmap_version}</version>
+ </dependency>
<!-- Common Annotations API -->
<dependency>
<groupId>javax.annotation</groupId>
diff --git a/dubbo-distribution/dubbo-all/pom.xml
b/dubbo-distribution/dubbo-all/pom.xml
index 519fbc7..d3d4a18 100644
--- a/dubbo-distribution/dubbo-all/pom.xml
+++ b/dubbo-distribution/dubbo-all/pom.xml
@@ -331,6 +331,10 @@
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
</dependency>
+ <dependency>
+ <groupId>org.roaringbitmap</groupId>
+ <artifactId>RoaringBitmap</artifactId>
+ </dependency>
<!-- Temporarily add this part to exclude transitive dependency -->
<dependency>
diff --git
a/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/RpcException.java
b/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/RpcException.java
index 3228deb..d578151 100644
---
a/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/RpcException.java
+++
b/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/RpcException.java
@@ -38,6 +38,7 @@ public /**final**/ class RpcException extends
RuntimeException {
public static final int LIMIT_EXCEEDED_EXCEPTION = 7;
public static final int TIMEOUT_TERMINATE = 8;
public static final int REGISTRY_EXCEPTION = 9;
+ public static final int ROUTER_CACHE_NOT_BUILD = 10;
private static final long serialVersionUID = 7815426752583648734L;
/**
* RpcException cannot be extended, use error code for exception type to
keep compatibility