Hero Image
记一次搜索迁移故障复盘

临近双十一,集团已经开始陆陆续续的封网,没有大的需求,都是些零零散散的需求在推进,整个人比较轻松。正好自己负责的系统有些历史架构问题可以在这个时候开始做起来了。 系统冗余了一份商货关系数据(商品和货品的绑定关系),当然这个份数据是个大宽表,除了绑定关系之外还冗余了商品 SKU 的一些信息供展示和搜索。 DB 的数据量级在1000W+(单表),但是最近业务数据增长的有点快,数据量级很快超过了2000W+,DB搜索的性能已经比较低了,很多走索引的查询也要100ms以上。长期来看肯定要迁移数据,切换到分表。但是业务数据还在不停的增长,因此我的方案是先切换搜索引擎,解决DB查询慢的问题,然后后面找机会再进行分表的切换。 背景 其实在阿里内部使用搜索有两种备选方案。一种是使用 OpenSearch。好处是,运维简单、开发起来也很方便。但是功能有些限制,比如:最大分页数不能超过 5000。另一种是自建 HA3(openSearch 本质上是在 HA3 的基础上做了一个通用的解决方案,做了很多限制),这个需要人工去运维 HA3 服务,多了些理解、运维的成本,但好处是自由度很高。由于我这个搜索并不需要很多的定制化搜索的能力,因此我用的是 OpenSearch 这套方案。 所以接下来的事情就简单多了,直接把 DB 的数据导入到 OpenSearch,然后创建好索引,那就已经完成了搜索引擎的搭建,接下来就是使用 SDK 了,具体细节就不讲了,应该只花了一个星期不到就完成了搜索的切换。(只需要做两件事一是把请求入参转换成 OpeanSearch 的语法,二是把搜索的结果转换成原来的数据对象) 当时没有测试资源,我自己测了一个两天,把原来接口的入参都一个个的试了一遍,确保字段都能进行有效的搜索之后,就上线了。当然,为了避免出现问题之后能即使止血,我加了个开关。发布的时候也灰度了5个小时,当然这些并不能避免问题,只是希望出现问题之后能即使止血,及时发现。 问题 按理说我加了止血方案,同时又灰度了5小时,应该没大问题。但是很不幸,我唯独评估漏了一个项,线上 OpenSearch 所需要的资源。这就导致我的 OpenSearch 搜索服务被限流了。(OpenSearch 有 LCU 的限制,如果 LCU 超了申请时候的配额会引发查询限流)。 虽然上游重试下就可以解决这个问题,但是因为临近双十一,大家都希望系统稳定点,所以直接跟老板反馈了。。 复盘 从表面上看这个是一个很小的问题,查询限流是很常见的事情。但是回顾自己整个迁移的方案,我发现自己漏了两项。一是没有梳理上游依赖这个接口的场景,因为这个接口之前直接查的 DB,之前不会限流迁移之后会限流。所以,如果上游根本没有重试的机制,那会导致上游业务处理失败。这次还好可以让上游手动重试,万一上游还没法手动重试,那只能去订正数据了。二是根本没有做压测和资源评估这件事,当然这跟我不熟悉 OpenSearch 有关,但是压测这事儿是很大的疏漏,因为我的自测只是保障了接口功能上的正常,却忘掉了真实场景下会有突发流量的问题。 总结 所以我总结总结了下,做迁移还是有些原则性的东西需要去做,否则方案就不完整。 上下游依赖梳理。需要弄清楚影响范围,这样可以在上线前做好相应的预案,上线时出了问题也能及时找到相应的人处理。 资源评估。系统依赖的中间件需要哪些资源,比如:CPU、磁盘这些东西,最好有个量化的指标,这样也便于去申请。 保证功能完整。迁移前后的功能一定是要保证一致的,这只最重要得一个大前提。 压测模拟。因为在小流量和大流量的场景下需要保障的东西完全不同,最常见的就是接口限流,RT变高等问题这些情况都是在小流量场景下没法复现的。 流量可监控。整个发布过程中我都是被动的等待问题出现,在迁移的过程中一定要有一些接口级别的异常监控,能帮我们主动发现问题。

    Hero Image
    21天Python学习计划

    前言 从毕业以来我把自己的主要精力都花在学习中间件的设计实现、分布式算法的原理及使用场景、架构设计理论(方法论)。我认为自己java水平并没有很高,如果自己工作中用到的语言都学不好就跑去学其他语言相关的东西,那就是”捡了芝麻,丢了西瓜“。当然,更多的是没有使用场景,在工作中我基本上不可能跨语言编程来完成我的工作。可能是我年少时太无知,随着我编程经验的增长,我开始否定自己之前的看法。我发现一些曾经的小语种也变得流行起来,而且在一些细分领域上他们似乎更加有优势。 举个例子来说,我发现在日常工作中我经常会和excel打交道,这个东西很烦,各种格式问题。我工作中最长见一个场景是运营给我excel让我帮忙刷数据。这种情况下使用java来完成说实话真的很烦,你需要去解析excel然后再去刷。我们知道excel烦就烦在我们需要对其格式进行调试,有时候一个单元格里面放的到底是字符串还是数字这个需要我们去调试。那这种情况下我们可以使用一些工具把excel里面的数据转化一下,让工具帮我们把数据处理成一个更容易解析的格式,那我们就可以直接得到结构化的数据,能省去很多debug的时间 那就直入主题,讲讲今天的主角Python。之所以选择针对Python进行一个21天的学习计划,主要目的是为了能学会Python的数据处理方式。 这套学习方法肯定不太会注重机器学习这方面,因为我不是做这方面的人才。这套学习计划主要针对Python基础语法、数据处理类库(Pandas、Numpy)的学习。如果有兴趣的朋友可以跟着学学看看,发现有些不合理的地方可以联系我改正。 第一周 学习内容 词法、句法、oop、基础数据类型的常见操作(尤其注意要熟练字符串的拼接、替换)。文件io(csv、txt、json、excel),原生os层面的io以及三方框架的io。 词法: 变量的定义、变量的种类、变量的作用范围 句法: 方法的定义、条件判断语句、选择语句、循环语句 oop: 类的定义、类的实例化、类的创建与销毁(只需要知道语句层面,不需要深入到虚拟机层面) 文件io: 这里只是学习API层面,如何使用好os、pandas库来读取文件。 相关资料 像计算机科学家那样思考 Python中文版第二版.pdf (这本书网上很多地方都有,自行下载就好) python读取文件的几种方式 Pandas库的学习 检验方式 字符串拼接、替换、查找判断。 利用循环打印杨辉三角。 利用递归解决斐波拉契数列问题。 文件读写,对数据进行排序、分组、去重等操作。 第二周 学习内容 数据分析工具的使用(numpy、pandas)、socket使用、http使用、正则的使用 检验方式 使用socket库搭建聊天服务器(需要用oop来抽象,并对服务端的socket通讯做到抽象隔离,让他能用任何方式来替换其实现。比如:目前基于系统socket库来实现,我们要保证换成其他库来实现仍然能保持最小改动量,只需要实现一个类就行。类似java的多态机制) http GET POST 等的请求的使用 爬取一个页面,获取页面上的所有图片、链接,并分别按照字符串倒序排序 使用数据分析工具分析日志 (统计错误次数、找到错误次数最多的错误信息) 相关资料 python 实现TCP socket通信和 HTTP服务器、服务器和客户端通信python实例 python实现http请求 轻松了解python正则表达式 (超详细,附举例) Python之numpy详细教程 如何在Pandas中实现类似于SQL查询的数据操作? Python pandas用法最全整理 第三周 学习内容 基于scitoll-learn来搭建机器学习模型 检验方式 跑通学习案例 相关资料 机器学习最佳Python库:Scikit-learn入门指南 关于机器学习中的Scikit-Learn,你不知道的10个实用功能 -在找到第一份数据科学工作前,我们可以通过哪些方法来获取工作经验? 写在最后 这里强调下,这篇文章主要侧重于把Python当做一个工具来使用。很多Python重点应用的地方我省略掉了,原因就是我不是专门从事这个语言的开发者,我更多的诉求只是利用它的便利性来当做一个小工具。如果真的想学习机器学习等方面的知识,本篇文章不提供任何参考价值。

      Hero Image
      Nacos-Config模块源码讲解

      前言 我们知道Nacos其实是由两个重要的模块组成,一是 Naming 模块,另一个就是今天要讲的 Config 模块。 版本说明 Nacos:2.1.1 jdk:1.8 代码分支:develope 配置中心基本原理 配置中心有三个角色分别是是配置发布者、配置服务端、配置客户端。整个流程由配置发布者发起,先发布一条配置到服务端,然后客户端监听这条配置,当配置有变动时由服务端推送到客用户端。 配置发布 我们从web层着手,可以很快的发现入口为com.alibaba.nacos.config.server.controller.ConfigController#publishConfig,这个方法看上去很长,但其实大体逻辑很简单,示意代码如下所示。 public class ConfigController { /** * Adds or updates non-aggregated data. * <p> * request and response will be used in aspect, see * {@link com.alibaba.nacos.config.server.aspect.CapacityManagementAspect} and * {@link com.alibaba.nacos.config.server.aspect.RequestLogAspect}. * </p> * @throws NacosException NacosException. */ @PostMapping @Secured(action = ActionTypes.WRITE, signType = SignType.CONFIG) public Boolean publishConfig(HttpServletRequest request, HttpServletResponse response, @RequestParam(value = "dataId") String dataId, @RequestParam(value = "group") String group, @RequestParam(value = "tenant", required = false, defaultValue = StringUtils.EMPTY) String tenant, @RequestParam(value = "content") String content, @RequestParam(value = "tag", required = false) String tag, @RequestParam(value = "appName", required = false) String appName, @RequestParam(value = "src_user", required = false) String srcUser, @RequestParam(value = "config_tags", required = false) String configTags, @RequestParam(value = "desc", required = false) String desc, @RequestParam(value = "use", required = false) String use, @RequestParam(value = "effect", required = false) String effect, @RequestParam(value = "type", required = false) String type, @RequestParam(value = "schema", required = false) String schema) throws NacosException { final String srcIp = RequestUtil.getRemoteIp(request); final String requestIpApp = RequestUtil.getAppName(request); if (StringUtils.isBlank(srcUser)) { srcUser = RequestUtil.getSrcUserName(request); } //check type if (!ConfigType.isValidType(type)) { type = ConfigType.getDefaultType().getType(); } // encrypted Pair<String, String> pair = EncryptionHandler.encryptHandler(dataId, content); content = pair.getSecond(); // check tenant ParamUtils.checkTenant(tenant); ParamUtils.checkParam(dataId, group, "datumId", content); ParamUtils.checkParam(tag); Map<String, Object> configAdvanceInfo = new HashMap<>(10); MapUtil.putIfValNoNull(configAdvanceInfo, "config_tags", configTags); MapUtil.putIfValNoNull(configAdvanceInfo, "desc", desc); MapUtil.putIfValNoNull(configAdvanceInfo, "use", use); MapUtil.putIfValNoNull(configAdvanceInfo, "effect", effect); MapUtil.putIfValNoNull(configAdvanceInfo, "type", type); MapUtil.putIfValNoNull(configAdvanceInfo, "schema", schema); ParamUtils.checkParam(configAdvanceInfo); if (AggrWhitelist.isAggrDataId(dataId)) { LOGGER.warn("[aggr-conflict] {} attempt to publish single data, {}, {}", RequestUtil.getRemoteIp(request), dataId, group); throw new NacosException(NacosException.NO_RIGHT, "dataId:" + dataId + " is aggr"); } final Timestamp time = TimeUtils.getCurrentTime(); String betaIps = request.getHeader("betaIps"); ConfigInfo configInfo = new ConfigInfo(dataId, group, tenant, appName, content); configInfo.setType(type); String encryptedDataKey = pair.getFirst(); configInfo.setEncryptedDataKey(encryptedDataKey); if (StringUtils.isBlank(betaIps)) { if (StringUtils.isBlank(tag)) { persistService.insertOrUpdate(srcIp, srcUser, configInfo, time, configAdvanceInfo, false); ConfigChangePublisher.notifyConfigChange( new ConfigDataChangeEvent(false, dataId, group, tenant, time.getTime())); } else { persistService.insertOrUpdateTag(configInfo, tag, srcIp, srcUser, time, false); ConfigChangePublisher.notifyConfigChange( new ConfigDataChangeEvent(false, dataId, group, tenant, tag, time.getTime())); } } else { // beta publish configInfo.setEncryptedDataKey(encryptedDataKey); persistService.insertOrUpdateBeta(configInfo, betaIps, srcIp, srcUser, time, false); ConfigChangePublisher.notifyConfigChange( new ConfigDataChangeEvent(true, dataId, group, tenant, time.getTime())); } ConfigTraceService.logPersistenceEvent(dataId, group, tenant, requestIpApp, time.getTime(), InetUtils.getSelfIP(), ConfigTraceService.PERSISTENCE_EVENT_PUB, content); return true; } } 整个方法就是先组装ConfigInfo,然后委托com.alibaba.nacos.config.server.service.repository.PersistService#insertOrUpdateTag进行处理,这里如果使用的是mysql那么实现类就是com.alibaba.nacos.config.server.service.repository.extrnal.ExternalStoragePersistServiceImpl,这个是一个插入数据库的动作就不细讲了,只是要注意的是这个方法包含新增和更新,当插入失败时会变成更新操作。而操作完数据库之后会发送一个ConfigDataChangeEvent事件,这个事件的处理类是com.alibaba.nacos.config.server.service.notify.AsyncNotifyService,示意代码如下所示。

        Hero Image
        Nacos总览

        前言 说到微服务,注册中心是避不开的话题。在微服务中注册中心的主要功能就是提供服务发现、服务注册这两个功能。市面上的注册中心种类非常多,今天我就讲讲阿里的开源产品 Nacos。 快速构建 源码下载 地址 本地编译 参考官网的Download source code from Github 。 启动 参考官网Start Server 。本地启动的时候一定要加-m standalone Naming+Config=Nacos 从Nacos官网上我们其实也了解到,Nacos分为两个部分。一个是命名服务,另一个是分布式配置,接下来会分开介绍 分布式系统中的命名服务 命名服务的定义看这里 。简单来说在分布式系统中的命名服务的作用主要是给系统中的每个RPC服务分配一个资源符号,命名服务通过这个资源符号能得到这个RPC服务的详细信息以及所有提供这些服务的机器信息。 为什么需要命名服务? 在进行服务调用时有两个信息至关重要,一是我要调用的最新的服务信息(RPC信息、地址信息)。如果不知道服务信息我们是无法调用的,而如果我们的服务信息不是最新的,那又会因为服务信息的不一致导致调用异常(试想下,如果你还在用旧的地址信息,而后端服务这个时候换了机器部署,那就会出现调用异常的情况。同理,后端服务信息发生变更,比如:方法签名发生变更。使用旧的信息仍然会出现调用失败的情况)。一是保证调到正常提供服务的机器上(可能有版本发布或是机器异常的情况,需要剔除掉这些信息)。正式基于这两点,我们需要一个服务能够管理服务的注册信息。 命名服务的主要功能点 信息推送 前面我们讲到了,必须让应用能时刻保持最新的服务信息,而实现这一功能的手段就是消息推送。消息推送的形式有很多种,比如:TCP全双工、UDP广播、回调,这些方式都可以实现。 心跳检测 当然,为了提升服务的可用性,异常服务的剔除也是必不可少的。而实现它的方式就是心跳检测,基于一个固定的时间去发送“ping命令”检查服务的情况。 Nacos中的命名服务 数据结构 上图是Dubbo启动后会注册两个服务,以及它的订阅者。可以看到Nacos的数据结构和Zookeeper的数据结构有很大的不一样。Nacos是基于Key-Value的结构,而Zookeeper是基于一种树形结构。 数据类型 和 Zookeeper 一样 Nacos 也分为两种节点,临时节点和永久节点,临时节点会和创建这个节点的应用建立心跳检测机制用于保证这个机器可达(服务健康检查也是基于这个机制)。而永久节点一经创建就不会被 Nacos 删除,除非你通过删除操作来删除。 分布式系统中的配置服务 配置服务也叫配置中心,用于管理整个系统的配置。提供简单的读写接口、实时更新到系统中各个机器上、自动载入到系统中。 为什么需要配置服务 在开发过程中我们会经常用到配置,比如:应用启动参数配置、业务参数配置等等。对于系统启动配置,我们希望有多个配置,能够针对不同的环境自动匹配(类似Spring的@Profile,可以让配置绑定特定环境)。针对应用业务配置,我们希望配置能被很好的管理,他的生命周期能和应用独立开(在以前配置都是散落在各个机器上,每个机器上都有一份,而且是散落在各个地方,我们没有一个很好的办法去同时修改这些配置文件)。当然,我们现在的业务系统很多业务逻辑,经常会遇到需要禁用某个功能或是有多个功能,这些功能的切换或者功能的禁用都需要基于配置,而传统的配置很难做到统一变更管理,统一的推送。 Nacos中的配置服务 Nacos 基于上面三点做了一个很好的实现,但是在使用时也需要注意配置并不支持高并发的变更(当然也没有高并发的实际场景,所有配置服务都要注意不能高并发更新)。 写在最后 本篇文章只从功能特点上来看 Nacos,主要为了引入分布式系统中的命名服务和配置服务这两个概念,后面会专门从源码实现层面讲。附录里面最后3个链接是 Nacos 和其他同类产品的对比,因为作者们的视角不同,因此需要大家看的时候舍弃部分重复的内容。 附录 Nacos Spring Boot 服务注册发现框架比较 服务注册与发现框架对比 微服务架构 - 服务发现、配置中心对比 - Nacos

          Hero Image
          Nacos-Naming模块源码讲解

          前言 上篇文章列举了Nacos大致的功能列表,但没有深入讲解 Nacos 的实现细节。我们知道其实还有很多同类的产品比如Zookeeper、Eureka、Etcd、Consul等等,这些产品的功能在大体上都和 Nacos 很相似,最主要的区别就在于它的实现。今天我们来深入了解 Nacos 的 Naming 模块的实现。 版本说明 Nacos:2.1.1 jdk:1.8 代码分支:develope 服务注册 根据服务注册url可以定位到服务注册的web层入口为com.alibaba.nacos.naming.controllers.InstanceController#register,示意代码如下所示。 public class InstanceController{ /** * Register new instance. * * @param request http request * @return 'ok' if success * @throws Exception any error during register */ @CanDistro @PostMapping @Secured(action = ActionTypes.WRITE) public String register(HttpServletRequest request) throws Exception { final String namespaceId = WebUtils .optional(request, CommonParams.NAMESPACE_ID, Constants.DEFAULT_NAMESPACE_ID); final String serviceName = WebUtils.required(request, CommonParams.SERVICE_NAME); NamingUtils.checkServiceNameFormat(serviceName); final Instance instance = HttpRequestInstanceBuilder.newBuilder() .setDefaultInstanceEphemeral(switchDomain.isDefaultInstanceEphemeral()).setRequest(request).build(); getInstanceOperator().registerInstance(namespaceId, serviceName, instance); NotifyCenter.publishEvent(new RegisterInstanceTraceEvent(System.currentTimeMillis(), "", false, namespaceId, NamingUtils.getGroupName(serviceName), NamingUtils.getServiceName(serviceName), instance.getIp(), instance.getPort())); return "ok"; } } 可以看到整个过程只有两个步骤。一是注册Instance,一是通过NotyfyCenter发布事件。我们拆开看,先看注册Instance的过程。

            Hero Image
            Dubbo3.0-DubboBootstrap模型讲解

            前言 Dubbo3.0有很多新的类(较2.0来说),本文主要就3.0中DubboBootstrap及其相关的类进行讲解,主要为想了解Dubbo3.0启动模型的同学做参考。 版本说明 Dubbo版本:3.0.9 DubboBootstrap的组成部分 DubboBootstrap的关键属性有4个,分别是ApplicationModel、ConfigManager、Environment、ApplicationDeployer。他们功能分别是应用的定义(可以理解为领域对象中的Entity)、配置管理、环境变量管理、应用声明周期管理。 ApplicationModel讲解 讲到ApplicationModel就不得不提到Dubbo的另外两个领域模型FrameworkModel和ModuleModel。 模型之间的依赖关系 这三个模型紧密相连,整个Dubbo共用一个FrameworkModel,一个FrameworkModel中允许有多个ApplicationModel存在,同时在一个ApplicationModel里面又允许有多个ModuleModel同时存在。同时,ApplicationModel又保留FrameworkModel的引用,ModuleModel也保留了ApplicationModel的引用。(也就是说FrameworkModel和ApplicationModel之间保持了双向引用,ApplicationModel和ModuleModel之间也保持了双向引用) 模型共有的特性 这三个模型都集成自ScopeModel,其中我想重点讲的是classLoaders、extensionScope、extensionDirector这三个。 每个Model其实都有自己的classLoader,他们都通过自己的classLoader来加载类(在没有特殊指定的情况下默认是使用AppClassLoader) extensionScope和extensionDirector配合使用来控制SPI扩展点的的实例(ExtensionLoader)。ExtensionLoader现在不是全局单例的,而是在ExtensionScope维度下是单例的(可以理解是按照ExtensionScope维度对SPI扩展点的加载做了隔离,不能跨scope获取SPI实例)。 模型的差别 这三个模型大致功能相同,不同的点在于功能维度不同。FrameworkModel主要提供的是已经暴露过的Provider信息的查询。ApplicationModel和ModuleModel的差别主要体现在Environment、ServiceRepository、ConfigManager、Deployer上。 ApplicationModel中可以查询整个应用的配置、环境变量,以及应用的初始化。 ModuleModel中只能查当前进程里面的相关应用配置、环境变量,同时他的ModuleModel也只负责服务的暴露和引用。 ConfigManager讲解 Config和Environment是有差别的,前者代表的是写在应用内的配置(可以理解为必须要的配置,比如:想要应用启动那一定要配置registry的配置,想要暴露一个服务就必须定义一个ServiceConfig),后者代表外部化配置(可以理解为一个值的获取方式,比如:RegistryConfig的地址是${spring.zookeeper.url},而它的取值方式可以是多种多样的,从spring、环境变量、JVM变量等等。我们可以通过这种方式,让部分属性值做到根据不同环境来设置不同的值) Dubbo将配置做了抽象,比如:ProtocolConfig、RegistryConfig、ReferenceConfig、ServiceConfig、ApplicationConfig等等。每个Config代表的是Dubbo的配置,比如:注册中心的配置在RegistryConfig中、服务暴露的配置在ServiceConfig中等等。当然也没什么说的,这个类主要是提供配置的查询功能,唯一要注意的是ModuleConfigManager和ConfigManager其实就是能加载的配置不一样(在各自构造函数中限制了能加载那些配置)。 Environment讲解 Dubbo的环境变量有几种来源方式,比如:系统环境变量、Spring环境变量。详情可以参考小马哥的文章 ApplicationDeployer讲解 ApplicationDeployer和ModuleDeployer两个类的创建都需要Model(ApplicationDeployer的创建需要ApplicationModel,ModuleDeployer的创建需要ModuleModel),而我们知道Model是有双向引用的,因此这两个Deployer可以通过Model互通。 在DubboBootstrap中我们也能看到这两个类的互调,比如:org.apache.dubbo.config.deploy.DefaultApplicationDeployer#start -> org.apache.dubbo.config.deploy.DefaultApplicationDeployer#doStart -> org.apache.dubbo.config.deploy.DefaultApplicationDeployer#startModules 通过这个调用链我们可以看到ApplicationDeployer调用ModuleDeployer,org.apache.dubbo.config.deploy.DefaultModuleDeployer#start我们从这个方法里面也能看到很多ApplicationDeployer的身影。 两个Deployer的差别 ApplicationDeployer负责应用的启动和初始化,ModuleDeployer负责服务的暴露和引用。 这两个类的差别其实就是功能上的差别,其他的都是由于功能差别导致的(比如内部属性的差异、接口的差异,这些其实都是由他们的职责差异导致的)。

              Hero Image
              Dubbo3.0-服务自省

              前言 Dubbo3.0中有个特别的机制“服务自省”,对于这个机制很多人都很陌生,我用自己的理解来讲讲这个机制。 版本说明 Dubbo版本:3.0.9 定义 所谓的服务自省就是一个元数据同步机制。在Dubbo3.0中,原来面向RPC定义的注册中心数据被拆分为服务元数据和服务地址信息,Dubbo的流量管理能力依赖服务的元数据,因此需要一个机制能很好的去把Provider端的元数据同步到Consumer端。 机制讲解 服务自省是Consumer端获取Provider端服务元数据的一个机制,那么肯定Provider端和Consumer端都有对应的处理。 Provider暴露MetadataService Provider需要提供服务的元数据,所以在Provider端启动时会先暴露所有的业务服务,当所有业务服务都准备好之后(启动并且已经向注册中心注册),暴露MetadataService供Consumer端来查询元数据(此时业务服务也元数据服务都能正常调用,注册中心也有了元数据和应用地址数据)。 Consumer监听MetadataService Consumer端需要在调用时决定调用哪个后端服务(也就是说要进行流量管理),因此在Consumer端启动时先通过应用名查询注册中心得到后端Provider的地址,并完成订阅。拿到Provider地址后,直接发起元数据查询,通过MetadataService得到服务端元数据,并且会有一个监听器,监听服务端元数据的变化(即reversion机制),一旦服务端元数据发生变化,Consumer就会重新发起MetadataService调用获取元数据。 源码解析 Provider暴露MetadataService Dubbo暴露服务的起点是org.apache.dubbo.config.ServiceConfig#export,当然实际的暴露动作是委托给org.apache.dubbo.rpc.Protocol#export。但这是单个服务的暴露,MetadataService肯定是要在所有业务服务暴露完之后再暴露,因此我们需要从DubboBootstrap中去寻找答案(因为现在Dubbo的启动都是由这个类完成)。 DubboBootstrap这个类的结构我们这次不讲,我们需要了解的是Dubbo3.0有两个Deployer类负责启动。ApplicationDeployer负责应用的初始化和启动,ModuleDeployer负责服务的暴露和引用。因此我们可以到ModuleDeployer中找答案。 org.apache.dubbo.config.deploy.DefaultModuleDeployer#start 从这个模板方法里面我们可以看到在暴露和引用服务之后我们会调用org.apache.dubbo.config.deploy.DefaultModuleDeployer#onModuleStarted后续的调用链路如下。 org.apache.dubbo.config.deploy.DefaultApplicationDeployer#notifyModuleChanged -> org.apache.dubbo.config.deploy.DefaultApplicationDeployer#prepareApplicationInstance -> org.apache.dubbo.config.deploy.DefaultApplicationDeployer#exportMetadataService -> org.apache.dubbo.config.metadata.ExporterDeployListener#onModuleStarted -> org.apache.dubbo.config.metadata.ConfigurableMetadataServiceExporter#export 最终暴露元数据服务的动作就在org.apache.dubbo.config.metadata.ConfigurableMetadataServiceExporter#export中。 Consumer监听MetadataService 这里和Provider的分析过程刚好相反,因为Consumer需要让每个ReferenceConfig都要监听服务端的元数据变化,否则流量管理能力会失效。所以,我们直接从org.apache.dubbo.config.ReferenceConfig#get看就好。 首先初始化过程会调用init方法,我们直接跳过一些其他的配置代码,我们知道代理类的创建其实也是委托给org.apache.dubbo.rpc.Protocol#refer,具体来说调用链路如下。 org.apache.dubbo.registry.integration.RegistryProtocol#refer -> org.apache.dubbo.registry.integration.RegistryProtocol#doRefer -> org.apache.dubbo.registry.integration.RegistryProtocol#interceptInvoker -> org.apache.dubbo.registry.client.migration.MigrationRuleListener#onRefer -> org.apache.dubbo.registry.client.migration.MigrationRuleHandler#doMigrate -> org.apache.dubbo.registry.client.migration.MigrationRuleHandler#refreshInvoker -> org.apache.dubbo.registry.client.migration.ServiceDiscoveryMigrationInvoker#migrateToApplicationFirstInvoker -> org.apache.dubbo.registry.client.migration.MigrationInvoker#refreshServiceDiscoveryInvoker -> org.apache.dubbo.registry.integration.RegistryProtocol#getServiceDiscoveryInvoker -> org.apache.dubbo.registry.integration.RegistryProtocol#doCreateInvoker -> org.apache.dubbo.registry.integration.DynamicDirectory#subscribe ->org.apache.dubbo.registry.client.ServiceDiscoveryRegistry#subscribe -> org.apache.dubbo.registry.client.ServiceDiscoveryRegistry#doSubscribe -> org.apache.dubbo.registry.client.ServiceDiscoveryRegistry#subscribeURLs 最终在org.apache.dubbo.registry.client.ServiceDiscoveryRegistry#subscribeURLs中我们能看到整个Consumer端的服务自省过程。 关键代码如下: public class ServiceDiscoveryRegistry { protected void subscribeURLs(URL url, NotifyListener listener, Set<String> serviceNames) { serviceNames = toTreeSet(serviceNames); String serviceNamesKey = toStringKeys(serviceNames); //... Lock appSubscriptionLock = getAppSubscription(serviceNamesKey); try { appSubscriptionLock.lock(); ServiceInstancesChangedListener serviceInstancesChangedListener = serviceListeners.get(serviceNamesKey); if (serviceInstancesChangedListener == null) { serviceInstancesChangedListener = serviceDiscovery.createListener(serviceNames); serviceInstancesChangedListener.setUrl(url); for (String serviceName : serviceNames) { List<ServiceInstance> serviceInstances = serviceDiscovery.getInstances(serviceName); if (CollectionUtils.isNotEmpty(serviceInstances)) { serviceInstancesChangedListener.onEvent(new ServiceInstancesChangedEvent(serviceName, serviceInstances)); } } serviceListeners.put(serviceNamesKey, serviceInstancesChangedListener); } if (!serviceInstancesChangedListener.isDestroyed()) { serviceInstancesChangedListener.setUrl(url); listener.addServiceListener(serviceInstancesChangedListener); serviceInstancesChangedListener.addListenerAndNotify(protocolServiceKey, listener); serviceDiscovery.addServiceInstancesChangedListener(serviceInstancesChangedListener); } else { logger.info(String.format("Listener of %s has been destroyed by another thread.", serviceNamesKey)); serviceListeners.remove(serviceNamesKey); } } finally { appSubscriptionLock.unlock(); } } } 从上面看到org.apache.dubbo.registry.client.event.listener.ServiceInstancesChangedListener被创建,并且在监听ServiceInstancesChangedEvent,通过org.apache.dubbo.registry.client.event.listener.ServiceInstancesChangedListener#onEvent -> org.apache.dubbo.registry.client.event.listener.ServiceInstancesChangedListener#doOnEvent 这里就是Consumer端发生服务自省的地方。从代码中可以看到对MetadataService的调用是通过org.apache.dubbo.registry.client.ServiceDiscovery#getRemoteMetadata(java.lang.String, java.util.List<org.apache.dubbo.registry.client.ServiceInstance>)

                Hero Image
                Dubbo3.0-SpringBoot+Triple

                前言 前面讲过使用Dubbo原生API的方式使用Triple,我们发现这个方式不太方便,实际开发中可能并不是很常见这种使用方式,今天我们讲讲在如何在日常项目中使用Triple协议。 快速搭建 环境准备 操作系统:Windowns JDK版本:1.8 编辑工具:IntelliJ IDEA Community Edition 2021.1.2 x64 Zookeeper:3.5.6(3.4或3.5以上的版本都行) 创建项目 项目有三个模块,分别是api、consumer、provider。proto文件都放在api中,consumer放消费端的代码,provider放服务提供者的代码。 api 依赖 <dependencies> <dependency> <groupId>org.apache.dubbo</groupId> <artifactId>dubbo</artifactId> <version>3.0.9</version> </dependency> <dependency> <groupId>io.grpc</groupId> <artifactId>grpc-all</artifactId> <version>1.44.1</version> </dependency> </dependencies> proto插件 这里也是定义在pom中 <build> <extensions> <extension> <groupId>kr.motd.maven</groupId> <artifactId>os-maven-plugin</artifactId> <version>1.6.1</version> </extension> </extensions> <plugins> <plugin> <groupId>org.xolstice.maven.plugins</groupId> <artifactId>protobuf-maven-plugin</artifactId> <version>0.6.1</version> <configuration> <protocArtifact>com.google.protobuf:protoc:${protoc.version}:exe:${os.detected.classifier} </protocArtifact> <pluginId>grpc-java</pluginId> <pluginArtifact>io.grpc:protoc-gen-grpc-java:${grpc.version}:exe:${os.detected.classifier} </pluginArtifact> <protocPlugins> <protocPlugin> <id>dubbo</id> <groupId>org.apache.dubbo</groupId> <artifactId>dubbo-compiler</artifactId> <version>0.0.4.1</version> <mainClass>org.apache.dubbo.gen.tri.Dubbo3TripleGenerator</mainClass> </protocPlugin> </protocPlugins> </configuration> <executions> <execution> <goals> <goal>compile</goal> <goal>test-compile</goal> <goal>compile-custom</goal> <goal>test-compile-custom</goal> </goals> </execution> </executions> </plugin> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>${maven-compiler-plugin.version}</version> <configuration> <source>1.8</source> <target>1.8</target> </configuration> </plugin> </plugins> </build> proto文件 地址:看这里

                  Hero Image
                  Dubbo3.0-Triple协议

                  前言 Dubbo3.0推出了一个新的协议Triple,我们今天简单看看它的用法。 Triple总览 官方说明:看这里 ,总结为以下几点: 打通 gRPC 生态。 网关友好。(dubbo协议 只能 HTTP 转 泛化Dubbo 调用后端服务,现在可以接入Ingress方案) 异步和流式支持。(dubbo协议 多用于响应式的场景,不适合传大规模的数据。3.0增加流式场景,支持大规模数据的传输) 快速搭建 环境准备 操作系统:Windowns JDK版本:1.8 编辑工具:IntelliJ IDEA Community Edition 2021.1.2 x64 Zookeeper:3.5.6(3.4或3.5以上的版本都行) 创建项目 项目有三个模块,分别是api、consumer、provider。proto文件都放在api中,consumer放消费端的代码,provider放服务提供者的代码。 api 依赖 <dependencies> <dependency> <groupId>org.apache.dubbo</groupId> <artifactId>dubbo</artifactId> <version>3.0.9</version> </dependency> <dependency> <groupId>io.grpc</groupId> <artifactId>grpc-all</artifactId> <version>1.44.1</version> </dependency> </dependencies> proto插件 这里也是定义在pom中 <build> <extensions> <extension> <groupId>kr.motd.maven</groupId> <artifactId>os-maven-plugin</artifactId> <version>1.6.1</version> </extension> </extensions> <plugins> <plugin> <groupId>org.xolstice.maven.plugins</groupId> <artifactId>protobuf-maven-plugin</artifactId> <version>0.6.1</version> <configuration> <protocArtifact>com.google.protobuf:protoc:${protoc.version}:exe:${os.detected.classifier} </protocArtifact> <pluginId>grpc-java</pluginId> <pluginArtifact>io.grpc:protoc-gen-grpc-java:${grpc.version}:exe:${os.detected.classifier} </pluginArtifact> <protocPlugins> <protocPlugin> <id>dubbo</id> <groupId>org.apache.dubbo</groupId> <artifactId>dubbo-compiler</artifactId> <version>0.0.4.1</version> <mainClass>org.apache.dubbo.gen.tri.Dubbo3TripleGenerator</mainClass> </protocPlugin> </protocPlugins> </configuration> <executions> <execution> <goals> <goal>compile</goal> <goal>test-compile</goal> <goal>compile-custom</goal> <goal>test-compile-custom</goal> </goals> </execution> </executions> </plugin> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>${maven-compiler-plugin.version}</version> <configuration> <source>1.8</source> <target>1.8</target> </configuration> </plugin> </plugins> </build> proto文件 地址:看这里

                    Hero Image
                    避坑:@Around与@Transactional混用导致事务不回滚

                    前言 上个月,同事出于好奇在群里问AOP的环绕通知与事务注解混合用会不会导致出现异常不回滚的情况。这个问题我一下子回答不上来,因为平时没这样用过,在好奇心的驱使下,我调试了半天终于得到结果,今天我就展开讲讲。(源码解读在最后面,感兴趣的可以看看。) 结论 首先告诉大家的是,同时使用AOP环绕通知和事务注解之后,最终生成的拦截器链的相对顺序是事务的拦截器在前面,AOP环绕通知的拦截器在后面。 在事务的实现中将拦截器的执行过程包裹在了try-catch块中,发生异常后根据配置来决定是否回滚事务。(详见org.springframework.transaction.interceptor.TransactionInterceptor#invoke),因此事务后面的拦截器都会影响事务的执行结果。如果在AOP环绕通知里面将拦截器链执行结果中的异常给吞掉,那么事务就会正常提交而不会回滚。 示例 业务代码 业务代码中直接抛出异常,代码如下所示。 package com.example.demo.aspect; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; /** * @Author Paul * @Date 2022/7/3 15:52 */ @Service public class CustomService { @Transactional(rollbackFor = Exception.class) public void echo(){ boolean s = true; if (s){ throw new RuntimeException("test"); } System.out.println("Hello------"); } } 环绕通知 环绕通知中捕捉异常并打印日志,代码如下所示。 package com.example.demo.aspect; import org.aspectj.lang.ProceedingJoinPoint; import org.aspectj.lang.annotation.Around; import org.aspectj.lang.annotation.Aspect; import org.aspectj.lang.annotation.Pointcut; import org.springframework.stereotype.Component; /** * @Author Paul * @Date 2022/7/3 15:49 */ @Component @Aspect public class CustomAspect { @Pointcut("execution(* com.example.demo.aspect..*(..))") public void pointcut(){ } @Around("pointcut()") public void around(ProceedingJoinPoint joinPoint){ try { joinPoint.proceed(); } catch (Throwable throwable) { throwable.printStackTrace(); } } } 测试类 这里是用get请求来测试(本来应该用 unit test 来测试的,但是懒得写代码了,手动测试和 UT 的效果一样),代码如下所示。

                      Hero Image
                      从零开始学gRPC(二)

                      前言 gRPC的大致功能点相信大家一定已经有些了解。但是,在日常开发中我们很少会使用原生的gRPC,更多的是希望gRPC能集成到现有应用的开发环境中。这篇文章会围绕如何以简单的方式将gRPC集成到SpringBoot展开讲解。 快速搭建 环境准备 操作系统:Windowns JDK版本:1.8 编辑工具:IntelliJ IDEA Community Edition 2021.1.2 x64 框架选型 目前开源切比较热门的grpc和SpringBoot整合框架有两种可选方案。分别是LogNet 和 yidongnan 。(这两者稍微有些区别,在最后我会有个对比) 这里我使用后者作为demo演示。 创建项目 自己创建一个maven项目就好,没有别的要求(All in one,暂时不需要分模块)。想要偷懒的同学也可以借鉴之前的项目 ,但是我建议自己动手尝试搭建比较好。 MAVEN依赖 注意:我只列出了关键依赖。而且注册中心选用的是zookeeper,要注意zookeeper的版本,可以参考 这里 。 <dependencies> <dependency> <groupId>net.devh</groupId> <artifactId>grpc-spring-boot-starter</artifactId> <version>2.13.1.RELEASE</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> <version>2.6.8</version> </dependency> <!-- spring-cloud-zookeeper-starter--> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-zookeeper-discovery</artifactId> <version>3.1.2</version> </dependency> </dependencies> 配置 server: port: 8081 spring: application: name: demo-grpc cloud: zookeeper: connect-string: localhost:2181 grpc: client: greeterService1: address: discovery:/demo-grpc negotiationType: PLAINTEXT server: port: 9091 这个地方的 negotiationType要注意点(不设置的话需要对客户端传输进行加密) SLB测试 本地启动zookeeper 本地启动两次(记得改端口,tomcat和grpc的都要改),测试代码如下。 测试结果展示 我前端都是通过一个地址访问的,请求通过gRPC的负载配置被分发到不同的端口去了,由此可见SLB生效了。 LogNet 与 yidongnan 共性:

                        Hero Image
                        从零开始学gRPC(一)

                        前言 gRPC作为当前最热门的RPC框架之一,以其独特的跨语言、跨平台特性,赢得许多公司的青睐。 老实说,之前我只是道听途说并没有认真去研究,今天我会根据官网的demo展开介绍整个gRPC的功能, 后面一篇会介绍gRPC如何整合到SpringCloud。 我这里只提供了搭建demo工程的资料,建议自己动手来操作。没有截图项目也是因为官方的资料相当齐全,没必要重复造轮子。 gRPC总览 在直接使用gRPC之前,我们先了解下它的所有特性。官方描述 我就不展开讲了,gRPC有以下几点主要功能: 使用Protocol Buffer 定义服务。 语言和平台的中立性 双向流式通讯 基于HTTP2.0身份认证、SLB、tracing、health check 组件可扩展。 快速搭建 环境准备 操作系统:Windowns JDK版本:1.8 编辑工具:IntelliJ IDEA Community Edition 2021.1.2 x64 创建项目 自己创建一个maven项目就好,没有别的要求(All in one,暂时不需要分模块)。想要偷懒的同学也可以借鉴我的项目,这里 已经给你准备好了。 但是,建议自己动手尝试搭建比较好,不然印象不太深刻,等于没学(有个词儿叫行动废人,虽然不好听但目的是希望能让你动手实践)。 说明: 这里之所以没用官方的demo是因为我本地实在编译不出来,而且官方给的项目太大,对于学习来说没必要全部下载。(我只想要examples模块,但是必须强制全部下载,就很烦) 当然,要下载的朋友看这里 。 官方的中文文档地址是这个 ,可以按照官方提供的文档做。但如果只是玩具,那我这种All on one的方式我觉得更简单。 Helloworld 我们直接拿官方的例子来讲解,官方代码地址 。里面有很多例子,我这里只讲部分。 proto文件 我们到官方地址 抄作业时要注意,整个目录的结构。 proto文件只能放在模块的src/main下面,注意位置还有名字是否弄错,不然生成不了代码。 代码 服务端:代码地址 客户端:代码地址 特性讲解 我们知道gRPC不仅仅是一个helloworld就能描述清楚的,我下面将官方的代码例子做个分类,顺便总结下。 Stream 例子 proto: 代码地址 代码目录: 代码地址 总结 上述代码想表达的意思是,服务端在收到全部的客户端数据之后再响应回客户端处理结果。实际情况可以是服务端先处理一部分,然后返回部分。 也可以是任意顺序, 重要的是了解如何使用Stream的相关API来做交互。 TLS 例子 proto: 代码地址 代码目录: 代码地址 (还是helloworld的例子) 总结 这个就没啥好说的,用了Http2.0的TLS特性,默认就是开启的。如果客户端要传明文则必须在channel中配置,如下所示。 SLB 和 Health Check 例子 proto: 代码地址 代码目录: 代码地址 官方并没提供,下面都是自己的尝试