Add route part

This commit is contained in:
ken.lj 2018-08-26 13:37:51 +08:00
parent 6d95c4d0b5
commit 191cae18c3
29 changed files with 936 additions and 101 deletions

View File

@ -34,5 +34,10 @@
<artifactId>dubbo-rpc-api</artifactId>
<version>${project.parent.version}</version>
</dependency>
<dependency>
<groupId>org.apache.dubbo</groupId>
<artifactId>dubbo-config-dynamic</artifactId>
<version>${project.parent.version}</version>
</dependency>
</dependencies>
</project>

View File

@ -22,6 +22,7 @@ import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.RpcException;
import java.util.List;
import java.util.Map;
/**
* Router. (SPI, Prototype, ThreadSafe)
@ -32,7 +33,6 @@ import java.util.List;
* @see org.apache.dubbo.rpc.cluster.Directory#list(Invocation)
*/
public interface Router extends Comparable<Router> {
/**
* get the router url.
*
@ -51,4 +51,20 @@ public interface Router extends Comparable<Router> {
*/
<T> List<Invoker<T>> route(List<Invoker<T>> invokers, URL url, Invocation invocation) throws RpcException;
default <T> Map<String, List<Invoker<T>>> preRoute(List<Invoker<T>> invokers, URL url, Invocation invocation) throws RpcException {
return null;
}
default boolean isRuntime() {
return true;
}
default String getKey() {
return "";
}
default boolean isForce() {
return false;
}
}

View File

@ -0,0 +1,107 @@
/*
* 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;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.extension.ExtensionLoader;
import org.apache.dubbo.common.utils.CollectionUtils;
import org.apache.dubbo.config.dynamic.DynamicConfiguration;
import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.RpcInvocation;
import org.apache.dubbo.rpc.cluster.router.InvokerTreeCache;
import org.apache.dubbo.rpc.cluster.router.TreeNode;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
/**
*
*/
public class RouterChain<T> {
List<Invoker<T>> fullInvokers;
private InvokerTreeCache<T> treeCache;
List<Router> routers;
public static <T> RouterChain<T> buildChain(DynamicConfiguration dynamicConfiguration) {
List<RouterFactory> extensionFactories = ExtensionLoader.getExtensionLoader(RouterFactory.class).getActivateExtension(dynamicConfiguration.getUrl(), (String[]) null);
List<Router> routers = extensionFactories.stream().map(factory -> factory.getRouter(dynamicConfiguration)).collect(Collectors.toList());
return new RouterChain<>(routers);
}
public RouterChain(List<Router> routers) {
this.routers = routers;
}
public void addRouter(Router router) {
this.routers.add(router);
}
public void sort() {
}
/**
* @param methodInvokers
* @param url
* @param invocation TODO has no been used yet
*/
public void preRoute(Map<String, List<Invoker<T>>> methodInvokers, URL url, Invocation invocation) {
if (CollectionUtils.isEmpty(routers)) {
return;
}
TreeNode root = treeCache.buildTree();
Router router = routers.get(0);
methodInvokers.forEach((method, invokers) -> {
TreeNode<T> node = new TreeNode<>("method", method, invokers, true);
root.addChild(node);
Invocation invocation1 = new RpcInvocation(method, new Class<?>[0], new Object[0]);
routeeee(router, 1, root, router.preRoute(invokers, url, invocation1), url, invocation1);
});
}
private void routeeee(Router router, int i, TreeNode parentNode, Map<String, List<Invoker<T>>> invokers, URL url, Invocation invocation) {
invokers.forEach((routerValue, list) -> {
TreeNode<T> node = new TreeNode<>(router.getKey(), routerValue, list, router.isForce());
parentNode.addChild(node);
Router nextRouter = routers.get(i);
if (CollectionUtils.isNotEmpty(list)) {
routeeee(nextRouter, i + 1, node, nextRouter.preRoute(list, url, invocation), url, invocation);
}
});
}
public List<Invoker<T>> route(List<Invoker<T>> invokers, URL url, Invocation invocation) {
// TODO will calculate every time
List<Invoker<T>> finalInvokers = treeCache.getInvokers(treeCache.getTree(), url, invocation);
for (Router router : routers) {
if (router.isRuntime()) {
finalInvokers = router.route(finalInvokers, url, invocation);
}
}
return finalInvokers;
}
public void notifyFullInvokers(Map<String, List<Invoker<T>>> invokers, URL url) {
preRoute(invokers, url, null);
}
}

View File

@ -19,6 +19,7 @@ package org.apache.dubbo.rpc.cluster;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.extension.Adaptive;
import org.apache.dubbo.common.extension.SPI;
import org.apache.dubbo.config.dynamic.DynamicConfiguration;
import org.apache.dubbo.rpc.Invocation;
/**
@ -41,4 +42,7 @@ public interface RouterFactory {
@Adaptive("protocol")
Router getRouter(URL url);
default Router getRouter(DynamicConfiguration dynamicConfiguration) {
return null;
}
}

View File

@ -16,9 +16,7 @@
*/
package org.apache.dubbo.rpc.cluster.directory;
import org.apache.dubbo.common.Constants;
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.rpc.Invocation;
@ -26,11 +24,9 @@ import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.Directory;
import org.apache.dubbo.rpc.cluster.Router;
import org.apache.dubbo.rpc.cluster.RouterFactory;
import org.apache.dubbo.rpc.cluster.router.MockInvokersSelector;
import org.apache.dubbo.rpc.cluster.RouterChain;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
/**
@ -48,22 +44,22 @@ public abstract class AbstractDirectory<T> implements Directory<T> {
private volatile URL consumerUrl;
private volatile List<Router> routers;
protected RouterChain<T> routerChain;
public AbstractDirectory(URL url) {
this(url, null);
}
public AbstractDirectory(URL url, List<Router> routers) {
this(url, url, routers);
public AbstractDirectory(URL url, RouterChain<T> routerChain) {
this(url, url, routerChain);
}
public AbstractDirectory(URL url, URL consumerUrl, List<Router> routers) {
public AbstractDirectory(URL url, URL consumerUrl, RouterChain<T> routerChain) {
if (url == null)
throw new IllegalArgumentException("url == null");
this.url = url;
this.consumerUrl = consumerUrl;
setRouters(routers);
setRouterChain(routerChain);
}
@Override
@ -71,19 +67,16 @@ public abstract class AbstractDirectory<T> implements Directory<T> {
if (destroyed) {
throw new RpcException("Directory already destroyed .url: " + getUrl());
}
List<Invoker<T>> invokers = doList(invocation);
List<Router> localRouters = this.routers; // local reference
if (localRouters != null && !localRouters.isEmpty()) {
for (Router router : localRouters) {
try {
if (router.getUrl() == null || router.getUrl().getParameter(Constants.RUNTIME_KEY, false)) {
invokers = router.route(invokers, getConsumerUrl(), invocation);
}
} catch (Throwable t) {
logger.error("Failed to execute router: " + getUrl() + ", cause: " + t.getMessage(), t);
}
}
try {
// runtime routers will be executed.
routerChain.route(invokers, getConsumerUrl(), invocation);
} catch (Throwable t) {
logger.error("Failed to execute router: " + getUrl() + ", cause: " + t.getMessage(), t);
}
return invokers;
}
@ -92,23 +85,24 @@ public abstract class AbstractDirectory<T> implements Directory<T> {
return url;
}
public List<Router> getRouters() {
return routers;
public RouterChain<T> getRouterChain() {
return routerChain;
}
protected void setRouters(List<Router> routers) {
public void setRouterChain(RouterChain<T> routerChain) {
this.routerChain = routerChain;
}
protected void addRouters(List<Router> routers) {
// copy list
routers = routers == null ? new ArrayList<Router>() : new ArrayList<Router>(routers);
// append url router
String routerkey = url.getParameter(Constants.ROUTER_KEY);
if (routerkey != null && routerkey.length() > 0) {
RouterFactory routerFactory = ExtensionLoader.getExtensionLoader(RouterFactory.class).getExtension(routerkey);
routers.add(routerFactory.getRouter(url));
}
// append mock invoker selector
routers.add(new MockInvokersSelector());
Collections.sort(routers);
this.routers = routers;
routers = routers == null ? new ArrayList<>() : new ArrayList<>(routers);
routers.forEach(this::addRouter);
}
protected void addRouter(Router router) {
routerChain.addRouter(router);
// FIXME append mock invoker selector
// routerChain.add(new MockInvokersSelector());
}
public URL getConsumerUrl() {

View File

@ -20,13 +20,15 @@ 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.Router;
import org.apache.dubbo.rpc.cluster.RouterChain;
import java.util.List;
/**
* StaticDirectory
*
* FIXME We should add separate router rules for StaticDirectory, not use RouterChain in AbstractDirectory. For example, machineRoom rule only. Ask LVS team for help.
*
*/
public class StaticDirectory<T> extends AbstractDirectory<T> {
@ -36,16 +38,16 @@ public class StaticDirectory<T> extends AbstractDirectory<T> {
this(null, invokers, null);
}
public StaticDirectory(List<Invoker<T>> invokers, List<Router> routers) {
this(null, invokers, routers);
public StaticDirectory(List<Invoker<T>> invokers, RouterChain<T> routerChain) {
this(null, invokers, routerChain);
}
public StaticDirectory(URL url, List<Invoker<T>> invokers) {
this(url, invokers, null);
}
public StaticDirectory(URL url, List<Invoker<T>> invokers, List<Router> routers) {
super(url == null && invokers != null && !invokers.isEmpty() ? invokers.get(0).getUrl() : url, routers);
public StaticDirectory(URL url, List<Invoker<T>> invokers, RouterChain<T> routerChain) {
super(url == null && invokers != null && !invokers.isEmpty() ? invokers.get(0).getUrl() : url, routerChain);
if (invokers == null || invokers.isEmpty())
throw new IllegalArgumentException("invokers == null");
this.invokers = invokers;

View File

@ -0,0 +1,71 @@
/*
* 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;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import java.util.Collections;
import java.util.List;
/**
*
*/
public class InvokerTreeCache<T> {
TreeNode<T> tree;
public TreeNode buildTree() {
tree = new TreeNode<>();
return tree;
}
public List<Invoker<T>> getInvokers(TreeNode<T> node, URL url, Invocation invocation) {
if (node.getChildren() == null) {
return node.getInvokers();
}
if (tree.getChildren().size() == 1) {
return getInvokers(tree.getChildren().get(0), url, invocation);
}
TreeNode<T> failoverNode = null;
for (TreeNode<T> n : tree.getChildren()) {
String key = n.getConditionKey();
if (TreeNode.FAILOVER_KEY.equals(key)) {
failoverNode = n;
continue;
}
if (n.getConditionValue().equals(invocation.getAttachment(key, url.getParameter(key)))) {
if (n.getInvokers() != null || (n.getInvokers() == null && n.isForce()))
return getInvokers(n, url, invocation);
}
}
if (failoverNode != null) {
return getInvokers(failoverNode, url, invocation);
}
return Collections.emptyList();
}
public TreeNode getTree() {
return tree;
}
}

View File

@ -18,13 +18,16 @@ package org.apache.dubbo.rpc.cluster.router;
import org.apache.dubbo.common.Constants;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.utils.CollectionUtils;
import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.Router;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
* A specific Router designed to realize mock feature.
@ -49,6 +52,36 @@ public class MockInvokersSelector implements Router {
return invokers;
}
@Override
public <T> Map<String, List<Invoker<T>>> preRoute(List<Invoker<T>> invokers, URL url, Invocation invocation) throws RpcException {
Map<String, List<Invoker<T>>> map = new HashMap<>();
if (CollectionUtils.isEmpty(invokers)) {
return map;
}
if (isRuntime()) {
map.put(TreeNode.FAILOVER_KEY, invokers);
return map;
}
return map;
}
@Override
public boolean isRuntime() {
return true;
}
@Override
public String getKey() {
return TreeNode.FAILOVER_KEY;
}
@Override
public boolean isForce() {
return false;
}
private <T> List<Invoker<T>> getMockedInvokers(final List<Invoker<T>> invokers) {
if (!hasMockProviders(invokers)) {
return null;

View File

@ -0,0 +1,103 @@
/*
* 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;
import org.apache.dubbo.common.utils.CollectionUtils;
import org.apache.dubbo.rpc.Invoker;
import java.util.ArrayList;
import java.util.List;
/**
*
*/
public class TreeNode<T> {
public static final String FAILOVER_KEY = "failover";
private String conditionKey;
private String conditionValue;
private boolean force;
private List<Invoker<T>> invokers;
private List<TreeNode<T>> children;
public TreeNode() {
this.children = new ArrayList<>();
}
public TreeNode(String conditionKey, String conditionValue, List<Invoker<T>> invokers, boolean force) {
this.conditionKey = conditionKey;
this.conditionValue = conditionValue;
this.invokers = invokers;
this.force = force;
this.children = new ArrayList<>();
}
public void traverse() {
}
public void addChild(TreeNode child) {
children.add(child);
}
public boolean isLeaf() {
if (CollectionUtils.isEmpty(children)) {
return true;
}
return false;
}
public String getConditionKey() {
return conditionKey;
}
public void setConditionKey(String conditionKey) {
this.conditionKey = conditionKey;
}
public String getConditionValue() {
return conditionValue;
}
public void setConditionValue(String conditionValue) {
this.conditionValue = conditionValue;
}
public List<Invoker<T>> getInvokers() {
return invokers;
}
public void setInvokers(List<Invoker<T>> invokers) {
this.invokers = invokers;
}
public List<TreeNode<T>> getChildren() {
return children;
}
public void setChildren(List<TreeNode<T>> children) {
this.children = children;
}
public boolean isForce() {
return force;
}
public void setForce(boolean force) {
this.force = force;
}
}

View File

@ -20,6 +20,7 @@ import org.apache.dubbo.common.Constants;
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.NetUtils;
import org.apache.dubbo.common.utils.StringUtils;
import org.apache.dubbo.common.utils.UrlUtils;
@ -27,6 +28,7 @@ import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.Router;
import org.apache.dubbo.rpc.cluster.router.TreeNode;
import java.text.ParseException;
import java.util.ArrayList;
@ -45,19 +47,26 @@ import java.util.regex.Pattern;
public class ConditionRouter implements Router, Comparable<Router> {
private static final Logger logger = LoggerFactory.getLogger(ConditionRouter.class);
private static Pattern ROUTE_PATTERN = Pattern.compile("([&!=,]*)\\s*([^&!=,\\s]+)");
private final URL url;
private final int priority;
private final boolean force;
private final Map<String, MatchPair> whenCondition;
private final Map<String, MatchPair> thenCondition;
protected static Pattern ROUTE_PATTERN = Pattern.compile("([&!=,]*)\\s*([^&!=,\\s]+)");
protected URL url;
protected int priority;
protected boolean force;
protected Map<String, MatchPair> whenCondition;
protected Map<String, MatchPair> thenCondition;
protected ConditionRouter() {
}
public ConditionRouter(URL url) {
this.url = url;
this.priority = url.getParameter(Constants.PRIORITY_KEY, 0);
this.force = url.getParameter(Constants.FORCE_KEY, false);
init(url.getParameterAndDecoded(Constants.RULE_KEY));
}
protected void init(String rule) {
try {
String rule = url.getParameterAndDecoded(Constants.RULE_KEY);
if (rule == null || rule.trim().length() == 0) {
throw new IllegalArgumentException("Illegal route rule!");
}
@ -174,6 +183,37 @@ public class ConditionRouter implements Router, Comparable<Router> {
return invokers;
}
@Override
public <T> Map<String, List<Invoker<T>>> preRoute(List<Invoker<T>> invokers, URL url, Invocation invocation) throws RpcException {
Map<String, List<Invoker<T>>> map = new HashMap<>();
if (CollectionUtils.isEmpty(invokers)) {
return map;
}
if (isRuntime()) {
map.put(TreeNode.FAILOVER_KEY, invokers);
return map;
}
return map;
}
@Override
public boolean isRuntime() {
// We always return true for previously defined Router, don't support cache.
return this.url.getParameter(Constants.RUNTIME_KEY, false);
}
@Override
public String getKey() {
return TreeNode.FAILOVER_KEY;
}
@Override
public boolean isForce() {
return force;
}
@Override
public URL getUrl() {
return url;
@ -229,7 +269,7 @@ public class ConditionRouter implements Router, Comparable<Router> {
return result;
}
private static final class MatchPair {
protected static final class MatchPair {
final Set<String> matches = new HashSet<String>();
final Set<String> mismatches = new HashSet<String>();

View File

@ -0,0 +1,50 @@
/*
* 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.condition.config;
/**
*
*/
public class ConditionRouterRule {
private String ruleBody;
private boolean runtime;
private boolean force;
public String getRuleBody() {
return ruleBody;
}
public void setRuleBody(String ruleBody) {
this.ruleBody = ruleBody;
}
public boolean isRuntime() {
return runtime;
}
public void setRuntime(boolean runtime) {
this.runtime = runtime;
}
public boolean isForce() {
return force;
}
public void setForce(boolean force) {
this.force = force;
}
}

View File

@ -0,0 +1,30 @@
/*
* 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.condition.config;
/**
*
*/
public class ConditionRuleParser {
public static ConditionRouterRule parse(String rawRule) {
ConditionRouterRule conditionRouterRule = new ConditionRouterRule();
conditionRouterRule.setRuleBody("host!=10.20.153.10,10.20.153.11=>");
return conditionRouterRule;
}
}

View File

@ -0,0 +1,102 @@
/*
* 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.condition.config;
import org.apache.dubbo.common.Constants;
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.config.dynamic.ConfigChangeEvent;
import org.apache.dubbo.config.dynamic.ConfigurationListener;
import org.apache.dubbo.config.dynamic.DynamicConfiguration;
import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.router.TreeNode;
import org.apache.dubbo.rpc.cluster.router.condition.ConditionRouter;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
*
*/
public class ConfigConditionRouter extends ConditionRouter implements ConfigurationListener {
private static final Logger logger = LoggerFactory.getLogger(ConfigConditionRouter.class);
private DynamicConfiguration configuration;
private ConditionRouterRule routerRule;
public ConfigConditionRouter(URL url) {
super(url);
}
public ConfigConditionRouter(DynamicConfiguration configuration) {
this.configuration = configuration;
this.priority = -2;
this.force = false;
try {
String app = configuration.getUrl().getParameter(Constants.APPLICATION_KEY);
String rawRule = configuration.getConfig(app + Constants.ROUTERS_SUFFIX, "dubbo", this);
ConditionRouterRule routerRule = ConditionRuleParser.parse(rawRule);
init(routerRule.getRuleBody());
} catch (Exception e) {
throw new IllegalStateException(e.getMessage(), e);
}
}
@Override
public void process(ConfigChangeEvent event) {
String rawRule = event.getNewValue();
ConditionRouterRule routerRule = ConditionRuleParser.parse(rawRule);
init(routerRule.getRuleBody());
}
@Override
public <T> Map<String, List<Invoker<T>>> preRoute(List<Invoker<T>> invokers, URL url, Invocation invocation) throws RpcException {
Map<String, List<Invoker<T>>> map = new HashMap<>();
if (CollectionUtils.isEmpty(invokers)) {
return map;
}
if (isRuntime()) {
map.put(TreeNode.FAILOVER_KEY, invokers);
return map;
}
map.put(TreeNode.FAILOVER_KEY, route(invokers, url, invocation));
return map;
}
@Override
public boolean isRuntime() {
return routerRule.isRuntime();
}
@Override
public String getKey() {
return "";
}
@Override
public boolean isForce() {
return routerRule.isForce();
}
}

View File

@ -0,0 +1,39 @@
/*
* 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.condition.config;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.extension.Activate;
import org.apache.dubbo.config.dynamic.DynamicConfiguration;
import org.apache.dubbo.rpc.cluster.Router;
import org.apache.dubbo.rpc.cluster.RouterFactory;
/**
*
*/
@Activate
public class ConfigConditionRouterFactory implements RouterFactory {
@Override
public Router getRouter(URL url) {
return null;
}
@Override
public Router getRouter(DynamicConfiguration dynamicConfiguration) {
return new ConfigConditionRouter(dynamicConfiguration);
}
}

View File

@ -0,0 +1,153 @@
/*
* 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.group;
import org.apache.dubbo.common.Constants;
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.utils.CollectionUtils;
import org.apache.dubbo.common.utils.StringUtils;
import org.apache.dubbo.config.dynamic.ConfigChangeEvent;
import org.apache.dubbo.config.dynamic.ConfigurationListener;
import org.apache.dubbo.config.dynamic.DynamicConfiguration;
import org.apache.dubbo.config.dynamic.DynamicConfigurationFactory;
import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.Router;
import org.apache.dubbo.rpc.cluster.router.TreeNode;
import org.apache.dubbo.rpc.cluster.router.group.model.GroupRouterRule;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Predicate;
import java.util.stream.Collectors;
/**
*
*/
public class GroupRouter implements Router, Comparable<Router>, ConfigurationListener {
private static final Logger logger = LoggerFactory.getLogger(GroupRouter.class);
private static final String GROUPRULE_DATAID = "global.routers";
private static final String ROUTE_GROUP = "route.group";
private static final String ROUTE_FAILOVER = "route.failover";
private int priority;
private URL url;
private DynamicConfiguration configuration;
private GroupRouterRule groupRouterRule;
public GroupRouter(URL url) {
this(ExtensionLoader.getExtensionLoader(DynamicConfigurationFactory.class).getAdaptiveExtension().getDynamicConfiguration(url));
this.url = url;
}
public GroupRouter(DynamicConfiguration configuration) {
this.priority = -1;
this.configuration = configuration;
init();
}
public void init() {
String rawRule = configuration.getConfig(GROUPRULE_DATAID, "dubbo", this);
this.groupRouterRule = GroupRuleParser.parse(rawRule);
}
@Override
public void process(ConfigChangeEvent event) {
String rawRule = event.getNewValue();
// remove, set groupRouterRule to null
// change, update groupRouterRule
}
@Override
public URL getUrl() {
return url;
}
@Override
public <T> List<Invoker<T>> route(List<Invoker<T>> invokers, URL url, Invocation invocation) throws RpcException {
if (CollectionUtils.isEmpty(invokers)) {
return invokers;
}
List<Invoker<T>> result = invokers;
String routeGroup = StringUtils.isEmpty(invocation.getAttachment(ROUTE_GROUP)) ? url.getParameter(ROUTE_GROUP) : invocation.getAttachment(ROUTE_GROUP);
if (StringUtils.isNotEmpty(routeGroup)) {
String providerApp = invokers.get(0).getUrl().getParameter(Constants.APPLICATION_KEY);
List<String> ips = groupRouterRule.filter(routeGroup, providerApp);
if (CollectionUtils.isNotEmpty(ips)) {
result = filterInvoker(invokers, invoker -> ipMatches(invoker.getUrl(), ips));
} else {
result = filterInvoker(invokers, invoker -> invoker.getUrl().getParameter(ROUTE_GROUP).equals(routeGroup));
}
}
if (StringUtils.isEmpty(routeGroup) || (CollectionUtils.isEmpty(result) && url.getParameter(ROUTE_FAILOVER, true))) {
result = filterInvoker(invokers, invoker -> StringUtils.isEmpty(invoker.getUrl().getParameter(ROUTE_GROUP)));
}
return result;
}
@Override
public <T> Map<String, List<Invoker<T>>> preRoute(List<Invoker<T>> invokers, URL url, Invocation invocation) throws RpcException {
Map<String, List<Invoker<T>>> map = new HashMap<>();
if (CollectionUtils.isEmpty(invokers)) {
return map;
}
if (isRuntime()) {
map.put(TreeNode.FAILOVER_KEY, invokers);
return map;
}
// FIXME Consider groupRouterRule
invokers.forEach(invoker -> {
String routeGroup = invoker.getUrl().getParameter(ROUTE_GROUP);
List<Invoker<T>> subInvokers = map.computeIfAbsent(routeGroup, k -> new ArrayList<>());
subInvokers.add(invoker);
});
return map;
}
@Override
public boolean isRuntime() {
return false;
}
public boolean isRuntime(Invocation invocation) {
return true;
}
private <T> List<Invoker<T>> filterInvoker(List<Invoker<T>> invokers, Predicate<Invoker<T>> predicate) {
return invokers.stream()
.filter(predicate)
.collect(Collectors.toList());
}
private boolean ipMatches(URL url, List<String> ips) {
return ips.contains(url.getHost());
}
@Override
public int compareTo(Router o) {
return 0;
}
}

View File

@ -0,0 +1,41 @@
/*
* 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.group;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.extension.Activate;
import org.apache.dubbo.config.dynamic.DynamicConfiguration;
import org.apache.dubbo.rpc.cluster.Router;
import org.apache.dubbo.rpc.cluster.RouterFactory;
/**
*
*/
@Activate
public class GroupRouterFactory implements RouterFactory {
public static final String NAME = "group";
@Override
public Router getRouter(URL url) {
return new GroupRouter(url);
}
public Router getRouter(DynamicConfiguration dynamicConfiguration) {
return new GroupRouter(dynamicConfiguration);
}
}

View File

@ -0,0 +1,29 @@
/*
* 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.group;
import org.apache.dubbo.rpc.cluster.router.group.model.GroupRouterRule;
/**
*
*/
public class GroupRuleParser {
public static GroupRouterRule parse(String rawRule) {
return null;
}
}

View File

@ -14,11 +14,15 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.dubbo.config.dynamic.parser.model;
package org.apache.dubbo.rpc.cluster.router.group.model;
import java.util.List;
/**
*
*/
public class RouterRule {
public class GroupRouterRule {
public List<String> filter(String routeGroup, String app) {
return null;
}
}

View File

@ -20,11 +20,13 @@ import org.apache.dubbo.common.Constants;
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.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.RpcContext;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.Router;
import org.apache.dubbo.rpc.cluster.router.TreeNode;
import javax.script.Bindings;
import javax.script.Compilable;
@ -34,6 +36,7 @@ import javax.script.ScriptEngineManager;
import javax.script.ScriptException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@ -114,6 +117,37 @@ public class ScriptRouter implements Router {
}
}
@Override
public <T> Map<String, List<Invoker<T>>> preRoute(List<Invoker<T>> invokers, URL url, Invocation invocation) throws RpcException {
Map<String, List<Invoker<T>>> map = new HashMap<>();
if (CollectionUtils.isEmpty(invokers)) {
return map;
}
if (isRuntime()) {
map.put(TreeNode.FAILOVER_KEY, invokers);
return map;
}
return map;
}
@Override
public boolean isRuntime() {
// ignore config in url
return true;
}
@Override
public String getKey() {
return TreeNode.FAILOVER_KEY;
}
@Override
public boolean isForce() {
return url.getParameter(Constants.FORCE_KEY, false);
}
@Override
public int compareTo(Router o) {
if (o == null || o.getClass() != ScriptRouter.class) {

View File

@ -1,3 +1,5 @@
file=org.apache.dubbo.rpc.cluster.router.file.FileRouterFactory
script=org.apache.dubbo.rpc.cluster.router.script.ScriptRouterFactory
condition=org.apache.dubbo.rpc.cluster.router.condition.ConditionRouterFactory
condition=org.apache.dubbo.rpc.cluster.router.condition.ConditionRouterFactory
configcondition=org.apache.dubbo.rpc.cluster.router.condition.config.ConfigConditionRouterFactory
group=org.apache.dubbo.rpc.cluster.router.group.GroupRouterFactory

View File

@ -27,10 +27,10 @@ import org.apache.dubbo.rpc.RpcInvocation;
import org.apache.dubbo.rpc.RpcResult;
import org.apache.dubbo.rpc.cluster.Directory;
import org.apache.dubbo.rpc.cluster.LoadBalance;
import org.apache.dubbo.rpc.cluster.RouterChain;
import org.apache.dubbo.rpc.cluster.RouterFactory;
import org.apache.dubbo.rpc.cluster.directory.StaticDirectory;
import org.apache.dubbo.rpc.cluster.support.AbstractClusterInvoker;
import org.junit.Assert;
import org.junit.Before;
import org.junit.BeforeClass;
@ -159,7 +159,8 @@ public class FileRouterEngineTest {
}
private void initDic(URL url) {
dic = new StaticDirectory<FileRouterEngineTest>(url, invokers, Arrays.asList(routerFactory.getRouter(url)));
RouterChain routerChain = new RouterChain(Arrays.asList(routerFactory.getRouter(url)));
dic = new StaticDirectory<FileRouterEngineTest>(url, invokers, routerChain);
}
static class MockClusterInvoker<T> extends AbstractClusterInvoker<T> {

View File

@ -29,14 +29,13 @@ import org.apache.dubbo.config.dynamic.ConfigChangeType;
import org.apache.dubbo.config.dynamic.ConfigType;
import org.apache.dubbo.config.dynamic.ConfigurationListener;
import org.apache.dubbo.config.dynamic.DynamicConfiguration;
import org.apache.dubbo.config.dynamic.parser.ConfigParser;
import org.apache.dubbo.registry.NotifyListener;
import org.apache.dubbo.registry.Registry;
import org.apache.dubbo.registry.integration.parser.ConfigParser;
import org.apache.dubbo.rpc.Invocation;
import org.apache.dubbo.rpc.Invoker;
import org.apache.dubbo.rpc.Protocol;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.RpcInvocation;
import org.apache.dubbo.rpc.cluster.Cluster;
import org.apache.dubbo.rpc.cluster.Configurator;
import org.apache.dubbo.rpc.cluster.ConfiguratorFactory;
@ -215,13 +214,10 @@ public class RegistryDirectory<T> extends AbstractDirectory<T> implements Notify
List<URL> configuratorUrls = new ArrayList<>();
List<URL> dynamicConfiguratorUrls = new ArrayList<>();
List<URL> appDynamicConfiguratorUrls = new ArrayList<>();
List<URL> dynamicRouterUrls = new ArrayList<>();
for (URL url : urls) {
String protocol = url.getProtocol();
String category = url.getParameter(Constants.CATEGORY_KEY, Constants.DEFAULT_CATEGORY);
if (Constants.DYNAMIC_ROUTERS_CATEGORY.equals(category)) {
dynamicRouterUrls.add(url);
} else if (Constants.ROUTERS_CATEGORY.equals(category)
if (Constants.ROUTERS_CATEGORY.equals(category)
|| Constants.ROUTE_PROTOCOL.equals(protocol)) {
routerUrls.add(url);
} else if (Constants.DYNAMIC_CONFIGURATORS_CATEGORY.equals(category)) {
@ -252,9 +248,7 @@ public class RegistryDirectory<T> extends AbstractDirectory<T> implements Notify
// routers
if (!routerUrls.isEmpty()) {
List<Router> routers = toRouters(routerUrls);
if (routers != null) { // null - do nothing
setRouters(routers);
}
addRouters(routers);
}
List<Configurator> localConfigurators = this.configurators; // local reference
// merge override parameters
@ -311,6 +305,8 @@ public class RegistryDirectory<T> extends AbstractDirectory<T> implements Notify
}
this.methodInvokerMap = multiGroup ? toMergeMethodInvokerMap(newMethodInvokerMap) : newMethodInvokerMap;
this.urlInvokerMap = newUrlInvokerMap;
// Route and build cache
routerChain.notifyFullInvokers(methodInvokerMap, getConsumerUrl());
try {
destroyUnusedInvokers(oldUrlInvokerMap, newUrlInvokerMap); // Close the unused Invoker
} catch (Exception e) {
@ -511,19 +507,6 @@ public class RegistryDirectory<T> extends AbstractDirectory<T> implements Notify
return providerUrl;
}
private List<Invoker<T>> route(List<Invoker<T>> invokers, String method) {
Invocation invocation = new RpcInvocation(method, new Class<?>[0], new Object[0]);
List<Router> routers = getRouters();
if (routers != null) {
for (Router router : routers) {
if (router.getUrl() != null) {
invokers = router.route(invokers, getConsumerUrl(), invocation);
}
}
}
return invokers;
}
/**
* Transform the invokers list into a mapping relationship with a method
*
@ -556,17 +539,7 @@ public class RegistryDirectory<T> extends AbstractDirectory<T> implements Notify
invokersList.add(invoker);
}
}
List<Invoker<T>> newInvokersList = route(invokersList, null);
newMethodInvokerMap.put(Constants.ANY_VALUE, newInvokersList);
if (serviceMethods != null && serviceMethods.length > 0) {
for (String method : serviceMethods) {
List<Invoker<T>> methodInvokers = newMethodInvokerMap.get(method);
if (methodInvokers == null || methodInvokers.isEmpty()) {
methodInvokers = newInvokersList;
}
newMethodInvokerMap.put(method, route(methodInvokers, method));
}
}
newMethodInvokerMap.put(Constants.ANY_VALUE, invokersList);
// sort and unmodifiable
for (String method : new HashSet<String>(newMethodInvokerMap.keySet())) {
List<Invoker<T>> methodInvokers = newMethodInvokerMap.get(method);

View File

@ -31,11 +31,11 @@ import org.apache.dubbo.config.dynamic.ConfigType;
import org.apache.dubbo.config.dynamic.ConfigurationListener;
import org.apache.dubbo.config.dynamic.DynamicConfiguration;
import org.apache.dubbo.config.dynamic.DynamicConfigurationFactory;
import org.apache.dubbo.config.dynamic.parser.ConfigParser;
import org.apache.dubbo.registry.NotifyListener;
import org.apache.dubbo.registry.Registry;
import org.apache.dubbo.registry.RegistryFactory;
import org.apache.dubbo.registry.RegistryService;
import org.apache.dubbo.registry.integration.parser.ConfigParser;
import org.apache.dubbo.registry.support.ProviderConsumerRegTable;
import org.apache.dubbo.rpc.Exporter;
import org.apache.dubbo.rpc.Invoker;
@ -44,6 +44,7 @@ import org.apache.dubbo.rpc.ProxyFactory;
import org.apache.dubbo.rpc.RpcException;
import org.apache.dubbo.rpc.cluster.Cluster;
import org.apache.dubbo.rpc.cluster.Configurator;
import org.apache.dubbo.rpc.cluster.RouterChain;
import org.apache.dubbo.rpc.protocol.InvokerWrapper;
import java.util.ArrayList;
@ -350,6 +351,7 @@ public class RegistryProtocol implements Protocol {
directory.setRegistry(registry);
directory.setProtocol(protocol);
directory.setDynamicConfiguration(dynamicConfiguration);
directory.setRouterChain(RouterChain.buildChain(dynamicConfiguration));
// all attributes of REFER_KEY
Map<String, String> parameters = new HashMap<String, String>(directory.getUrl().getParameters());
URL subscribeUrl = new URL(Constants.CONSUMER_PROTOCOL, parameters.remove(Constants.REGISTER_IP_KEY), 0, type.getName(), parameters);

View File

@ -14,14 +14,14 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.dubbo.config.dynamic.parser;
package org.apache.dubbo.registry.integration.parser;
import org.apache.dubbo.common.Constants;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.utils.StringUtils;
import org.apache.dubbo.config.dynamic.parser.model.ConfigItem;
import org.apache.dubbo.config.dynamic.parser.model.ConfiguratorConfig;
import org.apache.dubbo.config.dynamic.parser.model.ConfiguratorRule;
import org.apache.dubbo.registry.integration.parser.model.ConfigItem;
import org.apache.dubbo.registry.integration.parser.model.ConfiguratorConfig;
import org.apache.dubbo.registry.integration.parser.model.ConfiguratorRule;
import org.yaml.snakeyaml.TypeDescription;
import org.yaml.snakeyaml.Yaml;
import org.yaml.snakeyaml.constructor.Constructor;

View File

@ -14,7 +14,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.dubbo.config.dynamic.parser.model;
package org.apache.dubbo.registry.integration.parser.model;
import java.util.List;

View File

@ -14,7 +14,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.dubbo.config.dynamic.parser.model;
package org.apache.dubbo.registry.integration.parser.model;
import java.util.List;

View File

@ -14,7 +14,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.dubbo.config.dynamic.parser.model;
package org.apache.dubbo.registry.integration.parser.model;
import java.util.Map;

View File

@ -14,13 +14,13 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.dubbo.config;
package org.apache.dubbo.registry;
import org.apache.dubbo.common.Constants;
import org.apache.dubbo.common.URL;
import org.apache.dubbo.config.dynamic.parser.ConfigParser;
import org.apache.dubbo.config.dynamic.parser.model.ConfigItem;
import org.apache.dubbo.config.dynamic.parser.model.ConfiguratorConfig;
import org.apache.dubbo.registry.integration.parser.ConfigParser;
import org.apache.dubbo.registry.integration.parser.model.ConfigItem;
import org.apache.dubbo.registry.integration.parser.model.ConfiguratorConfig;
import org.junit.Assert;
import org.junit.Test;
import org.yaml.snakeyaml.TypeDescription;

View File

@ -14,7 +14,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.dubbo.config;
package org.apache.dubbo.registry;
import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;