首页 文章 精选 留言 我的

精选列表

搜索[玩转 Agent Bucket],共10002篇文章
优秀的个人博客,低调大师

玩转 Go 链路追踪

前言 链路追踪是每个微服务架构下必备的利器,go-zero 当然早已经为我们考虑好了,只需要在配置中添加配置即可使用。 关于 go-zero 如何追踪的原理追溯,之前已经有同学分享,这里我就不再多说,如果有想了解的同学去 https://mp.weixin.qq.com/s/hJEWcWc3PnGfWfbPCHfM9g 这个链接看就好了。默认会在 api 的中间件与 rpc 的 interceptor 添加追踪,如果有不了解 go-zero 默认如何使用默认的链路追踪的,请移步我的开源项目 go-zero-looklook 文档 https://github.com/Mikaelemmmm/go-zero-looklook/blob/main/doc/chinese/12-%E9%93%BE%E8%B7%AF%E8%BF%BD%E8%B8%AA.md。 今天我想讲的是,除了 go-zero 默认在 api 的 middleware 与 rpc 的 interceptor 中帮我们集成好的链路追踪,我们想自己在某些本地方法添加链路追踪代码或者我们想在 api 发送一个消息给 mq 服务时候想把整个链路包含 mq 的 producer、consumer 穿起来,在 go-zero 中该如何做。 场景 我们先简单讲一下我们的小 demo 的场景,一个请求进来调用 api 的 Login 方法,在 Login 方法中先调用 rpc 的 GetUserByMobile 方法,之后在调用 api 本地的 local 方法,紧接着调用 rabbitmq 传递消息到 mq 服务。 go-zero 默认集成了 jaeger、zinpink,这里我们就以 jaeger 为例 我们希望看到的链路是 rpc.GetUserByMobile" data-original="https://img-blog.csdnimg.cn/910d8fbfd04b4c6196b3e38479c1679b.png"> 也就是 api 衍生出来三条子链路,api.producerMq 有一条调用 mq.Consumer 的子链路。 我们想要将一个方法添加到链路中需要两个因素,一个 traceId,一个span,当我们在同一个 traceId 下开启 span 把相关的 span 都串联起来,如果想形成父子关系,就要把 span 之间相互串联起来,因为「微服务实践」公众号中讲解原理太多,我这里就简单提一下不涉及过多,如果不是特别熟悉原理可以看文章开头推荐的文章,这里我们只需要知道 traceId 与 spanId 关系就好。 核心业务代码 1、首先 API 中 LoginLogic 代码 type LoginLogic struct { logx.Logger ctx context.Context svcCtx *svc.ServiceContext } func NewLoginLogic(ctx context.Context, svcCtx *svc.ServiceContext) *LoginLogic { return &LoginLogic{ Logger: logx.WithContext(ctx), ctx: ctx, svcCtx: svcCtx, } } type MsgBody struct { Carrier *propagation.HeaderCarrier Msg string } func (l *LoginLogic) Login(req *types.RegisterReq) (*types.AccessTokenResp, error) { resp, err := l.svcCtx.UserRpc.GetUserByMobile(l.ctx, &usercenter.GetUserByMobileReq{ Mobile: req.Mobile, }) if err != nil { return &types.AccessTokenResp{}, nil } l.local() tracer := otel.GetTracerProvider().Tracer(trace.TraceName) spanCtx, span := tracer.Start(l.ctx, "send_msg_mq", oteltrace.WithSpanKind(oteltrace.SpanKindProducer)) carrier := &propagation.HeaderCarrier{} otel.GetTextMapPropagator().Inject(spanCtx, carrier) producer := rabbit.NewRabbitmqPublisher(RabbitmqDNS) msg := &MsgBody{ Carrier: carrier, Msg: req.Mobile, } b, err := json.Marshal(msg) if err != nil{ panic(err) } if err := producer.Publish(spanCtx, ExchangeName, RoutineKeys, b); err != nil { logx.Errorf("Publish Fail , msg :%s , err:%v", msg, err) } span.End() return &types.AccessTokenResp{ AccessExpire: resp.User.Id, }, err } func (l *LoginLogic) local() { tracer := otel.GetTracerProvider().Tracer(trace.TraceName) _ , span := tracer.Start(l.ctx, "local", oteltrace.WithSpanKind(oteltrace.SpanKindInternal)) defer span.End() // 执行你的代码 ..... } 2、rpc 中 GetUserByMobile 的代码 func (s *Logic) GetUserByMobile(context.Context, *usercenterPb.GetUserByMobileReq) (*usercenterPb.GetUserByMobileResp, error) { vo := &usercenterPb.UserVo{ Id: 1, } return &usercenterPb.GetUserByMobileResp{ User: vo, }, nil } 3、mq 中 Consumer 的代码 type MsgBody struct { Carrier *propagation.HeaderCarrier Msg string } func (c *consumer) Consumer(ctx context.Context, data []byte) error { var msg MsgBody if err := json.Unmarshal(data, &msg); err != nil { logx.Errorf(" consumer err : %v", err) } else { logx.Infof("consumerOne Consumer , msg:%+v", msg) wireContext := otel.GetTextMapPropagator().Extract(ctx, msg.Carrier) tracer := otel.GetTracerProvider().Tracer(trace.TraceName) _, span := tracer.Start(wireContext, "mq_consumer_msg", oteltrace.WithSpanKind(oteltrace.SpanKindConsumer)) defer span.End() } return nil } 代码详解 1、go-zero 默认集成 当一个请求进入 api 后,我们可以在 go-zero 源码中查看到 https://github.com/zeromicro/go-zero/blob/master/rest/engine.go#L92。go-zero 已经在 api 的 middleware 中帮我们添加了第一层 trace,当进入 Login 方法内,我们调用了 rpc 的 GetUserByMobile 方法,通过 go-zero 的源码 https://github.com/zeromicro/go-zero/blob/master/zrpc/internal/rpcserver.go#L55 可以看到在 rpc 的 interceptor 也默认帮我们添加好了,这两层都是 go-zero 默认帮我们做好的。 2、本地方法 当调用完 rpc 的 GetUserByMobile 之后,api 调用了本地的 local,如果我们想在整个链路上体现出来调用了本地 local 方法,那默认的 go-zero 是没有帮我们做的,需要我们手动来添加。 tracer := otel.GetTracerProvider().Tracer(trace.TraceName) _ , span := tracer.Start(l.ctx, "local", oteltrace.WithSpanKind(oteltrace.SpanKindInternal)) defer span.End() // 执行你的代码 ..... 我们通过上面代码拿到 tracer,ctx 之后开启一个 local 的 span,因为 start 时候会从 ctx 获取父 span 所以会将 local 方法与 Login 串联起父子调用关系,这样就将本次操作加入了这个链路 3、mq 的 producer 到 mq 的 consumer 我们在mq传递中如何串联起来这个链路呢?也就是形成 api.Login->api.producer->mq.Consumer。 想一下原理,虽然跨越了网络,api 可以通过 header 传递,rpc 可以通过 metadata 传递,那么 mq 是不是也可以通过 header、body 传递就可以了,按照这个想法来看下我门的代码。 tracer := otel.GetTracerProvider().Tracer(trace.TraceName) spanCtx , span := tracer.Start(l.ctx, "send_msg_mq", oteltrace.WithSpanKind(oteltrace.SpanKindProducer)) carrier := &propagation.HeaderCarrier{} otel.GetTextMapPropagator().Inject(spanCtx,carrier) producer := rabbit.NewRabbitmqPublisher(RabbitmqDNS) msg := &MsgBody{ Carrier: carrier, Msg: req.Mobile, } b , err := json.Marshal(msg) if err != nil{ panic(err) } if err := producer.Publish(spanCtx, ExchangeName, RoutineKeys, b); err != nil { logx.Errorf("Publish Fail, msg :%s, err:%v", msg, err) } span.End() 首先获取到了这个全局的 tracer,然后开启一个 producer 的 span,跟 local 方法一样,我们开启 producer 的 span 时候也是通过 ctx 获取到上一级父级 span,这样就可以将 producer 的 span 与 Login 形成父子 span 调用关系,那我们想将 producer 的 span 与 mq 的 consumer 中的 span 形成调用父子关系怎么做?我们将 api.producer 的 spanCtx 注入到 carrier 中,这里我们通过 mq 的 body 将 carrier 发送给 consumer,发送完成我们 stop 我们的 producer,那么 producer 的这层链路完成了。 随后我们来看 mq-consumer 在接收到 body 消息之后怎么做的。 type MsgBody struct { Carrier *propagation.HeaderCarrier Msg string } func (c *consumer) Consumer(ctx context.Context, data []byte) error { var msg MsgBody if err := json.Unmarshal(data, &msg); err != nil { logx.Errorf(" consumer err : %v", err) } else { logx.Infof("consumerOne Consumer , msg:%+v", msg) wireContext := otel.GetTextMapPropagator().Extract(ctx, msg.Carrier) tracer := otel.GetTracerProvider().Tracer(trace.TraceName) _, span := tracer.Start(wireContext, "mq_consumer_msg", oteltrace.WithSpanKind(oteltrace.SpanKindConsumer)) defer span.End() } return nil } consumer 接收到消息后反序列化出来 Carrier *propagation.HeaderCarrier,然后通过 otel.GetTextMapPropagator().Extract 取出来 api.producer 注入的 wireContext,在通过 tracer.Start、wireContext 创建 consumer 的 span,这样 consumer 就是 api.producer 的子 span,就形成了调用链路关系,最终我们得到的关系就是 rpc.GetUserByMobile" data-original="https://img-blog.csdnimg.cn/910d8fbfd04b4c6196b3e38479c1679b.png"> 让我们来调用一下 Logic 方法,看下 jaeger 中的链路如果与我们预想的链路一致,so happy~ 项目地址 go-zero 微服务框架:https://github.com/zeromicro/go-zero https://gitee.com/kevwan/go-zero go-zero 微服务最佳实践项目:https://github.com/Mikaelemmmm/go-zero-looklook 欢迎使用 go-zero 并 star 支持我们! 微信交流群 关注『微服务实践』公众号并点击 交流群 获取社区群二维码。

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

管家助力,轻松玩转MaxCompute

精彩视频回顾请点击:MaxCompute管家详解以下是直播内容精华整理,主要包括以下四个方面:1.背景速览;2.功能介绍;3.案例讲解;4.新功能预告。 一、背景速览 MaxCompute(原ODPS)是一项大数据计算服务,它能提供快速、完全托管的PB级数据仓库解决方案,使用户可以经济并高效的分析处理海量数据。在购买了MaxCompute之后会有相当多而繁琐的管理和维护工作,比如如何对项目进行更精细化的管理、如何将项目与配额进行关联等等,而MaxCompute管家可以帮助用户更好地完成这些工作,它是一个为用户提供作业信息查看、资源消耗查看(涵盖CU资源和存储资源)、项目查看及调整、配额组增删改查等涉及日常MaxCompute运维能力的管理平台。目前,全球包括美国、英国、德国、印度、日本、新加坡在内的18个国家或地区(详情见官网)购买

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

怎么玩转Java线程池?

一:简介 线程的使用在java中占有极其重要的地位,在jdk1.4极其之前的jdk版本中,关于线程池的使用是极其简陋的。在jdk1.5之后这一情况有了很大的改观。Jdk1.5之后加入了java.util.concurrent包,这个包中主要介绍java中线程以及线程池的使用。为我们在开发中处理线程的问题提供了非常大的帮助。 二:线程池 线程池的作用: 线程池作用就是限制系统中执行线程的数量。 根据系统的环境情况,可以自动或手动设置线程数量,达到运行的最佳效果;少了浪费了系统资源,多了造成系统拥挤效率不高。用线程池控制线程数量,其他线程排队等候。一个任务执行完毕,再从队列的中取最前面的任务开始执行。若队列中没有等待进程,线程池的这一资源处于等待。当一个新任务需要运行时,如果线程池中有等待的工作线程,就可以开始运行了;否则进入等待队列。 为什么要用线程池: 1.减少了创建和销毁线程的次数,每个工作线程都可以被重复利用,可执行多个任务。 2.可以根据系统的承受能力,调整线程池中工作线线程的数目,防止因为消耗过多的内存,而把服务器累趴下(每个线程需要大约1MB内存,线程开的越多,消耗的内存也就越大,最后死机)。 Java里面线程池的顶级接口是Executor,但是严格意义上讲Executor并不是一个线程池,而只是一个执行线程的工具。真正的线程池接口是ExecutorService。 比较重要的几个类: ExecutorService 真正的线程池接口。 ScheduledExecutorService 能和Timer/TimerTask类似,解决那些需要任务重复执行的问题。 ThreadPoolExecutor ExecutorService的默认实现。 ScheduledThreadPoolExecutor 继承ThreadPoolExecutor的ScheduledExecutorService接口实现,周期性任务调度的类实现。 要配置一个线程池是比较复杂的,尤其是对于线程池的原理不是很清楚的情况下,很有可能配置的线程池不是较优的,因此在Executors类里面提供了一些静态工厂,生成一些常用的线程池。 newSingleThreadExecutor 创建一个单线程的线程池。这个线程池只有一个线程在工作,也就是相当于单线程串行执行所有任务。如果这个唯一的线程因为异常结束,那么会有一个新的线程来替代它。此线程池保证所有任务的执行顺序按照任务的提交顺序执行。 2.newFixedThreadPool 创建固定大小的线程池。每次提交一个任务就创建一个线程,直到线程达到线程池的最大大小。线程池的大小一旦达到最大值就会保持不变,如果某个线程因为执行异常而结束,那么线程池会补充一个新线程。 newCachedThreadPool 创建一个可缓存的线程池。如果线程池的大小超过了处理任务所需要的线程, 那么就会回收部分空闲(60秒不执行任务)的线程,当任务数增加时,此线程池又可以智能的添加新线程来处理任务。此线程池不会对线程池大小做限制,线程池大小完全依赖于操作系统(或者说JVM)能够创建的最大线程大小。 4.newScheduledThreadPool 创建一个大小无限的线程池。此线程池支持定时以及周期性执行任务的需求。 实例 1:newSingleThreadExecutor MyThread.java publicclassMyThread extends Thread { @Override publicvoid run() { System.out.println(Thread.currentThread().getName() + "正在执行。。。"); } } TestSingleThreadExecutor.java publicclassTestSingleThreadExecutor { publicstaticvoid main(String[] args) { //创建一个可重用固定线程数的线程池 ExecutorService pool = Executors. newSingleThreadExecutor(); //创建实现了Runnable接口对象,Thread对象当然也实现了Runnable接口 Thread t1 = new MyThread(); Thread t2 = new MyThread(); Thread t3 = new MyThread(); Thread t4 = new MyThread(); Thread t5 = new MyThread(); //将线程放入池中进行执行 pool.execute(t1); pool.execute(t2); pool.execute(t3); pool.execute(t4); pool.execute(t5); //关闭线程池 pool.shutdown(); } } 输出结果 pool-1-thread-1正在执行。。。 pool-1-thread-1正在执行。。。 pool-1-thread-1正在执行。。。 pool-1-thread-1正在执行。。。 pool-1-thread-1正在执行。。。 2newFixedThreadPool TestFixedThreadPool.Java publicclass TestFixedThreadPool { publicstaticvoid main(String[] args) { //创建一个可重用固定线程数的线程池 ExecutorService pool = Executors.newFixedThreadPool(2); //创建实现了Runnable接口对象,Thread对象当然也实现了Runnable接口 Thread t1 = new MyThread(); Thread t2 = new MyThread(); Thread t3 = new MyThread(); Thread t4 = new MyThread(); Thread t5 = new MyThread(); //将线程放入池中进行执行 pool.execute(t1); pool.execute(t2); pool.execute(t3); pool.execute(t4); pool.execute(t5); //关闭线程池 pool.shutdown(); } } 输出结果 pool-1-thread-1正在执行。。。 pool-1-thread-2正在执行。。。 pool-1-thread-1正在执行。。。 pool-1-thread-2正在执行。。。 pool-1-thread-1正在执行。。。 3 newCachedThreadPool TestCachedThreadPool.java publicclass TestCachedThreadPool { publicstaticvoid main(String[] args) { //创建一个可重用固定线程数的线程池 ExecutorService pool = Executors.newCachedThreadPool(); //创建实现了Runnable接口对象,Thread对象当然也实现了Runnable接口 Thread t1 = new MyThread(); Thread t2 = new MyThread(); Thread t3 = new MyThread(); Thread t4 = new MyThread(); Thread t5 = new MyThread(); //将线程放入池中进行执行 pool.execute(t1); pool.execute(t2); pool.execute(t3); pool.execute(t4); pool.execute(t5); //关闭线程池 pool.shutdown(); } } 输出结果: pool-1-thread-2正在执行。。。 pool-1-thread-4正在执行。。。 pool-1-thread-3正在执行。。。 pool-1-thread-1正在执行。。。 pool-1-thread-5正在执行。。。 4newScheduledThreadPool TestScheduledThreadPoolExecutor.java publicclass TestScheduledThreadPoolExecutor { publicstaticvoid main(String[] args) { ScheduledThreadPoolExecutor exec = new ScheduledThreadPoolExecutor(1); exec.scheduleAtFixedRate(new Runnable() {//每隔一段时间就触发异常 @Override publicvoid run() { //throw new RuntimeException(); System.out.println("================"); } }, 1000, 5000, TimeUnit.MILLISECONDS); exec.scheduleAtFixedRate(new Runnable() {//每隔一段时间打印系统时间,证明两者是互不影响的 @Override publicvoid run() { System.out.println(System.nanoTime()); } }, 1000, 2000, TimeUnit.MILLISECONDS); } } 输出结果 8384644549516 8386643829034 8388643830710 8390643851383 8392643879319 8400643939383 三:ThreadPoolExecutor详解 ThreadPoolExecutor的完整构造方法的签名是:ThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueue workQueue, ThreadFactory threadFactory, RejectedExecutionHandler handler) . corePoolSize - 池中所保存的线程数,包括空闲线程。 maximumPoolSize-池中允许的最大线程数。 keepAliveTime - 当线程数大于核心时,此为终止前多余的空闲线程等待新任务的最长时间。 unit - keepAliveTime 参数的时间单位。 workQueue - 执行前用于保持任务的队列。此队列仅保持由 execute方法提交的 Runnable任务。 threadFactory - 执行程序创建新线程时使用的工厂。 handler - 由于超出线程范围和队列容量而使执行被阻塞时所使用的处理程序。 ThreadPoolExecutor是Executors类的底层实现。 在JDK帮助文档中,有如此一段话: “强烈建议程序员使用较为方便的Executors工厂方法Executors.newCachedThreadPool()(无界线程池,可以进行自动线程回收)、Executors.newFixedThreadPool(int)(固定大小线程池)Executors.newSingleThreadExecutor()(单个后台线程) 它们均为大多数使用场景预定义了设置。” 下面介绍一下几个类的源码: ExecutorService newFixedThreadPool (int nThreads):固定大小线程池。 可以看到,corePoolSize和maximumPoolSize的大小是一样的(实际上,后面会介绍,如果使用无界queue的话maximumPoolSize参数是没有意义的),keepAliveTime和unit的设值表名什么?-就是该实现不想keep alive!最后的BlockingQueue选择了LinkedBlockingQueue,该queue有一个特点,他是无界的。 public static ExecutorService newFixedThreadPool(int nThreads) { return new ThreadPoolExecutor(nThreads, nThreads, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>()); } ExecutorService newSingleThreadExecutor():单线程 public static ExecutorService newSingleThreadExecutor() { return new FinalizableDelegatedExecutorService (new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>())); } ExecutorService newCachedThreadPool():无界线程池,可以进行自动线程回收 这个实现就有意思了。首先是无界的线程池,所以我们可以发现maximumPoolSize为big big。其次BlockingQueue的选择上使用SynchronousQueue。可能对于该BlockingQueue有些陌生,简单说:该QUEUE中,每个插入操作必须等待另一个线程的对应移除操作。 public static ExecutorService newCachedThreadPool() { return new ThreadPoolExecutor(0, Integer.MAX_VALUE, 60L, TimeUnit.SECONDS, new SynchronousQueue<Runnable>()); } 先从BlockingQueue workQueue这个入参开始说起。在JDK中,其实已经说得很清楚了,一共有三种类型的queue。 所有BlockingQueue 都可用于传输和保持提交的任务。可以使用此队列与池大小进行交互: 如果运行的线程少于 corePoolSize,则 Executor始终首选添加新的线程,而不进行排队。(如果当前运行的线程小于corePoolSize,则任务根本不会存放,添加到queue中,而是直接抄家伙(thread)开始运行) 如果运行的线程等于或多于 corePoolSize,则 Executor始终首选将请求加入队列,而不添加新的线程。 如果无法将请求加入队列,则创建新的线程,除非创建此线程超出 maximumPoolSize,在这种情况下,任务将被拒绝。 queue上的三种类型。 排队有三种通用策略: 直接提交。工作队列的默认选项是 SynchronousQueue,它将任务直接提交给线程而不保持它们。在此,如果不存在可用于立即运行任务的线程,则试图把任务加入队列将失败,因此会构造一个新的线程。此策略可以避免在处理可能具有内部依赖性的请求集时出现锁。直接提交通常要求无界 maximumPoolSizes 以避免拒绝新提交的任务。当命令以超过队列所能处理的平均数连续到达时,此策略允许无界线程具有增长的可能性。 无界队列。使用无界队列(例如,不具有预定义容量的 LinkedBlockingQueue)将导致在所有corePoolSize 线程都忙时新任务在队列中等待。这样,创建的线程就不会超过 corePoolSize。(因此,maximumPoolSize的值也就无效了。)当每个任务完全独立于其他任务,即任务执行互不影响时,适合于使用无界队列;例如,在 Web页服务器中。这种排队可用于处理瞬态突发请求,当命令以超过队列所能处理的平均数连续到达时,此策略允许无界线程具有增长的可能性。 有界队列。当使用有限的 maximumPoolSizes时,有界队列(如 ArrayBlockingQueue)有助于防止资源耗尽,但是可能较难调整和控制。队列大小和最大池大小可能需要相互折衷:使用大型队列和小型池可以最大限度地降低 CPU 使用率、操作系统资源和上下文切换开销,但是可能导致人工降低吞吐量。如果任务频繁阻塞(例如,如果它们是 I/O边界),则系统可能为超过您许可的更多线程安排时间。使用小型队列通常要求较大的池大小,CPU使用率较高,但是可能遇到不可接受的调度开销,这样也会降低吞吐量。 BlockingQueue的选择。 例子一:使用直接提交策略,也即SynchronousQueue。 首先SynchronousQueue是无界的,也就是说他存数任务的能力是没有限制的,但是由于该Queue本身的特性,在某次添加元素后必须等待其他线程取走后才能继续添加。在这里不是核心线程便是新创建的线程,但是我们试想一样下,下面的场景。 我们使用一下参数构造ThreadPoolExecutor: new ThreadPoolExecutor( 2, 3, 30, TimeUnit.SECONDS, new SynchronousQueue<Runnable>(), new RecorderThreadFactory("CookieRecorderPool"), new ThreadPoolExecutor.CallerRunsPolicy()); new ThreadPoolExecutor( 2, 3, 30, TimeUnit.SECONDS, new SynchronousQueue(), new RecorderThreadFactory("CookieRecorderPool"), new ThreadPoolExecutor.CallerRunsPolicy()); 当核心线程已经有2个正在运行. 此时继续来了一个任务(A),根据前面介绍的“如果运行的线程等于或多于 corePoolSize,则Executor始终首选将请求加入队列,而不添加新的线程。”,所以A被添加到queue中。 又来了一个任务(B),且核心2个线程还没有忙完,OK,接下来首先尝试1中描述,但是由于使用的SynchronousQueue,所以一定无法加入进去。 此时便满足了上面提到的“如果无法将请求加入队列,则创建新的线程,除非创建此线程超出maximumPoolSize,在这种情况下,任务将被拒绝。”,所以必然会新建一个线程来运行这个任务。 暂时还可以,但是如果这三个任务都还没完成,连续来了两个任务,第一个添加入queue中,后一个呢?queue中无法插入,而线程数达到了maximumPoolSize,所以只好执行异常策略了。 所以在使用SynchronousQueue通常要求maximumPoolSize是无界的,这样就可以避免上述情况发生(如果希望限制就直接使用有界队列)。对于使用SynchronousQueue的作用jdk中写的很清楚:此策略可以避免在处理可能具有内部依赖性的请求集时出现锁。 什么意思?如果你的任务A1,A2有内部关联,A1需要先运行,那么先提交A1,再提交A2,当使用SynchronousQueue我们可以保证,A1必定先被执行,在A1么有被执行前,A2不可能添加入queue中。 例子二:使用无界队列策略,即LinkedBlockingQueue 这个就拿newFixedThreadPool来说,根据前文提到的规则: 如果运行的线程少于 corePoolSize,则 Executor 始终首选添加新的线程,而不进行排队。那么当任务继续增加,会发生什么呢? 如果运行的线程等于或多于 corePoolSize,则 Executor 始终首选将请求加入队列,而不添加新的线程。OK,此时任务变加入队列之中了,那什么时候才会添加新线程呢? 如果无法将请求加入队列,则创建新的线程,除非创建此线程超出 maximumPoolSize,在这种情况下,任务将被拒绝。这里就很有意思了,可能会出现无法加入队列吗?不像SynchronousQueue那样有其自身的特点,对于无界队列来说,总是可以加入的(资源耗尽,当然另当别论)。换句说,永远也不会触发产生新的线程!corePoolSize大小的线程数会一直运行,忙完当前的,就从队列中拿任务开始运行。所以要防止任务疯长,比如任务运行的实行比较长,而添加任务的速度远远超过处理任务的时间,而且还不断增加,不一会儿就爆了。 例子三:有界队列,使用ArrayBlockingQueue。 这个是最为复杂的使用,所以JDK不推荐使用也有些道理。与上面的相比,最大的特点便是可以防止资源耗尽的情况发生。 举例来说,请看如下构造方法: new ThreadPoolExecutor( 2, 4, 30, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(2), new RecorderThreadFactory("CookieRecorderPool"), new ThreadPoolExecutor.CallerRunsPolicy()); new ThreadPoolExecutor( 2, 4, 30, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(2), new RecorderThreadFactory("CookieRecorderPool"), new ThreadPoolExecutor.CallerRunsPolicy()); 假设,所有的任务都永远无法执行完。 对于首先来的A,B来说直接运行,接下来,如果来了C,D,他们会被放到queue中,如果接下来再来E,F,则增加线程运行E,F。但是如果再来任务,队列无法再接受了,线程数也到达最大的限制了,所以就会使用拒绝策略来处理。 keepAliveTime jdk中的解释是:当线程数大于核心时,此为终止前多余的空闲线程等待新任务的最长时间。 有点拗口,其实这个不难理解,在使用了“池”的应用中,大多都有类似的参数需要配置。比如数据库连接池,DBCP中的maxIdle,minIdle参数。 什么意思?接着上面的解释,后来向老板派来的工人始终是“借来的”,俗话说“有借就有还”,但这里的问题就是什么时候还了,如果借来的工人刚完成一个任务就还回去,后来发现任务还有,那岂不是又要去借?这一来一往,老板肯定头也大死了。 合理的策略:既然借了,那就多借一会儿。直到“某一段”时间后,发现再也用不到这些工人时,便可以还回去了。这里的某一段时间便是keepAliveTime的含义,TimeUnit为keepAliveTime值的度量。 RejectedExecutionHandler 另一种情况便是,即使向老板借了工人,但是任务还是继续过来,还是忙不过来,这时整个队伍只好拒绝接受了。 RejectedExecutionHandler接口提供了对于拒绝任务的处理的自定方法的机会。在ThreadPoolExecutor中已经默认包含了4中策略,因为源码非常简单,这里直接贴出来。 CallerRunsPolicy:线程调用运行该任务的 execute 本身。此策略提供简单的反馈控制机制,能够减缓新任务的提交速度。 public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { if (!e.isShutdown()) { r.run(); } } public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { if (!e.isShutdown()) { r.run(); } } 这个策略显然不想放弃执行任务。但是由于池中已经没有任何资源了,那么就直接使用调用该execute的线程本身来执行。 AbortPolicy:处理程序遭到拒绝将抛出运行时RejectedExecutionException public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { throw new RejectedExecutionException(); } public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { throw new RejectedExecutionException(); } 这种策略直接抛出异常,丢弃任务。 DiscardPolicy:不能执行的任务将被删除 public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { } public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { } 这种策略和AbortPolicy几乎一样,也是丢弃任务,只不过他不抛出异常。 DiscardOldestPolicy:如果执行程序尚未关闭,则位于工作队列头部的任务将被删除,然后重试执行程序(如果再次失败,则重复此过程) public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { if (!e.isShutdown()) { e.getQueue().poll(); e.execute(r); } } public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { if (!e.isShutdown()) { e.getQueue().poll(); e.execute(r); } } 该策略就稍微复杂一些,在pool没有关闭的前提下首先丢掉缓存在队列中的最早的任务,然后重新尝试运行该任务。这个策略需要适当小心。 设想:如果其他线程都还在运行,那么新来任务踢掉旧任务,缓存在queue中,再来一个任务又会踢掉queue中最老任务。 总结: keepAliveTime和maximumPoolSize及BlockingQueue的类型均有关系。如果BlockingQueue是无界的,那么永远不会触发maximumPoolSize,自然keepAliveTime也就没有了意义。 反之,如果核心数较小,有界BlockingQueue数值又较小,同时keepAliveTime又设的很小,如果任务频繁,那么系统就会频繁的申请回收线程。 public static ExecutorService newFixedThreadPool(int nThreads) { return new ThreadPoolExecutor(nThreads, nThreads, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>()); } 文章来源 https://yq.aliyun.com/articles/94989?utm_campaign=wenzhang&utm_medium=article&utm_source=QQ-qun&2017622&utm_content=m_23936

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

玩转树莓派——安装系统

纠结了很久,终于在今年生日的时候买了个树莓派 3。拿到以后少不了各种折腾,慢慢的把折腾过程写下来吧。 有关RaspBerry Pi 3,技术参数以及对应提升,就没必要在这里赘述了,官网介绍应有尽有:https://www.raspberrypi.org/ 万能的淘宝买回一个电路板~没有选择所谓的套餐,那背离了我折腾的初衷。当然,考虑到折腾的程度,散热片和漂亮的框还是要搞一个的,电源可以用iPad的,全部到手也不到300元。 通电之前的第一件事是准备系统。有两种不同的方式来“安装”系统:直接刷tf卡,或者使用NOOBS来下载安装。 直接刷tf卡没啥可说的,Linux/MAC类用dd命令,Windows可以用一个镜像写入程序。可参考:https://www.raspberrypi.org/documentation/installation/installing-images/README.md 需要说明的是,这个镜像工具不光能写,也能将你做好的tf卡读成一个镜像文件。不过,不像ghost那样会跳过空扇区,所以会变得巨大…… 上述方式适合各种树莓派以及第三方的镜像,比如游戏机之类的,以后玩了再写。 如果不想刷卡,或者网速很快,可以选择NOOBS方式来下载安装系统。 NOOBS是树莓派的引导安装系统,可以让树莓派引导到Recovery环境,然后选择安装/下载安装需要的操作系统。 默认的版本会包含最近的树莓派操作系统——RaspBian,所以稍大。这个版本的镜像在安装的时候就直接从本地安装RaspBian了。如果你只打算安装别的系统,或者选择在线下载安装RaspBian,可以直接下载NOOBS Lite版本。 下载完成后,格式化tf卡,做成启动分区,然后将压缩包解压到tf卡即可。 将tf卡插入树莓派,上电启动,即可开始安装操作系统。 作为一个Microsoft MVP,自然对微软的系统十分有兴趣,包括运行于树莓派的Windows 10 IoT Core。树莓派系统可以直接在NOOBS支持下,下载安装Windows 10 IoT。 本文转自HaoHu 51CTO博客,原文链接:http://blog.51cto.com/haohu/1853397,如需转载请自行联系原作者

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

玩转大数据的建议

大数据成为今年李克强总理政府工作报告中提升为“互联网+”行动计划的一部分,并由此推动移动互联网、云计算、物联网等与现代制造业结合,促进电子商务、工业互联网和互联网金融健康发展,引导互联网企业拓展国际市场。可以说,大数据的存储、挖掘、分析能力已经彻底改变了传统互联网的局面。 通过大量的数据应用,我们完全可以总结出当今大数据时代如下一些数据特征: 一是数据的原始性,社会需要的是原始数据;二是数据的完整性,只要是非保密的,政府应该把掌握的数据提供给社会;三是数据的标准性,需要一个国家级标准,每个部门提供的数据,公众可以同一个标准使用;四是数据的及时性,数据拿到以后第一时间提供给社会,让创新更加及时;五是数据的可获取性,部门不能在数据的提供方面设置障碍,甚至收费赚钱。 因此,对大数据要有一个国家层面的认识,特别是应发挥政府在大数据时代的作用。既然这样,就迫切需要建立大数据资源的共建共享机制。如何做到?首先应结合目前国家的简政放权可喜形势,依托政务信息资源共享平台,把行政部门的审批职能和政务服务集中起来,提供一站式网上办理,形成横向到各委办局,纵向到各省市县区的网上办事系统。目前,国、省、市、县四级平台可实现互联互通,审批和服务事项基本实现应尽必尽。全国各省网上办事大厅系统均实现零的突破,今后逐步实现各项办事目录的数据同步和办事过程数据的交换对接,形成了网上办事全流程数据的共建共享机制。 这种机制的共建共享,首先要以需求为导向,部门的自身需求才是共建共享的真正动力;其次要加强流程再造,实行网上办事;还要开展全程督查,在推进部门信息资源共建共享过程中建立监督检查机制。 最后,对我国建立数据资源部门共建共享机制提出如下几点建议:一是尽快成立国家大数据管理局那样的机构;二是对建立大数据资源部门共建共享机制尽快立法,同时尽快出台国家层面的移动互联网信息安全法律,还要依法对数据进行挖掘和分析,同时完善安全体系;三是出台相关数据开放标准规范,指导各省开展数据资源开放工作。 本文作者:佚名 来源:51CTO

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

Zadig 快速体验,玩转本地安装!

啥是本地安装? 听过 Zadig 支持All in One 安装、基于 Kubernetes 安装、基于 Helm 安装等多种安装方式,怎么又来了个本地安装?(这么多安装方式谁听起来头不大) 别急,先看一段对话: (本对话内容基于真实场景模拟,如有雷同,实属巧合) 目标人群 需要在本地快速体验上手,不需要数据持久化保存、不需要生产环境使用。 开源项目好奇宝宝: 开源的云原生持续交付平台?赶紧让我下载安装看看你有什么花活,安装太麻烦我就放弃了。 云原生开发工程师: 虽然你有 Helm 安装,可是做云原生工程师,日常修改 YAML 已经吐了,实在不想看那么多安装的参数,我只想在自己的电脑快速体验下,我有 Docker for desktop 可以启动集群,能不能来个一键脚本给个痛快? 技术 Leader: 团队最近在做云原生持续交付平台的调研选型,听说 Zadig 很强啊,在圈内也很火,赶紧让哪个谁谁谁快速安装了解下,是否适合团队使用,调研后再决定是否上生产。 实施工程师(技术 Leader の 小弟): Leader 需要这边安装体验下 Zadig,尽快出个报告,可是申请集群资源好麻烦啊,层层审批估计疫情到时候都结束了。好烦啊,自己电脑搞个虚拟机弄个集群或者安装个 minikube 不知道能不能行? 实施工程师(技术 Leader の 小弟): Leader 需要这边安装体验下 Zadig,尽快出个报告,可是申请集群资源好麻烦啊,层层审批估计疫情到时候都结束了。好烦啊,自己电脑搞个虚拟机弄个集群或者安装个 minikube 不知道能不能行? 没关系,以上统统安排! Zadig 作为一款开源云原生持续交付产品,支持多种安装方式,每种安装方式又适用于不同的使用场景,例如基于 Helm 命令的安装方式,适用于生产使用,而且对集群资源有一定要求。而对于资源无法满足要求但又对 Zadig 感兴趣的大量开发者来说,如何实现快速体验?成为我们团队需要关注和解决的一个问题。 于是我们推出了本地安装,帮助新人在本机尝鲜和快速体验 Zadig。 如何进行本地安装? 前提 使用 minikube、KinD 等工具在本地拉起一套 K8s 集群,参考如下: a.安装 minikube[1] b.安装 docker-desktop[2] c. 更多工具请参考其官方安装文档 确保本地 K8s 集群满足至少 4C8G 的资源配置,版本满足 v1.16~v1.22。 第一步:安装 Zadig 在本地集群中执行以下脚本: 若安装成功后需要集成外部系统(比如:代码源),请确保使用的IP 地址可外网访问。 1 export IP=<本机 IP 地址>2 export PORT=<任意一合法的 K8s Node Port>3 curl -SsL https://download.koderover.com/install?type=quickstart | bash 安装成功后系统会自动初始化登录账号和密码。 第二步:访问 Zadig 小贴士: 如果使用的是 KinD 拉起的集群,由于其自身特性,需要打通本机端口到 K8s 集群 NodePort 服务的通路,参考命令如下: 1 kubectl -n zadig port-forward svc/gateway-proxy 32000:80 访问 IP: PORT,使用默认账号密码 admin/zadig 登录成功后,即可愉快玩耍了~ One More Thing 我们计划在后续的更新中,支持内置的 demo 项目,在本地安装成功后即可直接体验工作流、环境、服务部署等功能,缩短从配置到使用的路径,做到开箱即用,降低体验 Zadig 的门槛。 参考链接 [1]https://minikube.sigs.k8s.io/docs/start/ [2]https://www.docker.com/products/docker-desktop/ Zadig,让工程师更专注创造!

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

玩转 HelloGitHub 的新姿势

本文不会涉及太多技术细节和源码,请放心食用 大家好,我是 HelloGitHub 的老荀,好久不见啊! 我在完成 HelloZooKeeper 系列之后,就很少“露面了”。但是我对开源和 HelloGitHub 的热情并没有丝毫的减少。这不,逮着个机会就来输出一波,防止被大家遗忘😂。 这次带来的是我写的一款在终端浏览 HelloGitHub 的工具:hg-tui,让你双手不离开键盘就能畅游在 HG 的开源世界。功能如下: 色彩丰富、平铺展示 关键字搜索月刊往期的项目 类 Vim 的快捷键操作方式 一键直达开源项目首页 支持 Linux、macOS、Windows 地址:https://github.com/kaixinbaba/hg-tui 下面我将分享自己发起这个开源项目的缘起、构思、再到开发的全部过程,最后分享一下,我通过做这个项目对开源的一些感悟。 一、起因 我本职是做 Java 开发,但架不住 Rust 太有意思了!所以最近在学 Rust 恰好前段时间看到 HG 讲解 tui.rs 的文章。 看完后手痒得厉害,就写了一篇 tui.rs 入门文章,但感觉还不过瘾就想写一个项目练手。 因为我平时经常上 HelloGitHub 找开源项目,所以就决定用 tui.rs 做一个终端浏览 HelloGitHub 官网的工具。 官网:https://hellogithub.com/ 二、构思 首先我希望这个应用能有以下功能: 有搜索框,可以按关键词搜索 HelloGitHub 中的任意项目 通过表格按列展示搜索结果 既然是终端应用,那操作方式肯定是使用键盘方式,快捷键我采用了一些大家熟知的 Vim 快捷键 浏览项目的途中,可以随时在浏览器中打开当前浏览的项目 有了这些主要功能点的思路,下面就要想想怎么设计一个界面了,我本职工作后端一碰到画界面就头疼,几经周折大概把界面设计成了这样: 又因为是 TUI 界面层级不能太深,所以再多弄个详情页面(用来浏览文字明细)或者弹窗页面(提示消息)就差不多了。 我又想到了 GitHub 为每一种编程语言都设计了一种颜色,我也可以把这些颜色应用在我的项目里,让整个终端界面看起来没那么单调,色彩更丰富。效果如下: 主界面: 详情页: 弹窗提示: 最后为了向 TUI 妥协,按期数或类别搜索,我是通过使用搜索前缀来和普通关键词搜索作出区别。 上面展示的这些差不多已经是这个项目的全部了 三、开发 3.1 技术选型 要实现上述的那些功能,就要从 Rust 的生态中选择合适的库了 下面这些是我在这个项目中使用到的: 基础设施:anyhow、thiserror、lazy_static、better-panic 绘制 UI:tui、crossterm HTTP client:reqwest 缓存:cached HTML 解析:nipper 工具:regex、crossbeam-channel 命令行:clap 虽然 Rust 还是编程界的小学生(2011 年启动),但是经过了这些年的发展,生态已经逐渐完善,工具库已经很丰富了。再加上 Rust 是系统级的语言,值得投入时间学习! 3.2 项目结构 项目结构规划(非全部) src ├── app.rs // 统一管理整个应用的状态 ├── cli.rs // 命令行解析 ├── draw.rs // 绘制 UI ├── events.rs // UI 事件、输入事件、通知 ├── fetch.rs // HTTP 请求 ├── main.rs // 入口 ├── parse.rs // HTML 解析 ├── utils.rs // 工具 └── widget // 自定义组件 ├── ... 合理的分文件(目录)开发,可以让每个功能模块 高内聚、低耦合,并且可以很容易地分开进行单元测试。 当然这些文件也不是在项目之初就已经一股脑地建立好的,都是在完善功能的路上一点点添加进来的~ 3.3 主要代码 因为是基于 tui.rs 开发的应用,所以主流程肯定是遵循该库的设计的,首先需要定义一个 App 用来保存整个项目的状态信息。 pub struct App { /// 用户输入框 pub input: InputState, /// 内容展示 pub content: ContentState, /// 弹窗提示 pub popup: PopupState, /// 状态栏 pub statusline: StatusLineState, /// 模式 pub mode: AppMode, /// 项目明细子页面 pub project_detail: ProjectDetailState, ... } 每一个状态字段,其实就是对应一个自定义组件.要在 tui.rs 中实现自定义组件(实现方式也是我自己的理解)也很简单只要三步,我以 Input 组件为例。 /// 用户输入框组件,组件本身没有字段,是一个无状态的对象 /// 无状态对象只关心 UI 怎么绘制,不存储数据 pub struct Input {} /// 组件的状态,每一个字段就是组件需要存储的数据 #[derive(Debug)] pub struct InputState { input: String, active: bool, pub mode: SearchMode, } /// 最后为 Input 组件实现 StatefulWidget trait impl StatefulWidget for Input { type State = InputState; // 指定关联类型为 InputState /// area 绘制的区域 /// buf 缓冲区(可以直接写入字符串,如果要高度定制的话,可以理解为画笔) /// state 从这个变量中直接取绘制过程中需要的数据 fn render(self, area: Rect, buf: &amp;mut Buffer, state: &amp;mut Self::State) { // 具体绘制的逻辑 ... } } 只要是面向用户的应用,都会处理各种各样的用户输入(事件)。Rust 中一般都使用 channel 来解耦处理各种各样的事件,再利用 Rust 强大的枚举支持,定义各种各样的事件(用户输入和非用户输入)即可。 /// 定义事件枚举 #[derive(Debug, Clone)] pub enum HGEvent { /// 用户事件(键盘事件) UserEvent(KeyEvent), /// 应用内部组件的通知事件 NotifyEvent(Notify), } #[derive(Debug, Clone, PartialEq)] pub enum Notify { /// 重绘界面 Redraw, /// 退出应用 Quit, /// 弹出窗口展示消息 Message(Message), /// tick,比如一些数据需要每隔一段时间自动更新的(比如:显示的时间) Tick, } /// 弹窗的消息,分为 错误、警告、提示 #[derive(Debug, Clone, PartialEq)] pub enum Message { Error(String), Warn(String), Tips(String), } 为了区分用户事件和通知,我使用了两个不同的 channel 分别处理这两类: lazy_static! { /// 因为通知队列希望被应用内部共享,所以使用了 lazy_static 方便使用 pub static ref NOTIFY: (Sender<hgevent>, Receiver<hgevent>) = bounded(1024); } 又因为不同的事件处理,并不应该互相阻塞,所以整个应用采用了最基础的多线程模型来提高性能,这里使用的也是标准库的多线程。 pub fn handle_key_event(event_app: Arc<mutex<app>&gt;) { let (sender, receiver) = unbounded(); ... std::thread::spawn(move || loop { // 单独一个线程接收用户事件 if let Ok(Event::Key(event)) = crossterm::event::read() { sender.send(HGEvent::UserEvent(event)).unwrap(); } }); std::thread::spawn(move || loop { // 单独一个线程处理用户事件 if let Ok(HGEvent::UserEvent(key_event)) = receiver.recv() { ... } }); } 其他剩下的就是业务逻辑,完整的代码可以直接看仓库 https://github.com/kaixinbaba/hg-tui 四、心路历程 一开始我做 hg-tui 项目的时候,仅仅是为了做个实际的项目把玩一下 tui.rs 这个框架,做好之后问题层出不穷,但我深知没有与生俱来的完美,只有不断的迭代才能让它越来越好,经过 100 多次的提交后,现在用着感觉顺手多了。毕竟作者是项目的第一个用户,自己用着不舒服其他人就更不喜欢了! 我想着既然要让别人用,一定要容易安装。接着我做了基于 GitHub Action 自动编译和发布,支持 Windows、Linux、macOS 直接下载就能用。 我还做了对 homebrew 安装的支持,但因为 Star 数不够没有收录到 homecore 要求:30 forks、30 watchers、75 stars 希望大家看到这里的话能给个 star✨ 地址:https://github.com/kaixinbaba/hg-tui 五、最后 hg-tui 它从出生那一刻起,体内流淌的就是开源的血。 它很小甚至是微不足道,我本不想开源,但蛋蛋的一段话让我改变了主意:开源不是完结,仅仅只是开始。 一个开源项目可能只是作者的一个灵光乍现,也可能只是为了解决自己实际工作生活中的小小痛点,没准用完就丢到角落里了。但开源出来或许就能找到有相同需求的人,从而延续这个项目的生命,或许这就是开源的本意吧。 以上就是我做这个项目的全部心得和收获,如果你们对 hg-tui 有什么建议和问题,欢迎给我提 issue 最后,如果你喜欢本文和项目的话,欢迎点赞和 Star 爱你们哟~

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

WebStorm

WebStorm

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

用户登录
用户注册