首页 文章 精选 留言 我的

精选列表

搜索[团队管理],共10009篇文章
优秀的个人博客,低调大师

Pipeline模式应用 | 京东云技术团队

本文记录Pipeline设计模式在业务流程编排中的应用 前言 Pipeline模式意为管道模式,又称为流水线模式。旨在通过预先设定好的一系列阶段来处理输入的数据,每个阶段的输出即是下一阶段的输入。 本案例通过定义PipelineProduct(管道产品),PipelineJob(管道任务),PipelineNode(管道节点),完成一整条流水线的组装,并将“原材料”加工为“商品”。其中管道产品负责承载各个阶段的产品信息;管道任务负责不同阶段对产品的加工;管道节点约束了管道产品及任务的关系,通过信号量定义了任务的执行方式。 依赖 工具依赖如下 <!-- 工具类大全 --> <dependency> <groupId>cn.hutool</groupId> <artifactId>hutool-all</artifactId> <version>最新版本</version> </dependency> 编程示例 1. 管道产品定义 package com.example.demo.pipeline.model; /** * 管道产品接口 * * @param <S> 信号量 * @author * @date 2023/05/15 11:49 */ public interface PipelineProduct<S> { } 2. 管道任务定义 package com.example.demo.pipeline.model; /** * 管道任务接口 * * @param <P> 管道产品 * @author * @date 2023/05/15 11:52 */ @FunctionalInterface public interface PipelineJob<P> { /** * 执行任务 * * @param product 管道产品 * @return {@link P} */ P execute(P product); } 3. 管道节点定义 package com.jd.baoxian.mall.market.service.pipeline.model; import java.util.function.Predicate; /** * 管道节点定义 * * @param <S> 信号量 * @param <P> 管道产品 * @author * @date 2023/05/15 11:54 */ public interface PipelineNode<S, P extends PipelineProduct<S>> { /** * 节点组装,按照上个管道任务传递的信号,执行 pipelineJob * * @param pipelineJob 管道任务 * @return {@link PipelineNode}<{@link S}, {@link P}> */ PipelineNode<S, P> flax(PipelineJob<P> pipelineJob); /** * 节点组装,按照传递的信号,判断当前管道的信号是否相等,执行 pipelineJob * * @param signal 信号 * @param pipelineJob 管道任务 * @return {@link PipelineNode}<{@link S}, {@link P}> */ PipelineNode<S, P> flax(S signal, PipelineJob<P> pipelineJob); /** * 节点组装,按照传递的信号,判断当前管道的信号是否相等,执行 pipelineJob * * @param predicate 信号 * @param pipelineJob 管道任务 * @return {@link PipelineNode}<{@link S}, {@link P}> */ PipelineNode<S, P> flax(Predicate<S> predicate, PipelineJob<P> pipelineJob); /** * 管道节点-任务执行 * * @param product 管道产品 * @return {@link P} */ P execute(P product); } 4. 管道产品、任务,节点的实现 4.1 管道产品 package com.example.demo.pipeline.factory; import com.example.demo.model.request.DemoReq; import com.example.demo.model.response.DemoResp; import com.example.demo.pipeline.model.PipelineProduct; import lombok.*; /** * 样例-管道产品 * * @author * @date 2023/05/15 14:04 */ @Data @Builder @NoArgsConstructor @AllArgsConstructor public class DemoPipelineProduct implements PipelineProduct<DemoPipelineProduct.DemoSignalEnum> { /** * 信号量 */ private DemoSignalEnum signal; /** * 产品-入参及回参 */ private DemoProductData productData; /** * 异常信息 */ private Exception exception; /** * 流程Id */ private String tradeId; @Data @Builder @NoArgsConstructor @AllArgsConstructor public static class DemoProductData { /** * 待验证入参 */ private DemoReq userRequestData; /** * 待验证回参 */ private DemoResp userResponseData; } /** * 产品-信号量 * * @author * @date 2023/05/15 13:54 */ @Getter public enum DemoSignalEnum { /** * */ NORMAL(0, "正常"), /** * */ CHECK_NOT_PASS(1, "校验不通过"), /** * */ BUSINESS_ERROR(2, "业务异常"), /** * */ LOCK_ERROR(3, "锁处理异常"), /** * */ DB_ERROR(4, "事务处理异常"), ; /** * 枚举码值 */ private final int code; /** * 枚举描述 */ private final String desc; /** * 构造器 * * @param code * @param desc */ DemoSignalEnum(int code, String desc) { this.code = code; this.desc = desc; } } } 4.2 管道任务(抽象类) package com.example.demo.pipeline.factory.job; import cn.hutool.core.util.ClassUtil; import cn.hutool.json.JSONUtil; import com.example.demo.pipeline.factory.DemoPipelineProduct; import com.example.demo.pipeline.model.PipelineJob; import lombok.extern.slf4j.Slf4j; /** * 管道任务-抽象层 * * @author * @date 2023/05/15 19:48 */ @Slf4j public abstract class AbstractDemoJob implements PipelineJob<DemoPipelineProduct> { /** * 公共执行逻辑 * * @param product 产品 * @return */ @Override public DemoPipelineProduct execute(DemoPipelineProduct product) { DemoPipelineProduct.DemoSignalEnum newSignal; try { newSignal = execute(product.getTradeId(), product.getProductData()); } catch (Exception e) { product.setException(e); newSignal = DemoPipelineProduct.DemoSignalEnum.BUSINESS_ERROR; } product.setSignal(newSignal); defaultLogPrint(product.getTradeId(), product); return product; } /** * 子类执行逻辑 * * @param tradeId 流程Id * @param productData 请求数据 * @return * @throws Exception 异常 */ abstract DemoPipelineProduct.DemoSignalEnum execute(String tradeId, DemoPipelineProduct.DemoProductData productData) throws Exception; /** * 默认的日志打印 */ public void defaultLogPrint(String tradeId, DemoPipelineProduct product) { if (!DemoPipelineProduct.DemoSignalEnum.NORMAL.equals(product.getSignal())) { log.info("流水线任务处理异常:流程Id=【{}】,信号量=【{}】,任务=【{}】,参数=【{}】", tradeId, product.getSignal(), ClassUtil.getClassName(this, true), JSONUtil.toJsonStr(product.getProductData()), product.getException()); } } } 4.3 管道节点 package com.example.demo.pipeline.factory; import cn.hutool.core.util.ClassUtil; import cn.hutool.json.JSONUtil; import com.example.demo.pipeline.model.PipelineJob; import com.example.demo.pipeline.model.PipelineNode; import lombok.extern.slf4j.Slf4j; import java.util.function.Predicate; /** * 审核-管道节点 * * @author * @date 2023/05/15 14:32 */ @Slf4j public class DemoPipelineNode implements PipelineNode<DemoPipelineProduct.DemoSignalEnum, DemoPipelineProduct> { /** * 下一管道节点 */ private DemoPipelineNode next; /** * 当前管道任务 */ private PipelineJob<DemoPipelineProduct> job; /** * 节点组装,按照上个管道任务传递的信号,执行 pipelineJob * * @param pipelineJob 管道任务 * @return {@link DemoPipelineNode} */ @Override public DemoPipelineNode flax(PipelineJob<DemoPipelineProduct> pipelineJob) { return flax(DemoPipelineProduct.DemoSignalEnum.NORMAL, pipelineJob); } /** * 节点组装,按照传递的信号,判断当前管道的信号是否相等,执行 pipelineJob * * @param signal 信号 * @param pipelineJob 管道任务 * @return {@link DemoPipelineNode} */ @Override public DemoPipelineNode flax(DemoPipelineProduct.DemoSignalEnum signal, PipelineJob<DemoPipelineProduct> pipelineJob) { return flax(signal::equals, pipelineJob); } /** * 节点组装,上个管道过来的信号运行 predicate 后是true的话,执行 pipelineJob * * @param predicate * @param pipelineJob * @return */ @Override public DemoPipelineNode flax(Predicate<DemoPipelineProduct.DemoSignalEnum> predicate, PipelineJob<DemoPipelineProduct> pipelineJob) { this.next = new DemoPipelineNode(); this.job = (job) -> { if (predicate.test(job.getSignal())) { return pipelineJob.execute(job); } else { return job; } }; return next; } /** * 管道节点-任务执行 * * @param product 管道产品 * @return */ @Override public DemoPipelineProduct execute(DemoPipelineProduct product) { // 执行当前任务 try { product = job == null ? product : job.execute(product); return next == null ? product : next.execute(product); } catch (Exception e) { log.error("流水线处理异常:流程Id=【{}】,任务=【{}】,参数=【{}】", product.getTradeId(), ClassUtil.getClassName(job, true), JSONUtil.toJsonStr(product.getProductData()), product.getException()); return null; } } } 5. 业务实现 通过之前的定义,我们已经可以通过Pipeline完成流水线的搭建,接下来以“审核信息提交”这一业务场景,完成应用。 5.1 定义Api、入参、回参 package com.example.demo.api; import com.example.demo.model.request.DemoReq; import com.example.demo.model.response.DemoResp; import com.example.demo.pipeline.factory.PipelineForManagerSubmit; import org.springframework.stereotype.Service; import javax.annotation.Resource; /** * 演示-API * * @author * @date 2023/08/06 16:27 */ @Service public class DemoManagerApi { /** * 管道-审核提交 */ @Resource private PipelineForManagerSubmit pipelineForManagerSubmit; /** * 审核提交 * * @param requestData 请求数据 * @return {@link DemoResp} */ public DemoResp managerSubmit(DemoReq requestData) { return pipelineForManagerSubmit.managerSubmitCheck(requestData); } } package com.example.demo.model.request; /** * 演示入参 * * @author * @date 2023/08/06 16:33 */ public class DemoReq { } package com.example.demo.model.response; import lombok.Data; /** * 演示回参 * * @author * @date 2023/08/06 16:33 */ @Data public class DemoResp { /** * 成功标识 */ private Boolean success = false; /** * 结果信息 */ private String resultMsg; /** * 构造方法 * * @param message 消息 * @return {@link DemoResp} */ public static DemoResp buildRes(String message) { DemoResp response = new DemoResp(); response.setResultMsg(message); return response; } } 5.2 定义具体任务 假定审核提交的流程需要包含:参数验证、加锁、解锁、事务提交 package com.example.demo.pipeline.factory.job; import cn.hutool.json.JSONUtil; import com.example.demo.model.request.DemoReq; import com.example.demo.pipeline.factory.DemoPipelineProduct; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; /** * 加锁-实现层 * * @author * @date 2023/05/17 17:00 */ @Service @Slf4j public class CheckRequestLockJob extends AbstractDemoJob { /** * 子类执行逻辑 * * @param tradeId 流程Id * @param productData 请求数据 * @return * @throws Exception 异常 */ @Override DemoPipelineProduct.DemoSignalEnum execute(String tradeId, DemoPipelineProduct.DemoProductData productData) throws Exception { DemoReq userRequestData = productData.getUserRequestData(); log.info("任务[{}]加锁,线程号:{}", JSONUtil.toJsonStr(userRequestData), tradeId); return DemoPipelineProduct.DemoSignalEnum.NORMAL; } } package com.example.demo.pipeline.factory.job; import cn.hutool.json.JSONUtil; import com.example.demo.model.request.DemoReq; import com.example.demo.pipeline.factory.DemoPipelineProduct; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; /** * 解锁-实现层 * * @author * @date 2023/05/17 17:00 */ @Service @Slf4j public class CheckRequestUnLockJob extends AbstractDemoJob { /** * 子类执行逻辑 * * @param tradeId 流程Id * @param productData 请求数据 * @return * @throws Exception 异常 */ @Override DemoPipelineProduct.DemoSignalEnum execute(String tradeId, DemoPipelineProduct.DemoProductData productData) throws Exception { DemoReq userRequestData = productData.getUserRequestData(); log.info("任务[{}]解锁,线程号:{}", JSONUtil.toJsonStr(userRequestData), tradeId); return DemoPipelineProduct.DemoSignalEnum.NORMAL; } } package com.example.demo.pipeline.factory.job; import cn.hutool.json.JSONUtil; import com.example.demo.model.request.DemoReq; import com.example.demo.pipeline.factory.DemoPipelineProduct; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; /** * 审核-参数验证-实现类 * * @author * @date 2023/05/15 19:50 */ @Slf4j @Component public class ManagerCheckParamJob extends AbstractDemoJob { /** * 执行基本入参验证 * * @param tradeId * @param productData 请求数据 * @return */ @Override DemoPipelineProduct.DemoSignalEnum execute(String tradeId, DemoPipelineProduct.DemoProductData productData) { /* * 入参验证 */ DemoReq userRequestData = productData.getUserRequestData(); log.info("任务[{}]入参验证,线程号:{}", JSONUtil.toJsonStr(userRequestData), tradeId); // 非空验证 // 有效验证 // 校验通过,退出 return DemoPipelineProduct.DemoSignalEnum.NORMAL; } } package com.example.demo.pipeline.factory.job; import cn.hutool.json.JSONUtil; import com.example.demo.model.request.DemoReq; import com.example.demo.model.response.DemoResp; import com.example.demo.pipeline.factory.DemoPipelineProduct; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; /** * 审核-信息提交-业务实现 * * @author * @date 2023/05/12 14:36 */ @Service @Slf4j public class ManagerSubmitJob extends AbstractDemoJob { /** * 子类执行逻辑 * * @param tradeId 流程Id * @param productData 请求数据 * @return * @throws Exception 异常 */ @Override DemoPipelineProduct.DemoSignalEnum execute(String tradeId, DemoPipelineProduct.DemoProductData productData) throws Exception { DemoReq userRequestData = productData.getUserRequestData(); try { /* * DB操作 */ log.info("任务[{}]信息提交,线程号:{}", JSONUtil.toJsonStr(userRequestData), tradeId); productData.setUserResponseData(DemoResp.buildRes("成功")); } catch (Exception ex) { log.error("审核-信息提交-DB操作失败,入参:{}", JSONUtil.toJsonStr(userRequestData), ex); throw ex; } return DemoPipelineProduct.DemoSignalEnum.NORMAL; } } 5.3 完成流水线组装 针对入回参转换,管道任务执行顺序及执行信号量的构建 package com.example.demo.pipeline.factory; import com.example.demo.model.request.DemoReq; import com.example.demo.model.response.DemoResp; import com.example.demo.pipeline.factory.job.CheckRequestLockJob; import com.example.demo.pipeline.factory.job.CheckRequestUnLockJob; import com.example.demo.pipeline.factory.job.ManagerCheckParamJob; import com.example.demo.pipeline.factory.job.ManagerSubmitJob; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import javax.annotation.PostConstruct; import java.util.Objects; import java.util.UUID; /** * 管道工厂入口-审核流水线 * * @author * @date 2023/05/15 19:52 */ @Slf4j @Service @RequiredArgsConstructor public class PipelineForManagerSubmit { /** * 审核-管道节点 */ private final DemoPipelineNode managerSubmitNode = new DemoPipelineNode(); /** * 审核-管道任务-提交-防刷锁-加锁 */ private final CheckRequestLockJob checkRequestLockJob; /** * 审核-管道任务-提交-防刷锁-解锁 */ private final CheckRequestUnLockJob checkRequestUnLockJob; /** * 审核-管道任务-参数验证 */ private final ManagerCheckParamJob managerCheckParamJob; /** * 审核-管道任务-事务操作 */ private final ManagerSubmitJob managerSubmitJob; /** * 组装审核的处理链 */ @PostConstruct private void assembly() { assemblyManagerSubmit(); } /** * 组装处理链 */ private void assemblyManagerSubmit() { managerSubmitNode // 参数验证及填充 .flax(managerCheckParamJob) // 防刷锁 .flax(checkRequestLockJob) // 事务操作 .flax(managerSubmitJob) // 锁释放 .flax((ignore) -> true, checkRequestUnLockJob); } /** * 审核-提交处理 * * @param requestData 入参 * @return */ public DemoResp managerSubmitCheck(DemoReq requestData) { DemoPipelineProduct initialProduct = managerSubmitCheckInitial(requestData); DemoPipelineProduct finalProduct = managerSubmitNode.execute(initialProduct); if (Objects.isNull(finalProduct) || Objects.nonNull(finalProduct.getException())) { return DemoResp.buildRes("未知异常"); } return finalProduct.getProductData().getUserResponseData(); } /** * 审核-初始化申请的流水线数据 * * @param requestData 入参 * @return 初始的流水线数据 */ private DemoPipelineProduct managerSubmitCheckInitial(DemoReq requestData) { // 初始化 return DemoPipelineProduct.builder() .signal(DemoPipelineProduct.DemoSignalEnum.NORMAL) .tradeId(UUID.randomUUID().toString()) .productData(DemoPipelineProduct.DemoProductData.builder().userRequestData(requestData).build()) .build(); } } 总结 本文重点为管道模式的抽象与应用,上述示例仅为个人理解。实际应用中,此案例长于应对各种规则冗杂的业务场景,便于规则编排。 待改进点: 各个任务其实隐含了执行的先后顺序,此项内容可进一步实现; 针对最后“流水线组装”这一步,可通过配置描述的方式,进一步抽象,从而将变动控制在每个“管道任务”的描述上,针对规则项做到“可插拔”式处理。 作者:京东保险 侯亚东 来源:京东云开发者社区 转载请注明来源

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

Tomcat目录结构 | 京东云技术团队

Tomcat目录结构图如下: 1、bin目录 存放一些可执行的二进制文件,****.sh 结尾的为linux下执行命令,****.bat 结尾的为windows下执行命令。 catalina.sh:真正启动tomcat文件,可以在里面设置jvm参数。 startup.sh:启动tomcat(需事先配置好JAVA_HOME环境变量才可启动,该命令源码实际执行的为catalina.sh start)。 shutdown.sh:关闭tomcat。 version.sh:查看tomcat版本相关信息。 2、conf目录 存放tomcat相关配置文件的。 2.1、catalina.policy 项目安全文件,用来防止欺骗代码或JSP执行带有像System.exit(0)这样的命令,可能影响容器的破坏。 只有当Tomcat用-security命令行参数启动时这个文件才会被使用,即启动tomcat时, startup.sh -security 。 2.2、catalina.proterties 配置tomcat启动相关信息文件 2.3、context.xml 监视并加载资源文件,当监视文件发生变化时,自动加载,通常不会去配置 2.4、jaspic-providers.xml和jaspic-providers.xsd 不常用文件 2.5、logging.properties tomcat日志文件配置,包括输出格式、日志级别等。 2.6、server.xml 核心配置文件:修改端口号,添加编码格式等 核心组件介绍: <1>Server:最顶层元素,而且唯一,代表整个tomcat容器。一个Server元素包含一个或者多个Service元素; <2>Service:对外提供服务的。一个Service元素包含多个Connector元素,但是只能包含一个Engine元素; <3>Connector:接收连接请求,创建Request和Response对象用于和请求端交换数据;然后分配线程让Engine来处理这个请求,并把产生的Request和Response对象传给Engine <4>Engine:Engine组件在Service组件中有且只有一个;Engine是Service组件中的请求处理组件。Engine组件从一个或多个Connector中接收请求并处理,并将完成的响应返回给Connector,最终传递给客户端。 <5>Host:代表特定的虚拟主机。 <Host name="localhost" appBase="webapps" unpackWARs="true" autoDeploy="true"> **name:**虚拟主机的主机名。比如 localhost 表示本机名称,实际应用时应该填写具体域名,比如 www.dog.com ,当然如果该虚拟主机是给内部人员访问的,也可以直接填写服务器的 ip 地址,比如 192.168.1.101; **appBase:**设置 Web 应用程序组的路径。appBase 属性的值可以是相对于 Tomcat 安装目录的相对路径,也可以是绝对路径,需要注意的是该路径必须是 Tomcat 有权限访问的; **unpackWARs:**是否自动展开war压缩包再运行Web应用程序,默认值为true; **autoDeplay:**是否允许自动部署,默认值是 true,表示 Tomcat 会自动检测 appBase 目录下面的文件变化从而自动应用到正在运行的 Web 应用程序; **deployOnStartup:**为true时,表示Tomcat在启动时检查Web应用,且检测到的所有Web应用视作新应用; <6>Context:该元素代表在特定虚拟主机Host上运行的一个Web应用,它是Host的子容器,每个Host容器可以定义多个Context元素。静态部署Web应用时使用。 <Context path="/" docBase="E:\Resource\test.war" reloadable="true"/> **path:**浏览器访问时的路径名,只有当自动部署完全关闭(deployOnStartup和autoDeploy都为false)或docBase不在appBase中时,才可以设置path属性。 **docBase:**静态部署时,docBase可以在appBase目录下,也可以不在;本例中,不在appBase目录下。 **reloadable:**设定项目有改动时,重新加载该项目。 2.7、tomcat-users.xml和tomcat-users.xsd tomcat-users.xml:tomcat用户配置文件,配置用户名,密码,用户具备权限 tomcat默认没有配置任何用户,只有配置好用户后才能使用以下Tomcat Manager三个功能: <role rolename="manager-gui"/> <role rolename="manager-script"/> <user username="tomcat" password="tomcat" roles="manager-gui"/> <user username="admin" password="123456" roles="manager-script"/> tomcat-users.xsd:对tomcat-users.xml文件的描述和约束 2.8、web.xml web应用相关通用配置,可以做下面这些事情。 配置servlet 添加过滤器,比如过滤敏感词汇 设置session过期时间,tomcat默认30分钟 注册了很多MIME类型,即文档类型。这些MIME类型是客户端与服务器之间说明文档类型的,如用户请求一个html网页,那么服务器还会告诉客户端浏览器响应的文档是text/html类型的,这就是一个MIME类型 配置系统欢迎页 3、lib目录 存放tomcat依赖jar包的。 其中ecj-x.x.x.jar起到了将.java文件编译成.class字节码文件的作用。 4、logs目录 存放tomcat运行时产生的日志文件。 在windows环境中,日志文件输出到catalina.xxxx-xx-xx.log文件中。 在linux环境中,日志文件输出到catalina.out文件中。 大体有以下几类: catalina.xxxx-xx-xx.log windows下日志文件输出内容 host-manager.xxxx-xx-xx.log 访问webapps下host-manager项目日志 localhost.xxxx-xx-xx.log tomcat启动时,自身访问服务,只记录tomcat访问日志,而非业务项目日志 localhost_access_log.xxxx-xx-xx.txt 表示访问tomcat下所有项目日志记录 manager.xxxx-xx-xx.log 访问webapps下manager项目日志 5、temp目录 用户存放tomcat在运行过程中产生的临时文件(清空不会对tomcat运行带来影响)。 6、webapps目录 用来存放应用程序,可以以文件夹、war包、jar包的形式发布应用。当然也可以将应用程序放在磁盘的任意位置,在配置文件中映射好即可。 默认自带以下5个项目: 7、work目录 用于存放tomcat在运行时的编译后文件(清空该目录下所有内容,重启tomcat,可达到清除缓冲的作用) 作者:京东科技 杨建 来源:京东云开发者社区 转载请注明来源

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

ReentrantLock源码解析 | 京东云技术团队

并发指同一时间内进行了多个线程。并发问题是多个线程对同一资源进行操作时产生的问题。通过加锁可以解决并发问题,ReentrantLock是锁的一种。 1 ReentrantLock 1.1 定义 ReentrantLock是Lock接口的实现类,可以手动的对某一段进行加锁。ReentrantLock可重入锁,具有可重入性,并且支持可中断锁。其内部对锁的控制有两种实现,一种为公平锁,另一种为非公平锁. 1.2 实现原理 ReentrantLock的实现原理为volatile+CAS。想要说明volatile和CAS首先要说明JMM。 1.2.1 JMM JMM(java 内存模型 Java Memory Model 简称JMM) 本身是一个抽象的概念,并不在内存中真实存在的,它描述的是一组规范或者规则,通过这组规范定义了程序中各个变量的访问方式. 由于 JMM 运行程序的实体是线程.而每个线程创建时JMM都会为其创建一个自己的工作内存(栈空间),工作内存是每个线程的私有 数据区域.而java内存模型中规定所有的变量都存储在主内存中,主内存是共享内存区域,所有线程都可以访问,但线程的变量的操作(读取赋值等)必须在自己的工作内存中去进行,首先要 将变量从主存拷贝到自己的工作内存中,然后对变量进行操作,操作完成后再将变量操作完后的新值写回主内存,不能直接操作主内存的变量,各个线程的工作内存中存储着主内存的变量拷贝的副本,因不同的线程间无法访问对方的工作内存,线程间的通信必须在主内存来完成。 如图所示:线程A对变量A的操作,只能是从主内存中拷贝到线程中,再写回到主内存中。 1.2.2 volatile volatile 是JAVA的关键字用于修饰变量,是java虚拟机的轻量同步机制,volatile不能保证原子性。 作用: 线程可见性:一个变量在某个线程里修改了它的值,如果使用了volatile关键字,那么别的线程可以马上读到修改后的值。 指令重排序:没加之前,指令是并发执行的,第一个线程执行到一半另一个线程可能开始执行了。加了volatile关键字后,不同线程是按照顺序一步一步执行的。1.2.3 CASCAS是Compare and Swap,就是比较和交换,而比较和交换是一个原子操作。线程基于CAS修改数据的方式:先获取主内存数据,在修改之前,先比较数据是否一致,如果一致修改主内存数据,如果不一致,放弃这次修改。 作用:CAS会使用现代处理器上提供的高效机器级别原子指令,这些原子指令以原子方式对内存执行读-改-写操作。1.2.4 AQSAQS的全称是AbstractQueuedSynchronizer(抽象的队列式的同步器),AQS定义了一套多线程访问共享资源的同步器框架。 AQS主要包含两部分内容:共享资源和等待队列。AQS底层已经对这两部分内容提供了很多方法。 共享资源:共享资源是一个volatile的int类型变量。 等待队列:等待队列是一个线程安全的队列,当线程拿不到锁时,会被park并放入队列。 2 源码解析 ReentrantLock在包java.util.concurrent.locks下,实现Lock接口。 2.1 lock方法 lock分为公平锁和非公平锁。 公平锁: final void lock() { acquire(1); } 非公平锁:上来先尝试将state从0修改为1,如果成功,代表获取锁资源。如果没有成功,调用acquire。state是AQS中的一个由volatile修饰的int类型变量,多个线程会通过CAS的方式修改state,在并发情况下,只会有一个线程成功的修改state。 final void lock() { //通过原子方式修改值 if (compareAndSetState(0, 1)) setExclusiveOwnerThread(Thread.currentThread()); else acquire(1); } /** * 获取锁的线程. */ private transient Thread exclusiveOwnerThread; /** * 设置拥有锁的线程 */ protected final void setExclusiveOwnerThread(Thread thread) { exclusiveOwnerThread = thread; 2.2 acquire方法 acquire是一个业务方法,里面并没有实际的业务处理,都是在调用其他方法。 public final void acquire(int arg) { //调用tryAcquire方法:尝试获取锁资源(非公平、公平),拿到锁资源,返回true,直接结束方法 if (!tryAcquire(arg) && //当没有获取锁资源后,会先调用addWaiter:会将没有获取到锁资源的线程封装为Node对象, //并且插入到AQS的队列的末尾. //继续调用acquireQueued方法,查看当前排队的Node是否在队列的前面,如果在前面,尝试获取锁资源 //如果没在前面,尝试将线程挂起,阻塞起来! acquireQueued(addWaiter(Node.EXCLUSIVE), arg)) selfInterrupt(); } 2.3 tryAcquire方法 tryAcquire分为公平和非公平两种。 公平: protected final boolean tryAcquire(int acquires) { //拿到当前线程 final Thread current = Thread.currentThread(); //拿到AQS的state int c = getState(); // 如果state == 0,说明没有线程占用着当前的锁资源 if (c == 0) { //如果没有线程排队,直接直接CAS尝试获取锁资源 if (!hasQueuedPredecessors() && compareAndSetState(0, acquires)) { //如果获取资源成功,将当前线程设置为持有锁资源的线程 setExclusiveOwnerThread(current); return true; } } //如果有线程持有锁资源,判断持有锁资源的线程是否是当前线程 else if (current == getExclusiveOwnerThread()) { //增加AQS的state的值 int nextc = c + acquires; if (nextc < 0) throw new Error("Maximum lock count exceeded"); setState(nextc); return true; } return false; } } 非公平: protected final boolean tryAcquire(int acquires) { return nonfairTryAcquire(acquires); } final boolean nonfairTryAcquire(int acquires) { //拿到当前线程 final Thread current = Thread.currentThread(); //拿到AQS的state int c = getState(); // 如果state == 0,说明没有线程占用着当前的锁资源 if (c == 0) { //获取锁资源 if (compareAndSetState(0, acquires)) { //将当前占用这个互斥锁的线程属性设置为当前线程 setExclusiveOwnerThread(current); return true; } } //如果有线程持有锁资源,判断持有锁资源的线程是否是当前线程 else if (current == getExclusiveOwnerThread()) { int nextc = c + acquires; if (nextc < 0) // overflow throw new Error("Maximum lock count exceeded"); setState(nextc); return true; } return false; } 2.4 addWaiter方法 在获取锁资源失败后,需要将当前线程封装为Node对象,并且插入到AQS队列的末尾。 private Node addWaiter(Node mode) { // 将当前线程封装为Node对象,mode为null,代表互斥锁 Node node = new Node(Thread.currentThread(), mode); // pred是tail节点 Node pred = tail; // 如果pred不为null,有线程正在排队 if (pred != null) { // 将当前节点的prev,指定tail尾节点 node.prev = pred; // 以CAS的方式,将当前节点变为tail节点 if (compareAndSetTail(pred, node)) { // 之前的tail的next指向当前节点 pred.next = node; return node; } } // 添加的流程为, 自己prev指向、tail指向自己、前节点next指向我 // 如果上述方式,CAS操作失败,导致加入到AQS末尾失败,如果失败,就基于enq的方式添加到AQS队列 enq(node); return node; } // enq,无论怎样都添加进入 private Node enq(final Node node) { for (;;) { // 拿到tail Node t = tail; // 如果tail为null,说明当前没有Node在队列中 if (t == null) { // 创建一个新的Node作为head,并且将tail和head指向一个Node if (compareAndSetHead(new Node())) tail = head; } else { // 和上述代码一致! node.prev = t; if (compareAndSetTail(t, node)) { t.next = node; return t; } } } } 2.5 acquireQueued方法 // acquireQueued方法 // 查看当前排队的Node是否是head的next, // 如果是,尝试获取锁资源, // 如果不是或者获取锁资源失败那么就尝试将当前Node的线程挂起(unsafe.park()) final boolean acquireQueued(final Node node, int arg) { boolean failed = true; try { for (;;) { // 拿到上一个节点 final Node p = node.predecessor(); if (p == head && // 说明当前节点是head的next tryAcquire(arg)) { // 竞争锁资源,成功:true,失败:false // 进来说明拿到锁资源成功 // 将当前节点置位head,thread和prev属性置位null setHead(node); // 帮助快速GC p.next = null; // 设置获取锁资源成功 failed = false; // 不管线程中断。 return interrupted; } // 如果不是或者获取锁资源失败,尝试将线程挂起 // 第一个事情,当前节点的上一个节点的状态正常! // 第二个事情,挂起线程 if (shouldParkAfterFailedAcquire(p, node) && // 通过LockSupport将当前线程挂起 parkAndCheckInterrupt()) } } finally { if (failed) cancelAcquire(node); } } // 确保上一个节点状态是正确的 private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) { // 拿到上一个节点的状态 int ws = pred.waitStatus; // 如果上一个节点为 -1 if (ws == Node.SIGNAL) // 返回true,挂起线程 return true; // 如果上一个节点是取消状态 if (ws > 0) { // 循环往前找,找到一个状态小于等于0的节点 do { node.prev = pred = pred.prev; } while (pred.waitStatus > 0); pred.next = node; } else { // 将小于等于0的节点状态该为-1 compareAndSetWaitStatus(pred, ws, Node.SIGNAL); } return false; } 2.6 unlock方法 释放锁资源,将state减1,如果state减为0了,唤醒在队列中排队的Node。 public final boolean release(int arg) { // 核心的释放锁资源方法 if (tryRelease(arg)) { // 释放锁资源释放干净了。 (state == 0) Node h = head; // 如果头节点不为null,并且头节点的状态不为0,唤醒排队的线程 if (h != null && h.waitStatus != 0) // 唤醒线程 unparkSuccessor(h); return true; } // 释放锁成功,但是state != 0 return false; } // 核心的释放锁资源方法 protected final boolean tryRelease(int releases) { // 获取state - 1 int c = getState() - releases; // 如果释放锁的线程不是占用锁的线程,抛异常 if (Thread.currentThread() != getExclusiveOwnerThread()) throw new IllegalMonitorStateException(); // 是否成功的将锁资源释放利索 (state == 0) boolean free = false; if (c == 0) { // 锁资源释放干净。 free = true; // 将占用锁资源的属性设置为null setExclusiveOwnerThread(null); } // 将state赋值 setState(c); // 返回true,代表释放干净了 return free; } // 唤醒节点 private void unparkSuccessor(Node node) { // 拿到头节点状态 int ws = node.waitStatus; // 如果头节点状态小于0,换为0 if (ws < 0) compareAndSetWaitStatus(node, ws, 0); // 拿到当前节点的next Node s = node.next; // 如果s == null ,或者s的状态为1 if (s == null || s.waitStatus > 0) { // next节点不需要唤醒,需要唤醒next的next s = null; // 从尾部往前找,找到状态正常的节点。(小于等于0代表正常状态) for (Node t = tail; t != null && t != node; t = t.prev) if (t.waitStatus <= 0) s = t; } // 经过循环的获取,如果拿到状态正常的节点,并且不为null if (s != null) // 唤醒线程 LockSupport.unpark(s.thread); } 3 使用实例 3.1 公平锁 1.代码: public class ReentrantLockTest { public static void main(String[] args) { ReentrantLock lock = new ReentrantLock(true); new Thread(()->test(lock),"线程A").start(); new Thread(()->test(lock),"线程B").start(); new Thread(()->test(lock),"线程C").start(); } public static void test(ReentrantLock lock){ for (int i = 0; i < 3;i++){ try { lock.lock(); System.out.println(Thread.currentThread().getName()+"获取了锁!"); Thread.sleep(200); }catch (Exception e){ e.printStackTrace(); }finally { lock.unlock(); } } } } 2.执行结果: 3.小结: 公平锁可以保证每个线程获取锁的机会是相等的。 3.2 非公平锁 1.代码: public class ReentrantLockTest { public static void main(String[] args) { ReentrantLock lock = new ReentrantLock(); new Thread(()->test(lock),"线程A").start(); new Thread(()->test(lock),"线程B").start(); new Thread(()->test(lock),"线程C").start(); } public static void test(ReentrantLock lock){ for (int i = 0; i < 3;i++){ try { lock.lock(); System.out.println(Thread.currentThread().getName()+"获取了锁!"); Thread.sleep(200); }catch (Exception e){ e.printStackTrace(); }finally { lock.unlock(); } } } } 2.执行结果: 3.小结: 非公平锁每个线程获取锁的机会是随机的。 3.3 忽略重复操作 1.代码: public class ReentrantLockTest { private ReentrantLock lock = new ReentrantLock(); public void doSomething(){ if(lock.tryLock()){ try { System.out.println(Thread.currentThread().getName()+"获取了锁!"); Thread.sleep(5); }catch (Exception e){ e.printStackTrace(); }finally { lock.unlock(); } } } public static void main(String[] args) throws Exception { ReentrantLockTest test = new ReentrantLockTest(); for (int i = 0; i < 10;i++){ new Thread(()->{test.doSomething();},"线程"+i).start(); Thread.sleep(1); } } } 2.执行结果: 3.小结: 当线程持有锁时,不会重复执行,可以用来防止定时任务重复执行或者页面事件多次触发时不会重复触发。 3.4 超时不执行 1.代码: public class ReentrantLockTest { public static void main(String[] args) { ReentrantLock lock = new ReentrantLock(); new Thread(()->test(lock),"线程A").start(); new Thread(()->test(lock),"线程B").start(); } public static void test(ReentrantLock lock){ try { if(lock.tryLock(2, TimeUnit.SECONDS)){ try { System.out.println(Thread.currentThread().getName()+"获取了锁!"); Thread.sleep(3000); }finally { lock.unlock(); } } }catch (Exception e){ e.printStackTrace(); } } 2.执行结果: 3.小结: 超时不执行可以防止由于资源处理不当长时间占用资源产生的死锁问题。 4 总结 并发是现在软件系统不可避免的问题,ReentrantLock是可重入的独占锁,比起synchronized功能更加丰富,支持公平锁实现,支持中断响应以及限时等待等,是处理并发问题很好的解决方案。 作者:京东物流 陈昌浩 来源:京东云开发者社区

资源下载

更多资源
Mario

Mario

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

腾讯云软件源

腾讯云软件源

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

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文件系统,支持十年生命周期更新。

用户登录
用户注册