高并发
快来分享你的内容吧~
- 04-05 14:33
- 如何用java去实现压测不同的qps值,以及提高qps值?问题描述如何用java去实现压测不同的qps值,多线程+for-loop的方式控制不了发送请求的时间吧?以及有哪些措施提高qps值?背景信息开发一个http server项目,并在里面实现一个高延迟API。该API要求如下:- 输入:圆的半径(浮点数)- 输出:圆的面积(浮点数)- API延迟要求:模拟业务处理耗时50-100ms - 当输出结果正确且延迟在此范围内时,才算成功。 - 延迟...查看全文程序员鱼皮:AI 的回答可以作为参考,但是没必要自己写这个压测程序,一般用 jmeter 就可以搞定。可以通过配置设置线程预启动,然后同时发送请求。但是一般也没必要,我们 qps 估算的一般都是平均值,比如 1000 个线程每秒都发一次请求,等个 100 秒,测试结果就很稳定了,可以避免刚开始启动时的性能损耗。可以看下代码生成器共享平台项目的性能优化章节。

- 2024-03-10·后端
我想请问一个问题,线程池源码中我想请问一个问题,线程池源码中的ThreadPoolExecutor中用了一个原子整形表示了线程池状态和线程池数量为了为了使一次原子操作完成这两个变量的同时修改但是这里面为什么jdk源码用的是AutomicInteger来生命的,而不是用性能更好的LongAdder呢?...查看全文程序员鱼皮:下面是比较专业的回答,重点还是在应用场景上,未必性能更高的技术一定就要用。在 Java 7 中,JDK 已经对线程池进行了优化。其中一个优化是使用 LongAdder 替换了原来的 AtomicLong 类,这是因为在高并发下,AtomicLong 会出现争用,降低并发性能。而 LongAdder 采用了分段累加的思想,多个线程累加时,不会互相争用,而是分别在各自的段中累加,最终再把各个段的值相加
- 2022-12-31·Java
- 2022-12-26
- 2022-10-29JDK1.6 以后对 Synchronized 做了哪些优化?查看全文编程导航:DK1.6 之前是重量级锁,他加锁底层是通过系统的mutex相关指令实现,会有用户态和内核态之间的切换,十分消耗性能。 JDK1.6之后对synchronized底层做了优化,引入了偏向锁、轻量级锁、在JVM层面实现加锁的逻辑,不依赖底层操作系统,就没有状态切换的消耗。同时引入自旋锁、适应性自旋锁、锁消除、锁粗化等技术来减少锁操作的开销。 所以在对象头的Mark word 锁主要存在四种状态依次是
18112分享
鱼皮哥,想问一下如何提高单体并鱼皮哥,想问一下如何提高单体并发量,是用一些并发包或者线程池或者用消息队列去削峰吗...查看全文程序员鱼皮:提高单机并发量有很多种方式,核心的问题是【分析影响你当前并发量的因素】到底是什么?对症下药而不是盲目或者纯凭理论(经验值)去猜测。常用的性能优化方法有:并发编程、增加机器(扩容)、负载均衡、机器升配、修改配置、读写分离、缓存、优化代码、优化技术选型、优化业务流程等。举一些例子:1. 单机 Tomcat 服务器处理请求的性能有限,可以修改 tomcat 参数(比如最大并发线程数)、可以提升 tomc
📚 学习实战 | 亿级流量高并发点赞系统(Spring Boot3 + Java21 进阶项目落地)
> 最近跟着亿级流量高并发点赞教程系统学习了分布式架构设计,总觉得光看理论不动手实战等于白学,索性就以抢票这个经典的高并发业务场景为目标,把教程里讲到的高可用设计、全链路限流熔断、异步解耦、可观测监控这些核心知识点,完整落地成了一套可直接运行、可压测调优、符合生产规范的微服务抢票系统。 > 从0到1拆分微服务、写核心抢票逻辑、做全链路稳定性防护、压测调优,前前后后打磨了一个多月,把Java 21虚拟线程、Redis Lua原子操作、Kafka消息可靠性保障、Sentinel熔断降级这些秋招面试高频考点都踩了一遍坑,现在把完整项目开源出来,给同样在学后端进阶、准备秋招的同学做个实战参考,也欢迎大家一起交流优化~ ## 💡 项目核心技术栈 完全贴合当下企业主流技术栈,紧跟Java技术前沿: | 核心组件 | 技术选型 | | :--- | :--- | | 基础框架 | Spring Boot 3.2.0 | | JDK版本 | Java 21(虚拟线程) | | 微服务体系 | Spring Cloud Gateway + Nacos + OpenFeign | | 缓存设计 | Caffeine + Redis 7 多级缓存 | | 消息队列 | Kafka 7.5.0 | | 数据库 | MySQL 8.0 + MyBatis-Plus | | 限流熔断 | Alibaba Sentinel | | 定时任务 | XXL-Job | | 监控体系 | Prometheus + Grafana | | 部署方案 | Docker + Kubernetes | | 性能压测 | k6 | ## ✨ 核心落地亮点(教程知识点实战+场景化优化) ### 1. 高并发抢票核心流程:Lua + Kafka + DB 最终一致性 针对抢票最核心的超卖问题,用Redis Lua脚本实现了原子性的库存扣减+限购校验,单脚本内完成所有前置校验和库存操作,避免并发问题;再通过Kafka异步解耦,把订单创建和DB库存扣减异步化,大幅提升抢票接口的吞吐量,同时保证数据最终一致性。 ### 2. 全链路熔断降级与限流防护 跟着教程里的高可用设计思路,实现了4层防护,彻底避免单点故障拖垮整个系统: - API层:基于Sentinel实现用户维度的接口限流,自定义注解开箱即用 - 服务间调用:OpenFeign全量配置降级工厂,服务异常自动降级返回 - 缓存层:Redis所有操作都做了熔断保护,避免缓存故障雪崩 - 消息层:Kafka发送失败实现两层降级(Redis Stream持久化 + 本地内存队列兜底),搭配XXL-Job定时重试,保证消息不丢失 ### 3. 性能优化细节落地 - 多级缓存设计:Caffeine本地缓存 + Redis分布式缓存,网关实现Sticky Session一致性哈希路由,大幅提升本地缓存命中率 - 全面拥抱Java 21虚拟线程,开启Spring Boot虚拟线程支持,大幅提升系统在高并发下的吞吐量 - 缓存失效通过Redis Pub/Sub广播通知,解决多实例本地缓存数据一致性问题 ### 4. 完整的工程化与生产级配套 - 提供一键启动/停止脚本,Docker compose一键拉起所有基础设施,开箱即用 - 完整的Kubernetes生产环境部署配置,包含命名空间、中间件、微服务、HPA自动扩缩容全套配置 - 配套k6性能压测脚本,覆盖登录、查询、抢票全场景,可直接执行压测验证性能 - 集成Prometheus + Grafana全栈监控,配套预定义的监控面板,实时观测系统运行状态 ## ⭐️ 完整源码地址 项目所有代码、配置、脚本都已完整开源,包含详细的文档说明,欢迎 Star 交流,一起优化完善: https://github.com/zunff/ticket-booking-backend 顺手把前端vibe coding了出来 https://github.com/zunff/ticket-booking-frontend
高并发下的超卖问题
### 超卖问题的处理需多角度进行流量的控制 #### 前端(防止用户的操作发出大量请求) 首先是页面的按钮防抖,防止用户发出大量请求 (有用的请求并没有那么多,如果可以,也可以直接在前端用随机的方式进行处理(直接随机给用户一个你没抢到.....)) 利用cdn缓存静态资源(减轻服务器的压力),手动推给cdn进行预热 这样大量的请求就会打到cdn而不是所有地区的请求都到服务器 #### 后端 #### 瞬时流量的承接(前端也吸收了一部分的流量) 1,负载均衡 nginx(这个东西一般来说是做为前端的静态资源代理器,当然在这个地方也可以进行单个集群的负载均衡以及限流+黑名单) 设置的内容: #### DB--防止超卖 ###### 乐观锁 乐观锁和悲观锁的内容不在这边进行叙述 #### DB--库存分组 将总库存也做集群,分为小库存,提高读写效率 #### DB--插入库存扣减流水 一直update行会出现问题,将update改为insert。插入数据的量不好控制会导致超卖 #### DB--热点行问题 数据库的热点行问题是指某些数据行被频繁访问或更新,导致这些行成为系统性能的瓶颈。(锁竞争) 解决方法之前的库存分组 #### db-redis缓存 用redis+lua的方式进行控制,先将大量写入数据放置到redis中进行处理,然后将redis的数据异步刷新到mysql中(在这里异步刷新的作用还有削峰填谷的作用)用来实现最终一致性, ### 除了这个db方案,还需要一个准时对账机制 lua脚本不单单需要进行库存的削减,还需要将流水zset进流水对列。定时拉取一段时间的流水和数据库的库存比较是否一致,不是则进行缓存。 #### 如果发生了不一致 我们更倾向于认为数据库的数据可靠()根据数据库进行补偿 redis中的流水数据可能因为缓存失效和数据丢失而导致不一致 #### 预防黑产 将行为异常的用户进行拉黑 #### 幂等性 无论执行多少次,结果不变 ### 兜底方案 直接关闭秒杀服务,秒杀服务的失败对于用户来说是可以接受的,直接止损
如何用java去实现压测不同的qps值,以及提高qps值?
### 问题描述 如何用java去实现压测不同的qps值,多线程+for-loop的方式控制不了发送请求的时间吧?以及有哪些措施提高qps值? ### 背景信息 开发一个http server项目,并在里面实现一个高延迟API。 该API要求如下: - 输入:圆的半径(浮点数) - 输出:圆的面积(浮点数) - API延迟要求:模拟业务处理耗时50-100ms - 当输出结果正确且延迟在此范围内时,才算成功。 - 延迟大于100ms则视为超时,算失败。 使用java写一个http客户端程序,用来调用http api,同时具备功能验证和性能压测需求。 - 功能验证:随机生成1万个数字,调用api,验证功能正确性 - 串行调用api,用txt文本输出每个请求的输入,输出,延迟,并标出结果有错误的case,计算出成功率 - 性能压测:客户端可以采用多线程+for-loop方式来压测api的性能。 - 统计不同压测qps值(100,1000,10000,...)时,统计api返回成功率(成功率=成功返回请求数/总请求数 * 100%) - 检验当可用性>99.9%时,api支持的最高qps是多少? - 性能优化: - 本题目中,影响压测qps上限的因素是哪些?有哪些解决方法? - 通过一定的优化措施,来不断提升最高qps上限。 ### 具体疑问 主要卡在了如何用java去实现压测不同的qps值,尝试了多线程+for-loop的方式,但是无法保证发送的总时间,以及api返回成功率中,成功返回请求数是指不超过100ms的返回,还是其他意思? ### 已有理解 目前查阅资料知道qps值可以通过设置线程池、配置tomcat、消息队列等提高 ### 预期目标 希望得到一个用java去实现压测不同的qps的示例,以及具体的提高qps的方法? ### 相关资料 https://codecopy.cn/post/w3lxt4?pw=f8Xzc5
一文了解进程和线程
# 1.前言 在计算机中,CPU是最核心的一个硬件资源,相当于我们人类的大脑一样,用来处理所有复杂的计算任务。而在运行程序时,我们处理某些运行时的数据,需要临时保存,这时候就需要用到内存来存储这些。但是存在内存中的数据毕竟是临时的,这时候我们需要把一些关键性的数据存储到外存(磁盘)中,这样以后可以随时查看了,而这些计算任务的调度,资源的分配,是由操作系统来统领的,我们打开的每一个应用程序,都是以进程的形式,运行于操作系统之上。 # 2.提升程序的运行速度 从前言中我们了解到一个程序简单的调度过程,我们发现,影响一个程序的速度主要有这三大方面: 1. CPU(计算能力强不强) 2. 内存(足够大的内存能够存储临时计算数据) 3. IO (能快速把数据持久化到磁盘上) 但是这三者之间的处理速度上,却有着十分大的差异,总体来说 CPU > 内存 >> IO,我们可以看出,IO的处理速度是最慢的,这也是影响一段程序速度的最主要的一块,在无数代先辈的努力下,这三者都在不断的发展: 1. CPU增加了高速缓存,平衡了和内存之间的速度差异,从早期的单核升级到如今的多核,极大提升了CPU的计算能力 2. 内存容量不断扩大,能够让计算机运行更多的程序 3. 操作系统从进程中,又诞生了更轻量级的线程,(乃止后面的超线程)为了能够分时复用CPU,平衡了CPU与IO之间的速度差异,硬件上存储设备诞生了固态硬盘等,让IO的性能有了明显的提升。 在这里我们主要来介绍进程到线程这段的发展 # 3. 从进程到线程 由于早期的CPU是单核的,但是用户却也能运行多个程序,这是怎么实现的?原来是因为早期的操作系统以“多进程”的形式运行程序。操作系统让每个进程占用一段时间CPU的使用权,当这段时间消耗完后,操作系统会重新选择一个进程,让它获取CPU的执行权,而这个切换过程通常为毫秒级别,对于用户来说完全感知不到这个任务切换,从而实现了同时运行多个程序。进程占用CPU处理任务的这段时间,我们称之为时间片。 但当某个进程在占用CPU时间片内,需要进行耗时很长的IO操作,为了提高CPU的利用率,这时候进程会让出CPU的时间片,让其他进程获取CPU时间片来执行任务,等自己完成IO操作,将数据读取到内存后,就可以重新获取CPU的时间片。 每个进程都有自己独立的内存空间,多个进程之间不共享彼此的数据,但是多个进程之间的任务切换存在比较大的开销,为了进一步的提高并发性能和CPU的利用率,进程内部诞生了线程,线程是进程内的一个执行单元,一个进程可以包含多个线程。相比于进程,线程之间的切换和通信成本更低。线程共享进程的内存空间和系统资源,这使得线程之间的数据交换更加容易。 线程是指“进程代码段”的一次顺序执行流程。线程是CPU调度的最小单位。一个进程可以有一个或多个线程,各个线程之间共享进程的内存空间、系统资源,进程仍然是操作系统资源分配的最小单位。 对于Java来说,Java代码都运行在JVM(Java虚拟机)之中,每当运行一个Java程序,就会启动一个JVM进程,在JVM内部,所有代码以线程来运行,在这个JVM进程中,起码会存在两个线程,一个是main线程,而另一个是GC垃圾回收线程。当程序运行完后,JVM进程也就结束了。 # 4. 进程与线程的区别 一个进程由一个或多个线程组成线程是CPU调度的最小单位,进程是操作系统分配资源的最小单位。线程的划分尺度小于进程。 创建和终止进程的开销通常较大,因为操作系统需要为进程分配独立的内存空间和系统资源。相比之下,创建和终止线程的开销较小,因为线程共享进程的资源。 操作系统在进行任务调度时,会为每个进程分配时间片。而线程作为进程内的执行单元,由进程来进行调度。在多核处理器系统中,多线程技术可以实现真正的并行执行,进一步提高系统性能。 #知识碎片
我想请问一个问题,线程池源码中
我想请问一个问题,线程池源码中的 ThreadPoolExecutor中 用了一个原子整形表示了线程池状态和线程池数量为了为了使一次原子操作完成这两个变量的同时修改 但是这里面为什么jdk源码用的是AutomicInteger来生命的,而不是用性能更好的LongAdder 呢?
【并发编程】自定义简单线程池
## 优质博文 [更好的使用 JAVA 线程池](https://my.oschina.net/andylucc/blog/648127) [深入理解Java线程池:ThreadPoolExecutor](http://www.ideabuffer.cn/2017/04/04/%E6%B7%B1%E5%85%A5%E7%90%86%E8%A7%A3Java%E7%BA%BF%E7%A8%8B%E6%B1%A0%EF%BC%9AThreadPoolExecutor/) [Java线程池实现原理及其在美团业务中的实践](https://tech.meituan.com/2020/04/02/java-pooling-pratice-in-meituan.html) ## 1、概念图 核心部分: - 阻塞队列`BlockingQueue`:暂存线程池中无法处理的任务 - 线程池`ThreadPool`:自定义的线程池,内部最多包含`coreSize`个工作线程执行任务 - 工作线程`WorkerThread`:执行传递过来的任务 - 拒绝策略`Rejectpolicy`:当阻塞队列已满时采用指定的策略拒绝任务  ## 2、流程分析 根据上面的概念图,进一步模拟一遍整个线程池执行的流程: 1. **初始化线程池**,指定线程池的参数如**核心线程数、阻塞队列容量、超时时间、拒绝策略**; 2. 并发生产**任务压入线程池**执行; 1. 工作线程数**未达到**设定的核心线程数。**新建工作线程**执行任务,并将工作线程加入到线程池中的**线程集合**中; 2. 工作线程数**达到了**设定的核心线程数。**尝试往阻塞队列中暂存任务**,当阻塞队列**已满**无法添加时,采用指定的**拒绝策略**对任务进行拒绝。 3. 工作线程执行完当前任务时,**循环从阻塞队列中获取任务**并执行直到消费完阻塞队列中的任务; 4. 当无任务时,将工作**线程回收**销毁。 ## 3、设计思路及实现 整体的设计思路应该由广到细,整体到局部。前面的概念图以及流程分析其实就算是一个整体的设计了,接下来便是局部的设计了。首先先列举一下需要的部分,分别为: - 线程池 - 工作线程 - 阻塞队列 - 拒绝策略 结合上面一二点的描述我们可以得出线程池中用到了工作线程和阻塞队列,而当阻塞队列满时需要根据拒绝策略进行任务拒绝,因此我们采取自下而上的方式逐一设计需要的几大主体。 ### 3.1、拒绝策略 --- 其实拒绝策略就是一段逻辑,通过调用者告知使用哪种方式进行任务拒绝。根据`OOP思想`,这一段逻辑我们可以封装成不同的方法,通过传入不同的标识选用不同的方法即可。这里使用了`Java1.8`出现的函数式编程进行设计,将这一个逻辑封装成一个函数式接口,调用者可直接使用Lambda表达式指定需要的拒绝策略,也可将逻辑封装成一个枚举类,直接传入对应的方法即可,这符合设计模式中的开闭原则,可维护性更高。 ```java @FunctionalInterface public interface RejectPolicy<T> { /** * @description 拒绝策略中的拒绝方法,可自定义设置适合的拒绝策略 * @author xbaozi * @date 2022/11/18 22:34 * @param queue 阻塞队列 * @param task 需要拒绝的任务 **/ void reject(BlockingQueue<T> queue, T task); } ``` ### 3.2、阻塞队列 --- 在该阻塞队列中采取了`公平的FIFO形式`,避免任务一直得不到消费出现**饿死情况**,因此内部需要维护一个**双向队列**。出于内存层面考虑,我们需要维护一个队列**最大容量**变量,用于判断队列是否已满,避免`OOM问题`出现。同时由于阻塞队列为多线程下的共享资源,我们需要对其上锁保证在并发消费下的原子性。最后为了阻塞队列的拓展性,队列中存放的内容采取泛型设计。 - **双向队列Deque**:Java内置的双向队列接口,实现类采用ArrayDeque,在大部分情况下会比LinkList性能要好一点; - **最大容量capacity**:基本整形变量,用于判断队列是否已满; - **锁对象LOCK**:采用可重入锁ReentrantLock实现,并分别设置消费者条件变量与生产者条件变量,对队列为空与队列已满两种情况进行隔离。 --- 在变量设计完成之后,我们还需要对队列中的**方法**进行设计。显而易见的是队列的**核心任务为存和取**,重点的是怎么存和取。比较容易想到的是**超时与无超时限制**的存取,但是这样子的话并没有用到我们的拒绝策略,因此应该还有一个方法是**尝试将任务存入队列**中,当队列满时采用指定的拒绝策略即可。 - `void put(T task)`:添加任务,这是一个**阻塞添加无超时**的方法,即这个方式在队列已满时会一直等待直到队列中出现闲余空间; - `boolean put(T task, long timeout, TimeUnit timeUnit)`:带超时限制的添加任务,当等待了指定时间后队列仍然无空间时,则会放弃当前任务退出等待; - `void tryPut(T task, RejectPolicy<T> rejectPolicy)`:**尝试添加任务**,如果任务队列满了会根据传入的**拒绝策略**对任务进行处理; - `T take()`:获取任务,这是一个**阻塞获取无超时**的方法,即队列为空时会一直等待直到队列中出现任务; - `T take(long timeout, TimeUnit unit)`:带**超时时间限制**返回获取到的任务,当等待了指定时间后队列仍然为空时,则会放弃获取退出等待。 ```java /** * @author xbaozi * @version 1.0 * @classname BlockingQueue * @date 2022-11-17 16:35 * @description 阻塞队列,使用泛型增加拓展性 */ @Slf4j(topic = "xbaoziplus.BlockingQueue") public class BlockingQueue<T> { // 任务队列,使用双向链表尾进头出 private final Deque<T> TASK_QUEUE = new ArrayDeque<>(); // 队列最大容量 private int capacity; // 锁对象 private final ReentrantLock LOCK = new ReentrantLock(); // 生产者条件变量,当队列中满了的话生产者线程需要进入该condition进行等待 private final Condition PRODUCER_WAIT_CONDITION = LOCK.newCondition(); // 消费者条件变量,当队列中为空时消费者线程需要进入该condition进行等待 private final Condition CONSUMER_WAIT_CONDITION = LOCK.newCondition(); // 有参构造器,初始化队列最大容量 public BlockingQueue(int capacity) { this.capacity = capacity; } /** * @param task 生产者产生的任务 * @description 添加任务,这是一个阻塞添加的方法 * @author xbaozi * @date 2022/11/17 16:57 **/ public void put(T task) { // 因为获取大小和添加任务不是原子操作,因此需要上锁保证原子性 LOCK.lock(); try { // 自旋判断队列是否已满,避免虚假唤醒 while (TASK_QUEUE.size() >= capacity) { try { log.error("队列已满,生产者等待将任务加入任务队列中……"); // 进入生产者条件变量中等待 PRODUCER_WAIT_CONDITION.await(); } catch (InterruptedException e) { e.printStackTrace(); } } // 任务队列出现闲余空间时,将任务采用尾插法添加至任务队列中 log.info("任务队列存在闲余空间,{}加入队列", task); TASK_QUEUE.addLast(task); // 唤醒消费者条件变量中线程,提示队列不为空,可以进行任务消费移除 CONSUMER_WAIT_CONDITION.signal(); } finally { LOCK.unlock(); } } /** * @param task 生产者产生的任务 * @param timeout 超时时间 * @param timeUnit 时间单位 * @description 带超时限制的添加任务 * @author xbaozi * @date 2022/11/18 21:38 **/ public boolean put(T task, long timeout, TimeUnit timeUnit) { LOCK.lock(); try { // 将超时时间转换成纳秒 long nanos = timeUnit.toNanos(timeout); while (TASK_QUEUE.size() >= capacity) { try { // 判断是否超时 if (nanos <= 0) { // 超时返回失败标识 log.error("添加任务{}超时", task); return false; } // 进入生产者条件变量中等待nanos秒或等待被唤醒 nanos = PRODUCER_WAIT_CONDITION.awaitNanos(nanos); } catch (InterruptedException e) { e.printStackTrace(); } } log.info("不超时,{}加入任务队列成功", task); // 尾插法插入任务 TASK_QUEUE.addLast(task); // 唤醒消费者条件变量中线程,提示队列不为空,可以进行任务消费移除 CONSUMER_WAIT_CONDITION.signal(); return true; } finally { LOCK.unlock(); } } /** * @description 添加任务,如果任务队列满了会根据传入的拒绝策略对任务进行处理 * @author xbaozi * @date 2022/11/18 23:15 * @param task 任务 * @param rejectPolicy 拒绝策略 **/ public void tryPut(T task, RejectPolicy<T> rejectPolicy) { LOCK.lock(); try { // 判断任务队列是否已满 if (TASK_QUEUE.size() >= capacity) { // 任务队列已满,采取设定的拒绝策略进行处理 rejectPolicy.reject(this, task); } else { // 任务队列未满,添加至任务队列中 log.info("{}加入任务队列成功", task); TASK_QUEUE.addLast(task); // 唤醒消费者条件变量中线程,提示队列不为空,可以进行任务消费移除 CONSUMER_WAIT_CONDITION.signal(); } } finally { LOCK.unlock(); } } /** * @return T 返回获取到的任务 * @description 获取任务,这是一个阻塞获取的方法 * @author xbaozi * @date 2022/11/17 17:05 **/ public T take() { // 因为判断队列是否为空和弹出任务不是原子操作,因此需要上锁保证原子性 LOCK.lock(); try { // 自旋判断队列是否为空,避免虚假唤醒 while (TASK_QUEUE.isEmpty()) { try { log.error("队列为空,消费者等待任务到来进行消费……"); // 进入消费者条件变量中等待 CONSUMER_WAIT_CONDITION.await(); } catch (InterruptedException e) { e.printStackTrace(); } } // 任务队列中产生了新任务时,从队列头部获取 T task = TASK_QUEUE.removeFirst(); log.info("队列中有任务{}取出消费", task); // 唤醒生产者条件变量中线程,提示队列已出现闲余空间,可以进行任务生产添加 PRODUCER_WAIT_CONDITION.signal(); // 返回任务 return task; } finally { LOCK.unlock(); } } /** * @description 带超时时间限制返回获取到的任务 * @author xbaozi * @date 2022/11/18 22:40 * @param timeout 超时时间 * @param unit 超时时间单位 **/ public T take(long timeout, TimeUnit unit) { LOCK.lock(); try { long nanos = unit.toNanos(timeout); // 判断任务队列中是否有任务可拿 while (TASK_QUEUE.isEmpty()) { if (nanos <= 0) { return null; } try { log.error("队列为空,消费者等待任务到来进行消费……"); // 进入消费者条件变量中等待 nanos = CONSUMER_WAIT_CONDITION.awaitNanos(nanos); } catch (InterruptedException e) { e.printStackTrace(); } } // 任务队列中有任务可拿时,从队列头部获取任务 T task = TASK_QUEUE.removeFirst(); log.info("队列中有任务{}取出消费", task); // 唤醒生产者条件变量中线程,提示队列已出现闲余空间,可以进行任务生产添加 PRODUCER_WAIT_CONDITION.signal(); // 返回任务 return task; } finally { LOCK.unlock(); } } /** * @description 获取阻塞队列中的任务数 * @author xbaozi * @date 2022/11/18 22:19 **/ public int size() { LOCK.lock(); try { return TASK_QUEUE.size(); } finally { LOCK.unlock(); } } } ``` ### 3.3、工作线程 --- 你可能在想着为什么还要自定义一个工作线程,直接用Thread不行吗?其实还真不行. 因为在工作线程中我们需要对task任务进行消费,而run方法并不支持传参,因此我们需要自定义一个**WorkerThread继承Thread**,并拓展一个成员变量task,通过构造器传参实现对任务的消费。 另外需要注意的是这个工作线程在设计的时候将其设定为线程池的内部类,因此代码在线程池中再一起贴出来 ### 3.4、线程池 --- 我们先根据需求来判断我们需要哪些成员变量。 首先我们需要核心线程来工作执行消费任务,因此需要一个**核心线程数**变量,并且需要一个**容纳线程的集合**存放工作线程; 其次我们在工作线程达到核心线程数时,需要将任务暂时存入阻塞队列中,因此需要一个**阻塞队列**的变量,值得注意的是在构造器初始化时应该传入队列容量在构造器中进行实例化,而不是传入一个阻塞队列对象,提供使用而不暴露实现; 紧接着的就是**超时时间**了,也可以将其忽略在执行方法时将其当做参数进行方法传参,这里放在成员变量中便于统一管理; 最后便是**拒绝策略**,在初始化线程池时就应该指定线程池的拒绝策略,在阻塞队列满时对任务进行拒绝。 ```java @Slf4j(topic = "xbaoziplus.MyThreadPool") public class ThreadPool { // 任务阻塞队列 private final BlockingQueue<Runnable> BLOCKING_QUEUE; // 线程集合 private final Set<WorkerThread> workers = new HashSet<>(); // 核心线程数 private int coreSize; // 获取任务的超时时间 private long timeout; // 获取任务超时时间的时间单位 private TimeUnit unit; // 拒绝策略 RejectPolicy<Runnable> rejectPolicy; public ThreadPool(int queueCapacity, int coreSize, int timeout, TimeUnit unit, RejectPolicy<Runnable> rejectPolicy) { this.BLOCKING_QUEUE = new BlockingQueue<>(queueCapacity); this.coreSize = coreSize; this.timeout = timeout; this.unit = unit; this.rejectPolicy = rejectPolicy; } /** * @param task 需要执行的任务 * @description 线程池接收任务执行 * @author xbaozi * @date 2022/11/17 17:52 **/ public void execute(Runnable task) { synchronized (workers) { // 判断工作线程数是否达到了核心线程数 if (workers.size() < coreSize) { // 工作线程数未达到核心线程数,新建一个工作线程 WorkerThread workerThread = new WorkerThread(task, workers.size() + "号工作线程"); log.info("未达到核心线程数,新建工作线程{}", workerThread.getName()); // 将工作线程添加到线程集合中 workers.add(workerThread); // 启动线程执行任务 workerThread.start(); } else { // 工作线程已达到核心线程数,将任务放入任务队列中暂存 // log.info("工作线程已达到核心线程数,{}进入任务队列暂存", task); // 无超时阻塞添加任务 // BLOCKING_QUEUE.put(task); // 设置超时时间添加任务 // BLOCKING_QUEUE.put(task, timeout, unit); // 尝试添加任务,任务添加失败时选择自定义的拒绝策略 log.info("工作线程已达到核心线程数,尝试添加任务{}到任务队列中暂存", task); BLOCKING_QUEUE.tryPut(task, rejectPolicy); } } } /** * @author xbaozi * @description 线程池中工作线程 * @date 2022/11/17 17:26 **/ //@Slf4j(topic = "xbaoziplus.MyThreadPool.WorkerThread") class WorkerThread extends Thread { // 需要执行的任务 private Runnable task; public WorkerThread(Runnable task, String name) { super(name); this.task = task; } @Override public void run() { // 自旋判断当前任务是否为空,不为空执行任务,为空时获取下一个任务接着执行 // while (task != null || (task = BLOCKING_QUEUE.take()) != null) { // 无超时限制的等待获取任务 // 有超时限制的等待获取任务 while (task != null || (task = BLOCKING_QUEUE.take(timeout, unit)) != null) { try { // 执行任务 log.info("{}正在执行……", task); task.run(); } catch (Exception e) { e.printStackTrace(); } finally { // 不能在这里赋值task = BLOCKING_QUEUE.take(),否则容易产生线程饥饿 task = null; } } // 执行完毕时,将当前工作线程移除,实现线程销毁效果 synchronized (workers) { log.info("任务执行完毕,线程{}已销毁", this.getName()); workers.remove(this); } } } } ``` ## 4、测试线程池 这里并没有封装一部分的拒绝策略给调用者进行选择,而是完全由调用者编写拒绝策略。这里一共列举了五种拒绝策略,分别为: - 死等,直到阻塞队列出现空间或工作线程空余; - 超时等待,等待一定时间后自行结束; - 直接放弃,当阻塞队列已满时直接对后续的任务进行放弃; - 抛异常,当阻塞队列已满时抛出异常提示调用者; - 自行执行任务,让调用者线程自行执行多出来的任务。 ```java @Slf4j(topic = "test.TestPool") public class TestPool { public static void main(String[] args) { // 新建线程池 ThreadPool pool = new ThreadPool(2, 2, 2, TimeUnit.SECONDS, (queue, task) -> { // 1. 死等 //queue.put(task); // 2) 超时等待 //queue.put(task, 1500, TimeUnit.MILLISECONDS); // 3) 直接放弃 log.debug("任务队列已满,放弃任务{}", task); // 4) 抛出异常 //try { // throw new RuntimeException("任务执行失败 " + task); //} catch (RuntimeException e) { // log.debug("任务队列已满,{}", e.getMessage()); //} // 5) 自行执行任务 // task.run(); }); // 模拟五个线程生产任务压入线程池中执行 for (int i = 0; i < 5; i++) { int index = i; pool.execute(() -> { try { // 模拟任务执行需要1s Thread.sleep(1000L); } catch (InterruptedException e) { e.printStackTrace(); } log.info("任务{}执行完毕", index); }); } } } ``` 
一周上线百万级高并发系统
地址:https://mp.weixin.qq.com/s/UjBlaZIJ0UePmA3A10Atmw
并发编程网
收录了一些软件架构开发技术的文章 地址:https://ifeve.com/
JDK1.6 以后对 Synchronized 做了哪些优化?
JDK1.6 以后对 Synchronized 做了哪些优化?
鱼皮哥,想问一下如何提高单体并
鱼皮哥,想问一下如何提高单体并发量,是用一些并发包或者线程池或者用消息队列去削峰吗
