ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

AQS之Semaphore详解

AQS之Semaphore详解 AQS之Semaphore详解一、Semaphore类的继承关系1. AbstractQueuedSynchronizer:提供了一个同步器的框架。2. SyncSemaphore的内部类提供了锁的具体实现。3. FairSyncSync的子类实现公平锁。4. NonfairSync:Sync的子类实现非公平锁。二、Semaphore的基本使用1. 使用场景2. 代码实现3. 运行结果4. 案例分析三、Semaphore的优缺点1. 简介2. 优点2.1、灵活性2.2、可重入性2.3、公平性2.4、可以用于多种场景3. 优点3.1、使用复杂3.2、容易造成死锁3.3、不支持条件等待3.4、可能出现饥饿3.5 小结四、源码分析1. 构造方法1.1 Semaphore(int permits)1.2 Semaphore(int permits, boolean fair)1.3 构造方法小结2. acquire方法2.1 AQS#acquireSharedInterruptibly2.2 Semaphore#tryAcquireShared2.3 AQS#doAcquireSharedInterruptibly2.4 AQS#doAcquireSharedInterruptibly#addWaiter2.5 AQS#addWaiter#enq2.6 AQS#shouldParkAfterFailedAcquire2.7 AQS#parkAndCheckInterrupt2.8 AQS#LockSupport.park2.9 流程图3. release方法3.1 AQS#releaseShared3.2 Semaphore#tryReleaseShared3.3 AQS#doReleaseShared3.4 AQS#doReleaseShared#unparkSuccessor3.5 AQS#LockSupport.unpark3.6 AQS#doReleaseShared#setHeadAndPropagate3.7 AQS#doReleaseShared#setHeadAndPropagate#isShared3.8 流程图一、Semaphore类的继承关系1. AbstractQueuedSynchronizer:提供了一个同步器的框架。AQS为Semaphore提供了基本的针对共享资源的获取失败入队出队阻塞唤醒的逻辑。Semaphore通过AQS的同步状态来表示可用的许可数并通过AQS的等待队列来管理等待获取许可的线程。当一个线程请求获取许可时如果许可数不足则该线程会被阻塞并加入到AQS的等待队列中。当有其他线程释放许可时AQS会从等待队列中选择一个线程唤醒使其重新尝试获取许可。这样就实现了对共享资源的获取失败入队出队阻塞唤醒的逻辑。2. SyncSemaphore的内部类提供了锁的具体实现。Sync是Semaphore的内部类它继承自AQS并重写了其中的方法用于实现Semaphore的同步逻辑。3. FairSyncSync的子类实现公平锁。FairSync是Sync的子类它实现了公平的获取许可的机制。当线程请求许可时FairSync会按照FIFO的顺序选择等待的线程获取许可。4. NonfairSync:Sync的子类实现非公平锁。NoFairSync是Sync的子类它实现了非公平的获取许可的机制。当线程请求许可时NoFairSync会直接尝试获取许可而不管是否有其他线程在等待。二、Semaphore的基本使用1. 使用场景假设有一个公共游泳池最多允许同时有5个人在游泳其他人需要等待。这就可以使用Semaphore来控制并发访问的线程数。2. 代码实现importjava.util.concurrent.Semaphore;publicclassSwimmingPool{privatestaticfinalintMAX_SWIMMERS5;privatestaticfinalSemaphoresemaphorenewSemaphore(MAX_SWIMMERS);publicstaticvoidmain(String[]args){for(inti1;i10;i){ThreadswimmernewThread(newSwimmer(i));swimmer.start();}}staticclassSwimmerimplementsRunnable{privateintid;publicSwimmer(intid){this.idid;}Overridepublicvoidrun(){try{semaphore.acquire();// 请求获取许可如果没有可用许可则阻塞等待System.out.println(Swimmer id starts swimming.);Thread.sleep(10000);// 模拟游泳过程semaphore.release();// 释放许可System.out.println(Swimmer id finishes swimming.);}catch(InterruptedExceptione){e.printStackTrace();}}}}3. 运行结果4. 案例分析在这个案例中有10个游泳者想要在游泳池中游泳但是游泳池最多允许同时有5个人在游泳。Semaphore的初始许可数为5当有游泳者调用acquire()方法请求获取许可时如果许可数不足则该游泳者会被阻塞等待。当有其他游泳者完成游泳并调用release()方法释放许可时Semaphore会选择一个等待的游泳者唤醒使其开始游泳。通过Semaphore的控制保证了同时游泳的人数不超过5个。三、Semaphore的优缺点1. 简介Semaphore是一种并发控制工具用于限制同时访问某个资源或执行某个任务的线程数量。它在多线程环境中起到了线程同步和互斥的作用。2. 优点2.1、灵活性Semaphore可以根据需要设置初始许可数允许多个线程同时访问共享资源或者限制同时执行的任务数量。2.2、可重入性Semaphore是可重入的同一个线程可以多次获取和释放许可。2.3、公平性Semaphore可以选择是否公平地分配许可。如果设置为公平模式那么等待时间最长的线程将优先获取许可。2.4、可以用于多种场景Semaphore可以用于解决生产者-消费者问题、连接池管理、并发线程数控制等多种并发场景。3. 优点3.1、使用复杂相对于其他线程同步和互斥的工具Semaphore的使用相对复杂需要手动调用acquire()和release()方法进行许可的获取和释放容易出现逻辑错误。3.2、容易造成死锁如果在使用Semaphore时没有正确地释放许可可能会导致线程间的死锁情况造成程序无法继续执行。3.3、不支持条件等待与ReentrantLock相比Semaphore不支持线程的条件等待和通知无法使用wait()和notify()方法进行线程间的通信。3.4、可能出现饥饿在公平模式下许可的获取是按照线程等待的先后顺序进行的这可能导致某些线程一直无法获取到许可出现饥饿现象。3.5 小结综上所述Semaphore是一种功能强大的并发控制工具能够灵活地管理共享资源的访问和任务的执行。然而由于其使用复杂、可能导致死锁、不支持条件等待以及可能出现饥饿等缺点开发人员在使用Semaphore时需要仔细考虑和处理这些问题。四、源码分析1. 构造方法publicSemaphore(intpermits){syncnewNonfairSync(permits);}publicSemaphore(intpermits,booleanfair){syncfair?newFairSync(permits):newNonfairSync(permits);}1.1 Semaphore(int permits)参数permits表示同时可以访问的线程数目。如果permits的值为1表示Semaphore对象可以用作互斥锁。当permits的值大于1时Semaphore对象可以用作资源池的控制器限制可以访问资源的线程的数量。1.2 Semaphore(int permits, boolean fair)参数fair表示是否采用公平锁策略。如果为true表示Semaphore对象采用公平锁策略即先进入等待队列的线程将先获得许可如果为false表示Semaphore对象采用非公平锁策略线程获取许可的顺序是不确定的。1.3 构造方法小结使用Semaphore的构造方法可以创建一个具有指定许可数目和锁策略的Semaphore对象。根据不同的应用场景可以选择适合的构造方法来创建Semaphore对象然后使用acquire()和release()方法来获取和释放许可。2. acquire方法外部调用加锁方法acquire()(支持可中断)中断概念如果不清楚的话可以参考Java之线程中断publicvoidacquire()throwsInterruptedException{sync.acquireSharedInterruptibly(1);}2.1 AQS#acquireSharedInterruptiblyacquire方法内部会调用acquireSharedInterruptibly方法可以看到AQS给我们提供了模板方法 tryAcquireSharedpublicfinalvoidacquireSharedInterruptibly(intarg)throwsInterruptedException{if(Thread.interrupted())thrownewInterruptedException();if(tryAcquireShared(arg)0)doAcquireSharedInterruptibly(arg);}注意这里没有实现加锁的逻辑模板方法留给子类实现2.2 Semaphore#tryAcquireShared调用子类Semaphore的tryAcquireShared方法这里以非公平为例protectedinttryAcquireShared(intacquires){returnnonfairTryAcquireShared(acquires);}//非公平finalintnonfairTryAcquireShared(intacquires){for(;;){intavailablegetState();intremainingavailable-acquires;if(remaining0||compareAndSetState(available,remaining))returnremaining;}}//公平protectedinttryAcquireShared(intacquires){for(;;){if(hasQueuedPredecessors())return-1;intavailablegetState();intremainingavailable-acquires;if(remaining0||compareAndSetState(available,remaining))returnremaining;}}获取剩余资源许可int available getState();得到剩余许可remaining小于0直接返回否则说明有许可CAS操作保证线程安全获取锁CAS成功则获取锁执行业务代码。如果没有获取锁tryAcquireShared(arg) 0回到acquireSharedInterruptibly方法publicfinalvoidacquireSharedInterruptibly(intarg)throwsInterruptedException{if(Thread.interrupted())thrownewInterruptedException();if(tryAcquireShared(arg)0)doAcquireSharedInterruptibly(arg);}2.3 AQS#doAcquireSharedInterruptibly没有可用资源只能调用doAcquireSharedInterruptibly走入队逻辑privatevoiddoAcquireSharedInterruptibly(intarg)throwsInterruptedException{finalNodenodeaddWaiter(Node.SHARED);booleanfailedtrue;try{for(;;){finalNodepnode.predecessor();if(phead){intrtryAcquireShared(arg);if(r0){setHeadAndPropagate(node,r);p.nextnull;// help GCfailedfalse;return;}}if(shouldParkAfterFailedAcquire(p,node)parkAndCheckInterrupt())thrownewInterruptedException();}}finally{if(failed)cancelAcquire(node);}}2.4 AQS#doAcquireSharedInterruptibly#addWaiteraddWaiter方法 添加到等待队列构造Node结点第一个入队的线程需要调用enq(node);创建队列privateNodeaddWaiter(Nodemode){NodenodenewNode(Thread.currentThread(),mode);// Try the fast path of enq; backup to full enq on failureNodepredtail;if(pred!null){node.prevpred;if(compareAndSetTail(pred,node)){pred.nextnode;returnnode;}}enq(node);returnnode;}创建一个需要入队的结点队列不为空的情况下CAS操作将其设置为tail尾结点。2.5 AQS#addWaiter#enqprivateNodeenq(finalNodenode){for(;;){Nodettail;if(tnull){// Must initializeif(compareAndSetHead(newNode()))tailhead;}else{node.prevt;if(compareAndSetTail(t,node)){t.nextnode;returnt;}}}}队列没有一个结点那么第一个线程需要构建队列。开始入队尾插法2.1 node.prev t; 修改当前结点的前驱指针2.2 compareAndSetTail(t, node) 将当前结点CAS置为tail尾结点2.3 t.next node;修改当前结点的后继指针addWaiter方法结束入队完成。final Node p node.predecessor();获取当前结点的前驱结点如果前驱结点是head,又会调用tryAcquireShared子类尝试获取资源如果没有可用许可调用shouldParkAfterFailedAcquire准备阻塞2.6 AQS#shouldParkAfterFailedAcquireprivatestaticbooleanshouldParkAfterFailedAcquire(Nodepred,Nodenode){intwspred.waitStatus;if(wsNode.SIGNAL)/* * This node has already set status asking a release * to signal it, so it can safely park. */returntrue;if(ws0){/* * Predecessor was cancelled. Skip over predecessors and * indicate retry. */do{node.prevpredpred.prev;}while(pred.waitStatus0);pred.nextnode;}else{/* * waitStatus must be 0 or PROPAGATE. Indicate that we * need a signal, but dont park yet. Caller will need to * retry to make sure it cannot acquire before parking. */compareAndSetWaitStatus(pred,ws,Node.SIGNAL);}returnfalse;}准备阻塞将当前结点的前驱结点pred.waitStatusCAS置为-1SIGNAL可唤醒staticfinalintSIGNAL-1;-1 表示后继的线程结点需要被唤醒2.7 AQS#parkAndCheckInterrupt调用parkAndCheckInterrupt()方法privatefinalbooleanparkAndCheckInterrupt(){LockSupport.park(this);returnThread.interrupted();}2.8 AQS#LockSupport.park最后一步LockSupport.park(this)真正阻塞将当前线程挂起。2.9 流程图3. release方法外部调用释放锁方法release()方法publicvoidrelease(){sync.releaseShared(1);}3.1 AQS#releaseShared通过内部类sync调用AQS模板方法releaseShared()方法publicfinalbooleanreleaseShared(intarg){if(tryReleaseShared(arg)){doReleaseShared();returntrue;}returnfalse;}tryReleaseShared方法AQS内部没有实现留给子类实现。3.2 Semaphore#tryReleaseShared来到Semaphore类tryReleaseShared 子类提供具体实现释放锁的逻辑。protectedfinalbooleantryReleaseShared(intreleases){for(;;){intcurrentgetState();intnextcurrentreleases;if(nextcurrent)// overflowthrownewError(Maximum permit count exceeded);if(compareAndSetState(current,next))returntrue;}}tryReleaseShared大致逻辑获取当前state变量值把资源的许可数量1通过CASfor循环保证操作更新成功。tryReleaseShared结束释放资源3.3 AQS#doReleaseSharedtryReleaseShared结束来到AQS的doReleaseShared方法。调用doReleaseShared唤醒阻塞队列线程。privatevoiddoReleaseShared(){for(;;){Nodehhead;if(h!nullh!tail){intwsh.waitStatus;if(wsNode.SIGNAL){if(!compareAndSetWaitStatus(h,Node.SIGNAL,0))continue;// loop to recheck casesunparkSuccessor(h);}elseif(ws0!compareAndSetWaitStatus(h,0,Node.PROPAGATE))continue;// loop on failed CAS}if(hhead)// loop if head changedbreak;}}如果头结点head等于Node.SIGNALCAS置为0唤醒前的准备工作。3.4 AQS#doReleaseShared#unparkSuccessordoReleaseShared内部调用unparkSuccessor方法privatevoidunparkSuccessor(Nodenode){intwsnode.waitStatus;if(ws0)compareAndSetWaitStatus(node,ws,0);Nodesnode.next;if(snull||s.waitStatus0){snull;for(Nodettail;t!nullt!node;tt.prev)if(t.waitStatus0)st;}if(s!null)LockSupport.unpark(s.thread);}获取到Node s node.next;头结点的下一个结点要唤醒的结点3.5 AQS#LockSupport.unpark执行LockSupport.unpark(s.thread)方法唤醒阻塞线程3.6 AQS#doReleaseShared#setHeadAndPropagate紧接着被唤醒的线程会来到doAcquireSharedInterruptibly方法的for循环里刚好等于头结点可以尝试获取资源if(phead){intrtryAcquireShared(arg);if(r0){//semaphore用不到这里setHeadAndPropagate(node,r);p.nextnull;// help GCfailedfalse;return;}}获取head结点将当前结点设置为head结点释放之前的头结点让gc回收。3.7 AQS#doReleaseShared#setHeadAndPropagate#isShared到这里还没完setHeadAndPropagate方法里会调用isShared() 判断当前链表是否是共享模式privatevoidsetHeadAndPropagate(Nodenode,intpropagate){Nodehhead;// Record old head for check belowsetHead(node);if(propagate0||hnull||h.waitStatus0||(hhead)null||h.waitStatus0){Nodesnode.next;if(snull||s.isShared())doReleaseShared();}}如果前驱节点正好是head,就可以尝试获取资源获取成功就可以执行业务逻辑只要资源数state0充足可以一直唤醒。3.8 流程图
返回列表