Dubbo3 基于 Kubernetes Informer 的服务发现原理解析

简介: ## List/Watch 机制介绍List/Watch机制是Kubernetes中实现集群控制模块最核心的设计之一,它采用统一的异步消息处理机制,保证了消息的实时性、可靠性、顺序性和性能等,为声明式风格的API奠定了良好的基础。`list`是调用`list API`获取资源列表,基于`HTTP`短链接实现。`watch`则是调用`watch API`监听资源变更事件,基于`HTTP

List/Watch 机制介绍

List/Watch机制是Kubernetes中实现集群控制模块最核心的设计之一,它采用统一的异步消息处理机制,保证了消息的实时性、可靠性、顺序性和性能等,为声明式风格的API奠定了良好的基础。

list是调用list API获取资源列表,基于HTTP短链接实现。

watch则是调用watch API监听资源变更事件,基于HTTP 长链接,通过Chunked transfer encoding(分块传输编码)来实现消息通知。

当客户端调用watch API时,kube-apiserverresponseHTTP Header中设置Transfer-Encoding的值为chunked,表示采用分块传输编码,客户端收到该信息后,便和服务端连接,并等待下一个数据块,即资源的事件信息。例如:

$ curl -i http://{kube-api-server-ip}:8080/api/v1/watch/endpoints?watch=yes
HTTP/1.1 200 OK
Content-Type: application/json
Transfer-Encoding: chunked
Date: Thu, 14 Seo 2022 20:22:59 GMT
Transfer-Encoding: chunked

{"type":"ADDED", "object":{"kind":"Endpoints","apiVersion":"v1",...}}
{"type":"ADDED", "object":{"kind":"Endpoints","apiVersion":"v1",...}}
{"type":"MODIFIED", "object":{"kind":"Endpoints","apiVersion":"v1",...}}

Dubbo 基于 Watch 的服务发现

在Dubbo3.1版本之前,dubbo-kubernetes服务发现是通过Fabric8 Kubernetes Java Client提供的watch API监听kube-apiserver中资源的create、update和delete事件,如下:

private void watchEndpoints(ServiceInstancesChangedListener listener, String serviceName) {
    Watch watch = kubernetesClient
        .endpoints()
        .inNamespace(namespace)
        .withName(serviceName)
        .watch(new Watcher<Endpoints>() {
            @Override
            public void eventReceived(Action action, Endpoints resource) {
                if (logger.isDebugEnabled()) {
                    logger.debug("Received Endpoint Event. Event type: " + action.name() +
                        ". Current pod name: " + currentHostname);
                }
                notifyServiceChanged(serviceName, listener);
            }
            @Override
            public void onClose(WatcherException cause) {
                // ignore
            }
        });
    ENDPOINTS_WATCHER.put(serviceName, watch);
}

private void notifyServiceChanged(String serviceName, ServiceInstancesChangedListener listener) {
    long receivedTime = System.nanoTime();
    ServiceInstancesChangedEvent event;
    event = new ServiceInstancesChangedEvent(serviceName, getInstances(serviceName));
    AtomicLong updateTime = SERVICE_UPDATE_TIME.get(serviceName);
    long lastUpdateTime = updateTime.get();
    if (lastUpdateTime <= receivedTime) {
        if (updateTime.compareAndSet(lastUpdateTime, receivedTime)) {
            listener.onEvent(event);
            return;
        }
    }
    if (logger.isInfoEnabled()) {
        logger.info("Discard Service Instance Data. " +
            "Possible Cause: Newer message has been processed or Failed to update time record by CAS. " +
            "Current Data received time: " + receivedTime + ". " +
            "Newer Data received time: " + lastUpdateTime + ".");
    }
}

@Override
public List<ServiceInstance> getInstances(String serviceName) throws NullPointerException {
    // 直接调用kube-apiserver
    Endpoints endpoints =
        kubernetesClient
            .endpoints()
            .inNamespace(namespace)
            .withName(serviceName)
            .get();
    return toServiceInstance(endpoints, serviceName);
}

监听到资源变化后,调用notifyServiceChanged方法从kube-apiserver全量拉取资源list数据,保持dubbo本地侧服务列表的实时性,并做相应的事件发布。

这样的操作存在很严重的问题,由于watch对应的回调函数会将更新的资源返回,dubbo社区考虑到维护成本较高,之前并没有在本地维护关于CRD资源的缓存,这样每次监听到变化后调用listkube-apiserver获取对应serviceName的endpoints信息,无疑增加了一次对kube-apiserver的直接访问。

clint-go为解决客户端需要自行维护缓存的问题,推出了informer机制。

Informer 机制介绍

Informer模块是Kubernetes中的基础组件,以List/Watch为基础,负责各组件与kube-apiserver的资源与事件同步。Kubernetes中的组件,如果要访问Kubernetes中的Object,绝大部分情况下会使用Informer中的Lister()方法,而非直接调用kube-apiserver。

以Pod资源为例,介绍下informer的关键逻辑(与下图步骤一一对应):

  1. Informer 在初始化时,Reflector 会先调用 List 获得所有的 Pod,同时调用Watch长连接监听kube-apiserver。
  2. Reflector 拿到全部 Pod 后,将Add Pod这个事件发送到 DeltaFIFO。
  3. DeltaFIFO随后pop这个事件到Informer处理。
  4. Informer向Indexer发布Add Pod事件。
  5. Indexer接到通知后,直接操作Store中的数据(key->value格式)。
  6. Informer触发EventHandler回调。
  7. 将key推到Workqueue队列中。
  8. 从WorkQueue中pop一个key。
  9. 然后根据key去Indexer取到val。根据当前的EventHandler进行Add Pod操作(用户自定义的回调函数)。
  10. 随后当Watch到kube-apiserver资源有改变的时候,再重复2-9步骤。


(来源于kubernetes/sample-controller)

Informer 关键设计

  • 本地缓存:Informer只会调用k8s List和Watch两种类型的API。Informer在初始化的时,先调用List获得某种resource的全部Object,缓存在内存中; 然后,调用Watch API去watch这种resource,去维护这份缓存; 最后,Informer就不再调用kube-apiserver。Informer抽象了cache这个组件,并且实现了store接口,后续获取资源直接通过本地的缓存来进行获取。
  • 无界队列:为了协调数据生产与消费的不一致状态,在客户端中通过实现了一个无界队列DeltaFIFO来进行数据的缓冲,当reflector获取到数据之后,只需要将数据推到到DeltaFIFO中,则就可以继续watch后续事件,从而减少阻塞时间,如上图2-3步骤所示。
  • 事件去重:在DeltaFIFO中,如果针对某个资源的事件重复被触发,则就只会保留相同事件最后一个事件作为后续处理,有resourceVersion唯一键保证,不会重复消费。
  • 复用连接:每一种资源都实现了Informer机制,允许监控不同的资源事件。为了避免同一个资源建立多个Informer,每个Informer使用一个Reflector与apiserver建立链接,导致kube-apiserver负载过高的情况,k8s中抽象了sharedInformer的概念,即共享的informer, 可以使同一类资源Informer共享一个Reflector。内部定义了一个map字段,用于存放所有Infromer的字段。针对同一资源只建立一个连接,减小kube-apiserver的负载。

Dubbo 引入 informer 机制后的服务发现

资源监听更换 Informer API

以Endpoints和Pods为例,将原本的Watch替换为Informer,回调函数分别为onAdd、onUpdate、onDelete,回调参数传的都是informer store中的资源全量值。

/**
 * 监听Endpoints
 */
private void watchEndpoints(ServiceInstancesChangedListener listener, String serviceName) {
    SharedIndexInformer<Endpoints> endInformer = kubernetesClient
            .endpoints()
            .inNamespace(namespace)
            .withName(serviceName)
            .inform(new ResourceEventHandler<Endpoints>() {
                @Override
                public void onAdd(Endpoints endpoints) {
                    if (logger.isDebugEnabled()) {
                        logger.debug("Received Endpoint Event. Event type: added. Current pod name: " + currentHostname + ". Endpoints is: " + endpoints);
                    }
                    notifyServiceChanged(serviceName, listener, toServiceInstance(endpoints, serviceName));
                }

                @Override
                public void onUpdate(Endpoints oldEndpoints, Endpoints newEndpoints) {
                    if (logger.isDebugEnabled()) {
                        logger.debug("Received Endpoint Event. Event type: updated. Current pod name: " + currentHostname + ". The new Endpoints is: " + newEndpoints);
                    }
                    notifyServiceChanged(serviceName, listener, toServiceInstance(newEndpoints, serviceName));
                }

                @Override
                public void onDelete(Endpoints endpoints, boolean deletedFinalStateUnknown) {
                    if (logger.isDebugEnabled()) {
                        logger.debug("Received Endpoint Event. Event type: deleted. Current pod name: " + currentHostname + ". Endpoints is: " + endpoints);
                    }
                    notifyServiceChanged(serviceName, listener, toServiceInstance(endpoints, serviceName));
                }
            });
    // 将endInformer存入ENDPOINTS_INFORMER,便于优雅下线统一管理
    ENDPOINTS_INFORMER.put(serviceName, endInformer);
}
/**
 * 监听Pods
 */
private void watchPods(ServiceInstancesChangedListener listener, String serviceName) {
    Map<String, String> serviceSelector = getServiceSelector(serviceName);
    if (serviceSelector == null) {
        return;
    }
    SharedIndexInformer<Pod> podInformer = kubernetesClient
            .pods()
            .inNamespace(namespace)
            .withLabels(serviceSelector)
            .inform(new ResourceEventHandler<Pod>() {
                @Override
                public void onAdd(Pod pod) {
                    if (logger.isDebugEnabled()) {
                        logger.debug("Received Pods Event. Event type: added. Current pod name: " + currentHostname + ". Pod is: " + pod);
                    }
                    // 不处理
                }
                @Override
                public void onUpdate(Pod oldPod, Pod newPod) {
                    if (logger.isDebugEnabled()) {
                        logger.debug("Received Pods Event. Event type: updated. Current pod name: " + currentHostname + ". new Pod is: " + newPod);
                    }
                    // TODO 后续要监听Pod label的变化
                    notifyServiceChanged(serviceName, listener, getInstances(serviceName));
                }

                @Override
                public void onDelete(Pod pod, boolean deletedFinalStateUnknown) {
                    if (logger.isDebugEnabled()) {
                        logger.debug("Received Pods Event. Event type: deleted. Current pod name: " + currentHostname + ". Pod is: " + pod);
                    }
                    // 不处理
                }
            });
    // 将podInformer存入PODS_INFORMER,便于优雅下线统一管理
    PODS_INFORMER.put(serviceName, podInformer);
}

/**
 * 通知订阅者Service改变
 */
private void notifyServiceChanged(String serviceName, ServiceInstancesChangedListener listener, List<ServiceInstance> serviceInstanceList) {
    long receivedTime = System.nanoTime();
    ServiceInstancesChangedEvent event;
    event = new ServiceInstancesChangedEvent(serviceName, serviceInstanceList);
    AtomicLong updateTime = SERVICE_UPDATE_TIME.get(serviceName);
    long lastUpdateTime = updateTime.get();
    if (lastUpdateTime <= receivedTime) {
        if (updateTime.compareAndSet(lastUpdateTime, receivedTime)) {
            // 发布事件
            listener.onEvent(event);
            return;
        }
    }
    if (logger.isInfoEnabled()) {
        logger.info("Discard Service Instance Data. " +
                "Possible Cause: Newer message has been processed or Failed to update time record by CAS. " +
                "Current Data received time: " + receivedTime + ". " +
                "Newer Data received time: " + lastUpdateTime + ".");
    }
}

getInstances() 优化

引入informer后,无需直接调用list接口,而是直接从informer的store中获取,减少对kube-apiserver的直接调用。

public List<ServiceInstance> getInstances(String serviceName) throws NullPointerException {
    Endpoints endpoints = null;
    SharedIndexInformer<Endpoints> endInformer = ENDPOINTS_INFORMER.get(serviceName);
    if (endInformer != null) {
        // 直接从informer的store中获取Endpoints信息
        List<Endpoints> endpointsList = endInformer.getStore().list();
        if (endpointsList.size() > 0) {
            endpoints = endpointsList.get(0);
        }
    }
    // 如果endpoints经过上面处理仍为空,属于异常情况,那就从kube-apiserver拉取
    if (endpoints == null) {
        endpoints = kubernetesClient
                .endpoints()
                .inNamespace(namespace)
                .withName(serviceName)
                .get();
    }

    return toServiceInstance(endpoints, serviceName);
}

结论

优化为Informer后,Dubbo的服务发现不用每次直接调用kube-apiserver,减小了kube-apiserver的压力,也大大减少了响应时间,助力Dubbo从传统架构迁移到Kubernetes中。

相关实践学习
深入解析Docker容器化技术
Docker是一个开源的应用容器引擎,让开发者可以打包他们的应用以及依赖包到一个可移植的容器中,然后发布到任何流行的Linux机器上,也可以实现虚拟化,容器是完全使用沙箱机制,相互之间不会有任何接口。Docker是世界领先的软件容器平台。开发人员利用Docker可以消除协作编码时“在我的机器上可正常工作”的问题。运维人员利用Docker可以在隔离容器中并行运行和管理应用,获得更好的计算密度。企业利用Docker可以构建敏捷的软件交付管道,以更快的速度、更高的安全性和可靠的信誉为Linux和Windows Server应用发布新功能。 在本套课程中,我们将全面的讲解Docker技术栈,从环境安装到容器、镜像操作以及生产环境如何部署开发的微服务应用。本课程由黑马程序员提供。 &nbsp; &nbsp; 相关的阿里云产品:容器服务 ACK 容器服务 Kubernetes 版(简称 ACK)提供高性能可伸缩的容器应用管理能力,支持企业级容器化应用的全生命周期管理。整合阿里云虚拟化、存储、网络和安全能力,打造云端最佳容器化应用运行环境。 了解产品详情: https://www.aliyun.com/product/kubernetes
目录
相关文章
|
缓存 Kubernetes Docker
GitLab Runner 全面解析:Kubernetes 环境下的应用
GitLab Runner 是 GitLab CI/CD 的核心组件,负责执行由 `.gitlab-ci.yml` 定义的任务。它支持多种执行方式(如 Shell、Docker、Kubernetes),可在不同环境中运行作业。本文详细介绍了 GitLab Runner 的基本概念、功能特点及使用方法,重点探讨了流水线缓存(以 Python 项目为例)和构建镜像的应用,特别是在 Kubernetes 环境中的配置与优化。通过合理配置缓存和镜像构建,能够显著提升 CI/CD 流水线的效率和可靠性,助力开发团队实现持续集成与交付的目标。
|
Kubernetes 网络协议 Nacos
OpenAI 宕机思考丨Kubernetes 复杂度带来的服务发现系统的风险和应对措施
Kubernetes 体系基于 DNS 的服务发现为开发者提供了很大的便利,但其高度复杂的架构往往带来更高的稳定性风险。以 Nacos 为代表的独立服务发现系统架构简单,在 Kubernetes 中选择独立服务发现系统可以帮助增强业务可靠性、可伸缩性、性能及可维护性,对于规模大、增长快、稳定性要求高的业务来说是一个较理想的服务发现方案。希望大家都能找到适合自己业务的服务发现系统。
783 104
|
Kubernetes API 调度
Kubernetes 架构解析:理解其核心组件
【8月更文第29天】Kubernetes(简称 K8s)是一个开源的容器编排系统,用于自动化部署、扩展和管理容器化应用。它提供了一个可移植、可扩展的环境来运行分布式系统。本文将深入探讨 Kubernetes 的架构设计,包括其核心组件如何协同工作以实现这些功能。
1309 3
|
Kubernetes 监控 API
深入解析Kubernetes及其在生产环境中的最佳实践
深入解析Kubernetes及其在生产环境中的最佳实践
981 93
|
Kubernetes Linux 虚拟化
入门级容器技术解析:Docker和K8s的区别与关系
本文介绍了容器技术的发展历程及其重要组成部分Docker和Kubernetes。从传统物理机到虚拟机,再到容器化,每一步都旨在更高效地利用服务器资源并简化应用部署。容器技术通过隔离环境、减少依赖冲突和提高可移植性,解决了传统部署方式中的诸多问题。Docker作为容器化平台,专注于创建和管理容器;而Kubernetes则是一个强大的容器编排系统,用于自动化部署、扩展和管理容器化应用。两者相辅相成,共同推动了现代云原生应用的快速发展。
4784 11
|
负载均衡 监控 Dubbo
Dubbo 原理和机制详解(非常全面)
本文详细解析了 Dubbo 的核心功能、组件、架构设计及调用流程,涵盖远程方法调用、智能容错、负载均衡、服务注册与发现等内容。欢迎留言交流。关注【mikechen的互联网架构】,10年+BAT架构经验倾囊相授。
Dubbo 原理和机制详解(非常全面)
|
运维 Kubernetes Cloud Native
Kubernetes云原生架构深度解析与实践指南####
本文深入探讨了Kubernetes作为领先的云原生应用编排平台,其设计理念、核心组件及高级特性。通过剖析Kubernetes的工作原理,结合具体案例分析,为读者呈现如何在实际项目中高效部署、管理和扩展容器化应用的策略与技巧。文章还涵盖了服务发现、负载均衡、配置管理、自动化伸缩等关键议题,旨在帮助开发者和运维人员掌握利用Kubernetes构建健壮、可伸缩的云原生生态系统的能力。 ####
|
缓存 负载均衡 Dubbo
Dubbo技术深度解析及其在Java中的实战应用
Dubbo是一款由阿里巴巴开源的高性能、轻量级的Java分布式服务框架,它致力于提供高性能和透明化的RPC远程服务调用方案,以及SOA服务治理方案。
670 6
|
存储 Kubernetes 安全
在K8S中,你用的flannel是哪个工作模式及fannel的底层原理如何实现数据报文转发的?
在K8S中,你用的flannel是哪个工作模式及fannel的底层原理如何实现数据报文转发的?
|
存储 Kubernetes 调度
深度解析Kubernetes中的Pod生命周期管理
深度解析Kubernetes中的Pod生命周期管理

热门文章

最新文章

推荐镜像

更多