服务注册与发布,书接上回,接着分析dubbo的服务注册与发布。
首先了解一下dubbo怎么解析自定义的标签的。
第一种:xml
在 dubbo-config 模块的 dubbo-config-spring 下面,找到resources下面的meta-inf下面的spring.handlers文件,该文件主要是指定标签校验为本地实现类,该文件里面的内容也是以key、value存储的,key就是我们在xml里面必须要添加的schema头,value就是解析的类,然后再去定义一个xsd文件,通过我们的定义按照这个文件进行约束供我们进行实现。
里面配置了一个dubbonamespacehandler 类,该类实现于spring-framework下面的beans模块的 namespacehandlersupport 抽象类,该抽象类又实现于 namespacehandler 接口,这个接口里面有一个init()方法,dubbonamespacehandler 通过init()方法会把我们在xml里面配置的标签通过 dubbobeandefinitionparser 给解析出来


在init()方法里面可以看到registerbeandefinitionparser()方法的第一个参数就是我们在xml里面配置的<dubbo:xxx> 标签,然后把这些标签解析为 dubbobeandefinitionparser 构造方法里面出入的class类型,比如第一个application,会把配置文件里面的<dubbo:application>转为 applicationconfig ,最后转为 beandefinition 能够让spring管理dubbo
public class dubbonamespacehandler extends namespacehandlersupport implements configurablesourcebeanmetadataelement {
static {
version.checkduplicate(dubbonamespacehandler.class);
}
@override
public void init() {
registerbeandefinitionparser("application", new dubbobeandefinitionparser(applicationconfig.class, true));
registerbeandefinitionparser("module", new dubbobeandefinitionparser(moduleconfig.class, true));
registerbeandefinitionparser("registry", new dubbobeandefinitionparser(registryconfig.class, true));
registerbeandefinitionparser("config-center", new dubbobeandefinitionparser(configcenterbean.class, true));
registerbeandefinitionparser("metadata-report", new dubbobeandefinitionparser(metadatareportconfig.class, true));
registerbeandefinitionparser("monitor", new dubbobeandefinitionparser(monitorconfig.class, true));
registerbeandefinitionparser("metrics", new dubbobeandefinitionparser(metricsconfig.class, true));
registerbeandefinitionparser("ssl", new dubbobeandefinitionparser(sslconfig.class, true));
registerbeandefinitionparser("provider", new dubbobeandefinitionparser(providerconfig.class, true));
registerbeandefinitionparser("consumer", new dubbobeandefinitionparser(consumerconfig.class, true));
registerbeandefinitionparser("protocol", new dubbobeandefinitionparser(protocolconfig.class, true));
registerbeandefinitionparser("service", new dubbobeandefinitionparser(servicebean.class, true));
registerbeandefinitionparser("reference", new dubbobeandefinitionparser(referencebean.class, false));
registerbeandefinitionparser("annotation", new annotationbeandefinitionparser());
}
......
}
dubbobeandefinitionparser 实现于 spring-framework的beans模块的 beandefinitionparser 接口,该接口有一个
beandefinition parse(element element, parsercontext parsercontext)方法,该方法主要是解析我们xml里面的标签。
在 dubbobeandefinitionparser 接口的主要实现代码
@override
public beandefinition parse(element element, parsercontext parsercontext) {
return parse(element, parsercontext, beanclass, required);
}
@suppresswarnings("unchecked")
private static rootbeandefinition parse(element element, parsercontext parsercontext, class<?> beanclass, boolean required) {
rootbeandefinition beandefinition = new rootbeandefinition();
beandefinition.setbeanclass(beanclass);
beandefinition.setlazyinit(false);
string id = resolveattribute(element, "id", parsercontext);
if (stringutils.isempty(id) && required) {
string generatedbeanname = resolveattribute(element, "name", parsercontext);
if (stringutils.isempty(generatedbeanname)) {
if (protocolconfig.class.equals(beanclass)) {
generatedbeanname = "dubbo";
} else {
generatedbeanname = resolveattribute(element, "interface", parsercontext);
}
}
if (stringutils.isempty(generatedbeanname)) {
generatedbeanname = beanclass.getname();
}
id = generatedbeanname;
int counter = 2;
while (parsercontext.getregistry().containsbeandefinition(id)) {
id = generatedbeanname + (counter++);
}
}
if (stringutils.isnotempty(id)) {
if (parsercontext.getregistry().containsbeandefinition(id)) {
throw new illegalstateexception("duplicate spring bean id " + id);
}
parsercontext.getregistry().registerbeandefinition(id, beandefinition);
beandefinition.getpropertyvalues().addpropertyvalue("id", id);
}
if (protocolconfig.class.equals(beanclass)) {
for (string name : parsercontext.getregistry().getbeandefinitionnames()) {
beandefinition definition = parsercontext.getregistry().getbeandefinition(name);
propertyvalue property = definition.getpropertyvalues().getpropertyvalue("protocol");
if (property != null) {
object value = property.getvalue();
if (value instanceof protocolconfig && id.equals(((protocolconfig) value).getname())) {
definition.getpropertyvalues().addpropertyvalue("protocol", new runtimebeanreference(id));
}
}
}
} else if (servicebean.class.equals(beanclass)) {
string classname = resolveattribute(element, "class", parsercontext);
if (stringutils.isnotempty(classname)) {
rootbeandefinition classdefinition = new rootbeandefinition();
classdefinition.setbeanclass(reflectutils.forname(classname));
classdefinition.setlazyinit(false);
parseproperties(element.getchildnodes(), classdefinition, parsercontext);
beandefinition.getpropertyvalues().addpropertyvalue("ref", new beandefinitionholder(classdefinition, id + "impl"));
}
} else if (providerconfig.class.equals(beanclass)) {
parsenested(element, parsercontext, servicebean.class, true, "service", "provider", id, beandefinition);
} else if (consumerconfig.class.equals(beanclass)) {
parsenested(element, parsercontext, referencebean.class, false, "reference", "consumer", id, beandefinition);
}
set<string> props = new hashset<>();
managedmap parameters = null;
for (method setter : beanclass.getmethods()) {
string name = setter.getname();
if (name.length() > 3 && name.startswith("set")
&& modifier.ispublic(setter.getmodifiers())
&& setter.getparametertypes().length == 1) {
class<?> type = setter.getparametertypes()[0];
string beanproperty = name.substring(3, 4).tolowercase() + name.substring(4);
string property = stringutils.cameltosplitname(beanproperty, "-");
props.add(property);
// check the setter/getter whether match
method getter = null;
try {
getter = beanclass.getmethod("get" + name.substring(3), new class<?>[0]);
} catch (nosuchmethodexception e) {
try {
getter = beanclass.getmethod("is" + name.substring(3), new class<?>[0]);
} catch (nosuchmethodexception e2) {
// ignore, there is no need any log here since some class implement the interface: environmentaware,
// applicationaware, etc. they only have setter method, otherwise will cause the error log during application start up.
}
}
if (getter == null
|| !modifier.ispublic(getter.getmodifiers())
|| !type.equals(getter.getreturntype())) {
continue;
}
if ("parameters".equals(property)) {
parameters = parseparameters(element.getchildnodes(), beandefinition, parsercontext);
} else if ("methods".equals(property)) {
parsemethods(id, element.getchildnodes(), beandefinition, parsercontext);
} else if ("arguments".equals(property)) {
parsearguments(id, element.getchildnodes(), beandefinition, parsercontext);
} else {
string value = resolveattribute(element, property, parsercontext);
if (value != null) {
value = value.trim();
if (value.length() > 0) {
if ("registry".equals(property) && registryconfig.no_available.equalsignorecase(value)) {
registryconfig registryconfig = new registryconfig();
registryconfig.setaddress(registryconfig.no_available);
beandefinition.getpropertyvalues().addpropertyvalue(beanproperty, registryconfig);
} else if ("provider".equals(property) || "registry".equals(property) || ("protocol".equals(property) && abstractserviceconfig.class.isassignablefrom(beanclass))) {
/**
* for 'provider' 'protocol' 'registry', keep literal value (should be id/name) and set the value to 'registryids' 'providerids' protocolids'
* the following process should make sure each id refers to the corresponding instance, here's how to find the instance for different use cases:
* 1. spring, check existing bean by id, see{@link servicebean#afterpropertiesset()}; then try to use id to find configs defined in remote config center
* 2. api, directly use id to find configs defined in remote config center; if all config instances are defined locally, please use {@link serviceconfig#setregistries(list)}
*/
beandefinition.getpropertyvalues().addpropertyvalue(beanproperty + "ids", value);
} else {
object reference;
if (isprimitive(type)) {
if ("async".equals(property) && "false".equals(value)
|| "timeout".equals(property) && "0".equals(value)
|| "delay".equals(property) && "0".equals(value)
|| "version".equals(property) && "0.0.0".equals(value)
|| "stat".equals(property) && "-1".equals(value)
|| "reliable".equals(property) && "false".equals(value)) {
// backward compatibility for the default value in old version's xsd
value = null;
}
reference = value;
} else if (onreturn.equals(property) || onthrow.equals(property) || oninvoke.equals(property)) {
int index = value.lastindexof(".");
string ref = value.substring(0, index);
string method = value.substring(index + 1);
reference = new runtimebeanreference(ref);
beandefinition.getpropertyvalues().addpropertyvalue(property + method, method);
} else {
if ("ref".equals(property) && parsercontext.getregistry().containsbeandefinition(value)) {
beandefinition refbean = parsercontext.getregistry().getbeandefinition(value);
if (!refbean.issingleton()) {
throw new illegalstateexception("the exported service ref " + value + " must be singleton! please set the " + value + " bean scope to singleton, eg: <bean id=\"" + value + "\" scope=\"singleton\" ...>");
}
}
reference = new runtimebeanreference(value);
}
beandefinition.getpropertyvalues().addpropertyvalue(beanproperty, reference);
}
}
}
}
}
}
namednodemap attributes = element.getattributes();
int len = attributes.getlength();
for (int i = 0; i < len; i++) {
node node = attributes.item(i);
string name = node.getlocalname();
if (!props.contains(name)) {
if (parameters == null) {
parameters = new managedmap();
}
string value = node.getnodevalue();
parameters.put(name, new typedstringvalue(value, string.class));
}
}
if (parameters != null) {
beandefinition.getpropertyvalues().addpropertyvalue("parameters", parameters);
}
return beandefinition;
}
这里面的判断有兴趣的可以自行debug研究,再来看看annotation是怎么实现的。
第二种:annotation
通过 dubbocomponentscan 注解,该注解spring的 @import 注解导入了 dubbocomponentscanregistrar.class ,该注解实现于 importbeandefinitionregistrar 接口,springboot里面很多地方都是通过该接口进行扩展的,这里不做深入研究。
服务发布
再来看下上面提到过的init()方法,我们知道服务发布是通过 @service 注解(该注解在2.7.7版本已被废弃,通过dubboservice 取代)实现的,会把我们 @service 注解配置的类或者<dubbo:service>标签配置的类转为servicebean

可以看到该类实现了五个接口
initializingbean(执行初始化方法)disposablebean(执行销毁方法)applicationcontextaware(获取applicationcontext上下文)applicationlistener<contextrefreshedevent>(事件监听,2.7.7已去掉)beannameaware(获取当前的beanname)applicationeventpublisheraware(发布事件)
这五个接口都是spring的

在spring上下文加载的时候触发 applicationlistener 接口的onapplicationevent()方法监听,首先判断该服务是否发布过了,没有发布的话就调用export()方法。
@override
public void onapplicationevent(contextrefreshedevent event) {
if (!isexported() && !isunexported()) {
if (logger.isinfoenabled()) {
logger.info("the service ready on spring started. service: " + getinterface());
}
export();
}
}
@override
public void export() {
// 这里调用父类
super.export();
// publish servicebeanexportedevent
publishexportevent();
}
再来看父类 serviceconfig 的export()方法,
public synchronized void export() {
// 检查或更新配置,比如完善默认配置:端口号、协议
checkandupdatesubconfigs();
// 是否要发布服务,可通过@service注解的export属性设置为fasle,就是不把该接口暴露出去,默认值为true
if (!shouldexport()) {
return;
}
// 是否延时发布,可通过@service注解的delay属性设置,这里提供延时加载的功能可能是考虑到spring容器还没初始化就加载该类会出错
if (shoulddelay()) {
delayexportexecutor.schedule(this::doexport, getdelay(), timeunit.milliseconds);
} else {
doexport();
}
}
在通过doexport()方法调用doexporturls(),首先加载我们配置的注册中心地址,通过url来驱动
registry://ip:2181/org.apache.dubbo.registry.registryservice/…
private void doexporturls() {
list<url> registryurls = loadregistries(true);
for (protocolconfig protocolconfig : protocols) {
string pathkey = url.buildkey(getcontextpath(protocolconfig).map(p -> p + "/" + path).orelse(path), group, version);
providermodel providermodel = new providermodel(pathkey, ref, interfaceclass);
applicationmodel.initprovidermodel(pathkey, providermodel);
doexporturlsfor1protocol(protocolconfig, registryurls);
}
}
registries就是我们配置的注册中心的地址

然后会拼接上我们配置的dubbo应用名称、协议名、端口号等一系列参数。

然后判断如果是服务提供者并且当前地址里面的registry参数不为空的话就把当前地址添加到registrylist注册中心集合里面。

遍历完再回到doexporturls()方法里面,接着去遍历我们配置的协议规则,目前我们只配置dubbo这一种协议,然后去构建我们的服务接口

进入到doexporturlsfor1protocol()方法,先把之前构建好的参数封装到map里面,这个map里面最多可以有三十多个参数,因为有些没有配置,这里只显示了20个

然后通过把map变为一个url,这个url就非常熟悉了,就是我们最终存放在zk的地址
dubbo://192.168.124.7:20880/com.tqz.dubbo.api.isayhelloservice?anyhost=true&application=springboot-dubbo-provider&bean.name=servicebean:com.tqz.dubbo.api.isayhelloservice&bind.ip=192.168.124.7&bind.port=20880&cluster=failsafe&deprecated=false&dubbo=2.0.2&dynamic=true&generic=false&interface=com.tqz.dubbo.api.isayhelloservice&methods=sayhello&pid=73828&qos-enable=false®ister=true&release=2.7.2&side=provider&timeout=50000×tamp=1616175325378

这时候url已经组装好了,接下来就要进行服务的发布了。
还是在doexporturlsfor1protocol()方法中,首先获取了一个服务发布的范围,这里一共分为两种,一种是在同一个jvm里面调用,没必要走远程通信会通过injvm://ip:port方式进行调用,另外一种是remote,也就是上面我们组装的url进行远程调用。
如果我们配置了注册中心的地址,默认两种都会进行发布。这里我们假设是远程调用,遍历registryurls,也就是我们之前的registry://ip:2181/org.apache.dubbo.registry.registryservice/…
string scope = url.getparameter(scope_key);
// don't export when none is configured
if (!scope_none.equalsignorecase(scope)) {
// 如果是本地发布
if (!scope_remote.equalsignorecase(scope)) {
exportlocal(url);
}
// 如果是远程调用
if (!scope_local.equalsignorecase(scope)) {
if (!isonlyinjvm() && logger.isinfoenabled()) {
logger.info("export dubbo service " + interfaceclass.getname() + " to url " + url);
}
if (collectionutils.isnotempty(registryurls)) {
for (url registryurl : registryurls) {
// registryurl等于:registry://ip:2181/org.apache.dubbo.registry.registryservice/...
//if protocol is only injvm ,not register
if (local_protocol.equalsignorecase(url.getprotocol())) {
continue;
}
url = url.addparameterifabsent(dynamic_key, registryurl.getparameter(dynamic_key));
url monitorurl = loadmonitor(registryurl);
if (monitorurl != null) {
url = url.addparameterandencoded(monitor_key, monitorurl.tofullstring());
}
......
// for providers, this is used to enable custom proxy to generate invoker
string proxy = url.getparameter(proxy_key);
if (stringutils.isnotempty(proxy)) {
registryurl = registryurl.addparameter(proxy_key, proxy);
}
invoker<?> invoker = proxyfactory.getinvoker(ref, (class) interfaceclass, registryurl.addparameterandencoded(export_key, url.tofullstring()));
delegateprovidermetadatainvoker wrapperinvoker = new delegateprovidermetadatainvoker(invoker, this);
exporter<?> exporter = protocol.export(wrapperinvoker);
exporters.add(exporter);
}
} else {
invoker<?> invoker = proxyfactory.getinvoker(ref, (class) interfaceclass, url);
delegateprovidermetadatainvoker wrapperinvoker = new delegateprovidermetadatainvoker(invoker, this);
exporter<?> exporter = protocol.export(wrapperinvoker);
exporters.add(exporter);
}
/**
* @since 2.7.0
* servicedata store
*/
metadatareportservice metadatareportservice = null;
if ((metadatareportservice = getmetadatareportservice()) != null) {
metadatareportservice.publishprovider(url);
}
}
}
this.urls.add(url);
上面的先不用关心,直接看
exporter<?> exporter = protocol.export(wrapperinvoker)这行,这时候的protocol是成员属性
private static final protocol protocol = extensionloader.getextensionloader(protocol.class).getadaptiveextension();
自适应扩展点,之前分析过 protocol 接口的export()方法和refer()方法都有 @adaptive 注解,会根据代理生成字节码protocol$adaptive,回过头来再看下这个类。
在我们上面的registryurls是registry://ip:2181/org.apache.dubbo.registry.registryservice/…,所以生成字节码类里面的export()方法里面的extname一定是registry
public class protocol$adaptive implements org.apache.dubbo.rpc.protocol {
public void destroy() {
throw new unsupportedoperationexception("the method public abstract void org.apache.dubbo.rpc.protocol.destroy() of interface org.apache.dubbo.rpc.protocol is not adaptive method!");
}
public int getdefaultport() {
throw new unsupportedoperationexception("the method public abstract int org.apache.dubbo.rpc.protocol.getdefaultport() of interface org.apache.dubbo.rpc.protocol is not adaptive method!");
}
public org.apache.dubbo.rpc.exporter export(org.apache.dubbo.rpc.invoker arg0) throws org.apache.dubbo.rpc.rpcexception {
if (arg0 == null) throw new illegalargumentexception("org.apache.dubbo.rpc.invoker argument == null");
if (arg0.geturl() == null)
throw new illegalargumentexception("org.apache.dubbo.rpc.invoker argument geturl() == null");
org.apache.dubbo.common.url url = arg0.geturl();
string extname = (url.getprotocol() == null ? "dubbo" : url.getprotocol());
if (extname == null)
throw new illegalstateexception("failed to get extension (org.apache.dubbo.rpc.protocol) name from url (" + url.tostring() + ") use keys([protocol])");
// extname等于:registry
org.apache.dubbo.rpc.protocol extension = (org.apache.dubbo.rpc.protocol) extensionloader.getextensionloader(org.apache.dubbo.rpc.protocol.class).getextension(extname);
return extension.export(arg0);
}
public org.apache.dubbo.rpc.invoker refer(java.lang.class arg0, org.apache.dubbo.common.url arg1) throws org.apache.dubbo.rpc.rpcexception {
if (arg1 == null) throw new illegalargumentexception("url == null");
org.apache.dubbo.common.url url = arg1;
string extname = (url.getprotocol() == null ? "dubbo" : url.getprotocol());
if (extname == null)
throw new illegalstateexception("failed to get extension (org.apache.dubbo.rpc.protocol) name from url (" + url.tostring() + ") use keys([protocol])");
org.apache.dubbo.rpc.protocol extension = (org.apache.dubbo.rpc.protocol) extensionloader.getextensionloader(org.apache.dubbo.rpc.protocol.class).getextension(extname);
return extension.refer(arg0, arg1);
}
}
在去dubbo的meta-inf里面看下org.apache.dubbo.rpc.protocol文件,在里面可以找到一个key为registry的,所以exprot()方法里面最后的extension应该等于 registryprotocol

进入到 registryprotocol 的export()方法,首先获取了注册中心的地址和服务提供者的地址
@override
public <t> exporter<t> export(final invoker<t> origininvoker) throws rpcexception {
url registryurl = getregistryurl(origininvoker);
// url to export locally
url providerurl = getproviderurl(origininvoker);
// subscribe the override data
// fixme when the provider subscribes, it will affect the scene : a certain jvm exposes the service and call
// the same service. because the subscribed is cached key with the name of the service, it causes the
// subscription information to cover.
// 服务更改进行重新,比如在控制台进行了更改
final url overridesubscribeurl = getsubscribedoverrideurl(providerurl);
final overridelistener overridesubscribelistener = new overridelistener(overridesubscribeurl, origininvoker);
overridelisteners.put(overridesubscribeurl, overridesubscribelistener);
providerurl = overrideurlwithconfig(providerurl, overridesubscribelistener);
//export invoker
// 启动一个netty服务
final exporterchangeablewrapper<t> exporter = dolocalexport(origininvoker, providerurl);
// url to registry
final registry registry = getregistry(origininvoker);
final url registeredproviderurl = getregisteredproviderurl(providerurl, registryurl);
providerinvokerwrapper<t> providerinvokerwrapper = providerconsumerregtable.registerprovider(origininvoker, registryurl, registeredproviderurl);
//to judge if we need to delay publish
boolean register = registeredproviderurl.getparameter("register", true);
if (register) {
register(registryurl, registeredproviderurl);
providerinvokerwrapper.setreg(true);
}
// deprecated! subscribe to override rules in 2.6.x or before.
registry.subscribe(overridesubscribeurl, overridesubscribelistener);
exporter.setregisterurl(registeredproviderurl);
exporter.setsubscribeurl(overridesubscribeurl);
//ensure that a new exporter instance is returned every time export
return new destroyableexporter<>(exporter);
}
接下来我们看看是怎么启动的,在dolocalexport()方法里,bounds就是一个 concurrenthashmap,这里使用了jdk1.8特性,判断当前的key是否存在,就如不存在就put到该map中,然后通过protocol.export()方法暴露出去。这里的protocol会经过层层的包装,
- qosprotocolwrapper(protocollistenerwrapper(protocolfilterwrapper(dubboprotocol)))
最终也就是我们配置的dubbo协议
@suppresswarnings("unchecked")
private <t> exporterchangeablewrapper<t> dolocalexport(final invoker<t> origininvoker, url providerurl) {
string key = getcachekey(origininvoker);
return (exporterchangeablewrapper<t>) bounds.computeifabsent(key, s -> {
invoker<?> invokerdelegate = new invokerdelegate<>(origininvoker, providerurl);
return new exporterchangeablewrapper<>((exporter<t>) protocol.export(invokerdelegate), origininvoker);
});
}
该方法前面的都不用关心,直接看最后两行,首先开启一个服务,暴露我们配置的20880协议,然后优化序列化
@override
public <t> exporter<t> export(invoker<t> invoker) throws rpcexception {
url url = invoker.geturl();
// export service.
string key = servicekey(url);
dubboexporter<t> exporter = new dubboexporter<t>(invoker, key, exportermap);
exportermap.put(key, exporter);
//export an stub service for dispatching event
boolean isstubsupportevent = url.getparameter(stub_event_key, default_stub_event);
boolean iscallbackservice = url.getparameter(is_callback_service, false);
if (isstubsupportevent && !iscallbackservice) {
string stubservicemethods = url.getparameter(stub_event_methods_key);
if (stubservicemethods == null || stubservicemethods.length() == 0) {
if (logger.iswarnenabled()) {
logger.warn(new illegalstateexception("consumer [" + url.getparameter(interface_key) +
"], has set stubproxy support event ,but no stub methods founded."));
}
} else {
stubservicemethodsmap.put(url.getservicekey(), stubservicemethods);
}
}
openserver(url);
optimizeserialization(url);
return exporter;
}
在openserver里面会调用createserver()方法,默认开启server关闭时发送readonly事件,以及心跳机制事件,default_remoting_server的默认值是netty,还有其他的例如:mina、grizzy。
private exchangeserver createserver(url url) {
url = urlbuilder.from(url)
// send readonly event when server closes, it's enabled by default
.addparameterifabsent(channel_readonlyevent_sent_key, boolean.true.tostring())
// enable heartbeat by default
.addparameterifabsent(heartbeat_key, string.valueof(default_heartbeat))
.addparameter(codec_key, dubbocodec.name)
.build();
string str = url.getparameter(server_key, default_remoting_server);
if (str != null && str.length() > 0 && !extensionloader.getextensionloader(transporter.class).hasextension(str)) {
throw new rpcexception("unsupported server type: " + str + ", url: " + url);
}
exchangeserver server;
try {
server = exchangers.bind(url, requesthandler);
} catch (remotingexception e) {
throw new rpcexception("fail to start server(url: " + url + ") " + e.getmessage(), e);
}
str = url.getparameter(client_key);
if (str != null && str.length() > 0) {
set<string> supportedtypes = extensionloader.getextensionloader(transporter.class).getsupportedextensions();
if (!supportedtypes.contains(str)) {
throw new rpcexception("unsupported client type: " + str);
}
}
return server;
}
然后通过
exchangers.bind(url, requesthandler) -> org.apache.dubbo.remoting.exchange.support.header.headerexchanger#bind -> org.apache.dubbo.remoting.transporters#bind(org.apache.dubbo.common.url, org.apache.dubbo.remoting.channelhandler...) -> org.apache.dubbo.remoting.transport.netty4.nettytransporter#bind
进入到该方法,会初始化 nettyserver
@override
public server bind(url url, channelhandler listener) throws remotingexception {
return new nettyserver(url, listener);
}
public nettyserver(url url, channelhandler handler) throws remotingexception {
super(url, channelhandlers.wrap(handler, executorutil.setthreadname(url, server_thread_pool_name)));
}
调用父类的构造方法,然后调用doopen()方法,该方法为抽象方法,由刚刚的实现类 nettyserver 去实现
public abstractserver(url url, channelhandler handler) throws remotingexception {
......
try {
doopen();
if (logger.isinfoenabled()) {
logger.info("start " + getclass().getsimplename() + " bind " + getbindaddress() + ", export " + getlocaladdress());
}
} catch (throwable t) {
throw new remotingexception(url.toinetsocketaddress(), null, "failed to bind " + getclass().getsimplename()
+ " on " + getlocaladdress() + ", cause: " + t.getmessage(), t);
}
......
}
由刚刚的实现类 nettyserver 去实现,到这里就把服务给暴露出去了。
@override
protected void doopen() throws throwable {
bootstrap = new serverbootstrap();
bossgroup = new nioeventloopgroup(1, new defaultthreadfactory("nettyserverboss", true));
workergroup = new nioeventloopgroup(geturl().getpositiveparameter(io_threads_key, constants.default_io_threads),
new defaultthreadfactory("nettyserverworker", true));
final nettyserverhandler nettyserverhandler = new nettyserverhandler(geturl(), this);
channels = nettyserverhandler.getchannels();
bootstrap.group(bossgroup, workergroup)
.channel(nioserversocketchannel.class)
.childoption(channeloption.tcp_nodelay, boolean.true)
.childoption(channeloption.so_reuseaddr, boolean.true)
.childoption(channeloption.allocator, pooledbytebufallocator.default)
.childhandler(new channelinitializer<niosocketchannel>() {
@override
protected void initchannel(niosocketchannel ch) throws exception {
// fixme: should we use gettimeout()?
int idletimeout = urlutils.getidletimeout(geturl());
nettycodecadapter adapter = new nettycodecadapter(getcodec(), geturl(), nettyserver.this);
ch.pipeline()//.addlast("logging",new logginghandler(loglevel.info))//for debug
.addlast("decoder", adapter.getdecoder())
.addlast("encoder", adapter.getencoder())
.addlast("server-idle-handler", new idlestatehandler(0, 0, idletimeout, milliseconds))
.addlast("handler", nettyserverhandler);
}
});
// bind
channelfuture channelfuture = bootstrap.bind(getbindaddress());
channelfuture.syncuninterruptibly();
channel = channelfuture.channel();
}
回过头来看看服务发布做了哪些事情
- 基于spirng进行解析配置文件存储到config
- 各种判断逻辑,保证配置信息安全性
- 组装url:registry:// -> zookeeper:// -> dubbo:// -> injvm
- 构建一个invoker(代理)
- registryprotocol.export()
- 各种包装(qos/filter/lisenter)
- dubboprotocol.export()发布服务
- 启动一个nettyserver#doopen()
总结
以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。
发表评论