首页 文章 精选 留言 我的

精选列表

搜索[浏览器扩展],共10013篇文章
优秀的个人博客,低调大师

图解kubernetes资源扩展机制实现(上)

k8s目前主要支持CPU和内存两种资源,为了支持用户需要按需分配的其他硬件类型的资源的调度分配,k8s实现了设备插件框架(device plugin framework)来用于其他硬件类型的资源集成,比如现在机器学习要使用GPU等资源,今天来看下其内部的关键实现 1. 基础概念 1.1 集成方式 1.1.1 DaemonSet与服务 当我们要集成本地硬件的资源的时候,我们可以在当前节点上通过DaemonSet来运行一个GRPC服务,通过这个服务来进行本地硬件资源的上报与分配 1.1.2 服务注册设计 当提供硬件服务需要与kubelet进行通信的时候,则首先需要进行注册,注册的方式,则是通过最原始的底层的socket文件,并且通过Linux文件系统的inotify机制,来实现服务的注册 1.2 插件服务感知 1.2.1 Watcher Watcher主要是负责感知当前节点上注册的服务,当发现新的要注册的插件服务,则会产生对应的事件,注册到当前的kubelet中 1.2.2 期望状态与实际状态 这里的状态主要是指的是否需要注册,因为kubelet与对应的插件服务是通过网络进行通信的,当网络出现问题、或者对应的插件服务故障,则可能会导致服务注册失败,但此时对应的服务的socket还依旧存在,即对应的插件服务依旧存在 此时就会有两种状态:期望状态与实际状态, 因为socket存在所以服务的期望状态其实是需要注册这个插件服务,但是实际上因为某些原因,这个插件服务并没有完成注册,后续会不断的通过期望状态,调整实际状态,从而达到一致 1.2.3 协调器 协调器则就是完成上述两种状态之间操作的核心,其通过调用对应插件的回调函数,其实就是调用对应的grpc接口,来完成期望状态与实际状态的一致性 1.2.4 插件控制器 针对每种类型的插件,都会有对应的控制器,其实也就是实现对应设备注册和反注册并且完成底层资源的分配(Allocate)和收集(ListWatch)操作 2. 插件服务发现 2.1 核心数据结构 type Watcher struct { // 感知插件服务注册的socket的路径 path string fs utilfs.Filesystem // inotify监测插件服务socket变化 fsWatcher *fsnotify.Watcher stopped chan struct{} // 存储期望状态 desiredStateOfWorld cache.DesiredStateOfWorld } 2.2 初始化 初始化其实就是创建对应的目录 func (w *Watcher) init() error { klog.V(4).Infof("Ensuring Plugin directory at %s ", w.path) if err := w.fs.MkdirAll(w.path, 0755); err != nil { return fmt.Errorf("error (re-)creating root %s: %v", w.path, err) } return nil } 2.3 插件服务发现核心 go func(fsWatcher *fsnotify.Watcher) { defer close(w.stopped) for { select { case event := <-fsWatcher.Events: //如果发现对应目录的文件的变化,则会触发对应的事件 if event.Op&fsnotify.Create == fsnotify.Create { err := w.handleCreateEvent(event) if err != nil { klog.Errorf("error %v when handling create event: %s", err, event) } } else if event.Op&fsnotify.Remove == fsnotify.Remove { w.handleDeleteEvent(event) } continue case err := <-fsWatcher.Errors: if err != nil { klog.Errorf("fsWatcher received error: %v", err) } continue case <-stopCh: // In case of plugin watcher being stopped by plugin manager, stop // probing the creation/deletion of plugin sockets. // Also give all pending go routines a chance to complete select { case <-w.stopped: case <-time.After(11 * time.Second): klog.Errorf("timeout on stopping watcher") } w.fsWatcher.Close() return } } }(fsWatcher) 2.4 补偿机制 其实补偿机制主要是在重新启动kubelet的时候,需要将之前已经存在的socket重新注册到当前的kubelet中 func (w *Watcher) traversePluginDir(dir string) error { return w.fs.Walk(dir, func(path string, info os.FileInfo, err error) error { if err != nil { if path == dir { return fmt.Errorf("error accessing path: %s error: %v", path, err) } klog.Errorf("error accessing path: %s error: %v", path, err) return nil } switch mode := info.Mode(); { case mode.IsDir(): if err := w.fsWatcher.Add(path); err != nil { return fmt.Errorf("failed to watch %s, err: %v", path, err) } case mode&os.ModeSocket != 0: event := fsnotify.Event{ Name: path, Op: fsnotify.Create, } //TODO: Handle errors by taking corrective measures if err := w.handleCreateEvent(event); err != nil { klog.Errorf("error %v when handling create event: %s", err, event) } default: klog.V(5).Infof("Ignoring file %s with mode %v", path, mode) } return nil }) } 2.5 注册事件回调 注册其实就只需要感知到的socket文件路径传递给期望状态进行管理 func (w *Watcher) handlePluginRegistration(socketPath string) error { if runtime.GOOS == "windows" { socketPath = util.NormalizePath(socketPath) } // 调用期望状态进行更新 klog.V(2).Infof("Adding socket path or updating timestamp %s to desired state cache", socketPath) err := w.desiredStateOfWorld.AddOrUpdatePlugin(socketPath) if err != nil { return fmt.Errorf("error adding socket path %s or updating timestamp to desired state cache: %v", socketPath, err) } return nil } 2.6 删除事件回调 注册其实就只需要感知到的socket文件路径传递给期望状态进行管理 func (w *Watcher) handleDeleteEvent(event fsnotify.Event) { klog.V(6).Infof("Handling delete event: %v", event) socketPath := event.Name klog.V(2).Infof("Removing socket path %s from desired state cache", socketPath) w.desiredStateOfWorld.RemovePlugin(socketPath) } 3.期望状态与实际状态 3.1 插件信息 插件信息其实只是存储了对应socket的路径和最近更新的时间 type PluginInfo struct { SocketPath string Timestamp time.Time } 3.2 期望状态 期望状态与实际状态在数据结构上都是一样的,因为本质上只是为了存储插件的当前的状态信息,即更新时间,这里不在赘述 type desiredStateOfWorld struct { socketFileToInfo map[string]PluginInfo sync.RWMutex } type actualStateOfWorld struct { socketFileToInfo map[string]PluginInfo sync.RWMutex } 4.OperationExecutor 目前k8s中支持两大类的插件的管理一类是DevicePlugin即我们本文说的这些都是这种概念,一类是CSIPlugin,其中针对每一类DRiver的处理其实内部都是不一样的,那其实在操作之前就要先感知到当前的Driver是那种类型的 OperationExecutor主要就是做这件事的,其根据不同的plugin类型,生成不同的要执行的操作,即对应的Plugin类型获取对应的handler,就生成了一个要执行的操作 4.1 生成注册插件回调函数 4.1.1 通过socket连接对应的插件服务 registerPluginFunc := func() error { client, conn, err := dial(socketPath, dialTimeoutDuration) if err != nil { return fmt.Errorf("RegisterPlugin error -- dial failed at socket %s, err: %v", socketPath, err) } defer conn.Close() ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() infoResp, err := client.GetInfo(ctx, &registerapi.InfoRequest{}) if err != nil { return fmt.Errorf("RegisterPlugin error -- failed to get plugin info using RPC GetInfo at socket %s, err: %v", socketPath, err) } 4.1.2 根据插件类型验证服务 handler, ok := pluginHandlers[infoResp.Type] if !ok { if err := og.notifyPlugin(client, false, fmt.Sprintf("RegisterPlugin error -- no handler registered for plugin type: %s at socket %s", infoResp.Type, socketPath)); err != nil { return fmt.Errorf("RegisterPlugin error -- failed to send error at socket %s, err: %v", socketPath, err) } return fmt.Errorf("RegisterPlugin error -- no handler registered for plugin type: %s at socket %s", infoResp.Type, socketPath) } if infoResp.Endpoint == "" { infoResp.Endpoint = socketPath } if err := handler.ValidatePlugin(infoResp.Name, infoResp.Endpoint, infoResp.SupportedVersions); err != nil { if err = og.notifyPlugin(client, false, fmt.Sprintf("RegisterPlugin error -- plugin validation failed with err: %v", err)); err != nil { return fmt.Errorf("RegisterPlugin error -- failed to send error at socket %s, err: %v", socketPath, err) } return fmt.Errorf("RegisterPlugin error -- pluginHandler.ValidatePluginFunc failed") } 4.1.3 注册插件到实际状态 err = actualStateOfWorldUpdater.AddPlugin(cache.PluginInfo{ SocketPath: socketPath, Timestamp: timestamp, }) if err != nil { klog.Errorf("RegisterPlugin error -- failed to add plugin at socket %s, err: %v", socketPath, err) } // 调用插件的注册回调函数 if err := handler.RegisterPlugin(infoResp.Name, infoResp.Endpoint, infoResp.SupportedVersions); err != nil { return og.notifyPlugin(client, false, fmt.Sprintf("RegisterPlugin error -- plugin registration failed with err: %v", err)) } 4.1.4 通知对应的服务注册成功 if err := og.notifyPlugin(client, true, ""); err != nil { return fmt.Errorf("RegisterPlugin error -- failed to send registration status at socket %s, err: %v", socketPath, err) } 4.2 通过socket构建注册client func dial(unixSocketPath string, timeout time.Duration) (registerapi.RegistrationClient, *grpc.ClientConn, error) { ctx, cancel := context.WithTimeout(context.Background(), timeout) defer cancel() c, err := grpc.DialContext(ctx, unixSocketPath, grpc.WithInsecure(), grpc.WithBlock(), grpc.WithContextDialer(func(ctx context.Context, addr string) (net.Conn, error) { return (&net.Dialer{}).DialContext(ctx, "unix", addr) }), ) if err != nil { return nil, nil, fmt.Errorf("failed to dial socket %s, err: %v", unixSocketPath, err) } return registerapi.NewRegistrationClient(c), c, nil } 今天就先到这里,下一章会继续介绍如何组合上述组件以及默认的回调管理机制的实现,进探究到这里谢谢大家,感谢分享点赞,反转又不花钱 k8s源码阅读电子书地址: https://www.yuque.com/baxiaoshi/tyado3 > 微信号:baxiaoshi2020 > 关注公告号阅读更多源码分析文章 > 更多文章关注 www.sreguide.com > 本文由博客一文多发平台 OpenWrite 发布

优秀的个人博客,低调大师

一次Zookeeper 扩展之殇

一、背景 基于公司发展硬性需求,生产VM服务器要统一迁移到ZStack 虚拟化服务器。检查自己项目使用的服务器,其中zookeeper集群中招,所以需要进行迁移。 二、迁移计划 为了使迁移不对业务产生影响,所以最好是采用扩容 -> 缩容 的方式进行。 说明: 1.原生产集群为VM-1,VM-2,VM-3组成一个3节点的ZK集群; 2.对该集群扩容,增加至6节点(新增ZS-1,ZS-2,ZS-3),进行数据同步完成; 3.进行缩容,下掉原先来的三个节点(VM-1,VM-2,VM-3); 4.替换nginx解析地址。 OK! 目标很明确,过程也很清晰,然后开干。 三、步骤 (过程已在测试环境验证无问题): 对新增的三台服务器进行zk环境配置,和老集群配置一样即可,最好使用同一版本(版主使用的是3.4.6); 对老节点的zoo.cfg 增加新

优秀的个人博客,低调大师

扩展Spring Cloud Feign 实现自动降级

自动降级目的 在Spring Cloud 使用feign 的时候,需要明确指定fallback 策略,不然会提示错误。 先来看默认的feign service 是要求怎么做的。feign service 定义一个 factory 和 fallback 的类 @FeignClient(value = ServiceNameConstants.UMPS_SERVICE, fallbackFactory = RemoteLogServiceFallbackFactory.class) public interface RemoteLogService {} 但是我们大多数情况的feign 降级策略为了保证幂等都会很简单,输出错误日志即可。类似如下代码,在企业中开发非常不方便 @Slf4j @Component public class Rem

优秀的个人博客,低调大师

synchronized的功能的扩展:重入锁

重入锁 重入锁可以说是synchronized,Object.wait(),Object.notify()的一种替代品。 在JDK5的早期版本,重入锁的新能要比synchronized好很多,在JDK6后对synchronized进行可很多优化,使得他和重入锁的性能差距并不大。 重入锁使用java.util.concurrent.locks.ReentrantLock类实现,下面我么来看下重入锁的简单使用案例: import java.util.*; import java.util.concurrent.locks.ReentrantLock; public class ReenterLock implements Runnable { public static ReentrantLock lock=new ReentrantLock(); public static int i=0; @Override public void run() { // TODO Auto-generated method stub for(int j=0;j<10000000;j++) { lock.lock(); try { i++; }finally { lock.unlock(); } } } public static void main(String[] args) throws InterruptedException { // TODO Auto-generated method stub ReenterLock tl=new ReenterLock(); Thread t1=new Thread(tl); Thread t2=new Thread(tl); t1.start(); t2.start(); t1.join(); t2.join(); System.out.println(i); } } 运行程序可以得到结果为20000000. 重入锁具有很高的灵活性,需要开发人员手动用lock()与unlock()函数来指定何时加锁何时解锁,但是需要注意的是,在离开临界区的时候要记得释放锁,否则其他线程就没有机会在访问临界区了。(之所以叫重入锁是因为这种锁可以反复进入,但是记得一个线程同时获得多少锁,也必须释放相同的次数) 重入锁除了使用上的灵活性,还有一些高级的功能。 中断响应 对于synchronized来说,如果一个线程在等待锁,那么结果只有两种情况,要么获得锁继续执行,要么保持等待。但是使用重入锁。则提供了另外一种可能,就是线程可以被中断。也就是说,在等待锁的过程中,线程可以取消对锁的请求。有些时候这么做是很有必要的,比如和好朋友约好去打球,等了半小时没到,街道电话得知朋友临时有事,不能来了,那么就打道回府了。这种情况对于处理死锁有一定的帮助。使用lockInterruptibly()函数表示重入锁可以响应中断。 锁申请限时等待 除了等待外部通知外,要避免死锁还有一种办法,就是限时等待。我么可以通过tryLock方法进行一次限时的等待。 try{ if(lock.tryLock(5,TimeUnit.SECONDS)){ Thread.sleep(6000); }else{ System.out.println("get lock failed"); } }catch(InterruptdeException e){ e.printStackTrace(); }finally{ if(lock.isHeldByCurrentThread()) lock.unlock();} } 上面代码中,设置的5秒的限时等待,由于睡眠了6秒,或导致请求锁失败。tryLock也可以不带参数进行运行,这种情况下,线程会尝试获得锁,如果锁未被其他线程占用,申请锁成功,返回true,否则申请失败,线程也不会进行等待,直接返回false。 公平锁 大多数情况下锁是不公平的,也就是说不一定先请求就先获得锁,使用synchronized关键字实现锁,锁就是不公平的,重入锁可以实现公平锁,避免现象,也就是说,只要你排队,就能获得锁,但是由于需要维护一个有序队列,导致公平锁的实现成本比较高,性能相对也非常低下。重入锁有如下一个构造函数:public ReentrantLock(boolean fair)当参数为true时锁是公平的的。 重入锁的好搭档:Condition条件 Condition与Object.wait(),Object.notify()方法的作用大致相同,只不过是用来与重入锁相关联的。通过Lock接口(出入锁就实现乐这一接口)的Condition newCondition()方法可以生成与当前重入锁对象绑定的Condition实例,利用Condition对象,我们可以让线程在特定的时刻进行等待,或者在某一个特定的时刻得到通知,继续执行。Condition接口提供的方法如下。 void await() throws InterrupttedExceptionI; void awaitUninterruptibly(); long awaitNanos(long nanosTimeout) throws InterrupttedExceptionI; boolean await(long time,TimeUnit unit) throws InterrupttedExceptionI; boolean awaitUntil(Date deadline) throws InterrupttedExceptionI; void signal(); void siganlAll(); 以上个方法的含义如下: await()方法会使当前线程等待,同时会释放锁,当其他线程中使用signal()或signalAll()方法的时候,线程会重新获得锁继续执行。或者线程被中断是也能跳出等待。 awaitUninterruptibly()方法与await()方法基本相同,但是他不会在等待过程中相应或者相应中断。 signal()用于唤醒一个在等待中的线程,同理signalAll()用于唤醒所有在等待中的线程。

优秀的个人博客,低调大师

php 利用pcntl扩展实现高并发

pcntl_fork官方描述:pcntl_fork — 在当前进程当前位置产生分支(子进程)。译注:fork是创建了一个子进程,父进程和子进程 都从fork的位置开始向下继续执行,不同的是父进程执行过程中,得到的fork返回值为子进程 号,而子进程得到的是0。 通过一段代码来理解pcntl_fork如何进行多进程并发执行的 $i = 2; while($i >= 0){ $pid = pcntl_fork(); if($pid > 0){ //当前处于父进程中 }elseif($pid == 0){ //这里是子进程逻辑 exit();//子进程跑完了必须退出。否则子进程会继续执行循环创建新的进程 }else{ //进程创建失败 } 当前进程通过pcntl_fork创建了一个子进程,创建完成后当前进程继续执行,此时$pid为创建子进程的进程id所以会进入pid > 0 的判断区块。同时,当父进程执行pcntl_fork时,php又会重新开启一个进程。该进程会拷贝父进程的变量,并且从pcntl_fork之后开始执行,因为进程为子进程,所以pid为0,会执行pid==0里面的逻辑。子进程逻辑跑完了建议退出,否则会继续执行while循环。 100个号码批量分成10个进程发送示例 //随机取100条数据 $mobileList = ['135***','139***',...,'152*****']; insert($data); function insert(array $phoneList){ $cnt = count($phoneList); //测试数组大小 $slice = 10; //需要调用的进程数量 $master = array_chunk($phoneList,floor($cnt/(10))); $childList = []; echo "进程开始\r\n"; while($slice > 0) { $pid = pcntl_fork(); if($pid > 0){ $childList[$pid] = 1; echo $pid."\r\n"; //$pid>0表示当前还在执行父进程的代码 //这里最好啥都不做,每次执行pcntl_fork都会执行这里的代码。 //这里的代码执行完之后 会将$pid设置为0,然后jump到pcntl_fork代码之后,重新做判断; }elseif($pid == 0){ //这里写我们的逻辑 foreach($master[$slice-1] as $val) { sendMessage();//发送短信 } //子进程执行完之后务必需要关闭; exit(); }else { //程序发生错误也需要关闭程序 exit(); } $slice--; } // 等待所有子进程结束后回收资源 while(!empty($childList)){ $childPid = pcntl_wait($status); if ($childPid > 0){ unset($childList[$childPid]); } } }

资源下载

更多资源
Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册