当前位置: 代码网 > it编程>编程语言>Java > dubbo服务发布过程是怎样的?从配置到netty启动全解析

dubbo服务发布过程是怎样的?从配置到netty启动全解析

2026年08月28日 Java 我要评论
服务注册与发布,书接上回,接着分析dubbo的服务注册与发布。首先了解一下dubbo怎么解析自定义的标签的。第一种:xml在 dubbo-config 模块的 dubbo-config-spring

服务注册与发布,书接上回,接着分析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&register=true&release=2.7.2&side=provider&timeout=50000&timestamp=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()

总结

以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。

(0)

相关文章:

版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。

发表评论

验证码:
Copyright © 2017-2026  代码网 保留所有权利. 粤ICP备2024248653号
站长QQ:2386932994 | 联系邮箱:2386932994@qq.com