保护共享资源

public class NoLockDemo {
    static int count = 0;

    public static void add() {
        for (int i = 0; i < 10000; i++) {
            count++;
        }
    }

    public static void main(String[] args) throws InterruptedException {
        Thread t1 = new Thread(NoLockDemo::add);
        Thread t2 = new Thread(NoLockDemo::add);
        t1.start();
        t2.start();
        t1.join();
        t2.join();
        System.out.println(count);
    }
}

预期结果:20000 实际结果经常小于 20000,数据发生丢失,就是没有锁保护共享变量。

通过加锁的方式保护共享资源

public class SyncLockDemo {
    static int count = 0;
    // 锁对象
    static final Object lock = new Object();

    public static void add() {
        for (int i = 0; i < 10000; i++) {
            synchronized (lock){
                count++;
            }
        }
    }

    public static void main(String[] args) throws InterruptedException {
        Thread t1 = new Thread(SyncLockDemo::add);
        Thread t2 = new Thread(SyncLockDemo::add);
        t1.start();
        t2.start();
        t1.join();
        t2.join();
        System.out.println(count); // 稳定等于20000
    }
}

通过不加锁的方式实现共享资源保护

import java.util.concurrent.atomic.AtomicInteger;

public class AtomicDemo {
    // 无锁共享变量
    private static AtomicInteger count = new AtomicInteger(0);

    public static void add() {
        for (int i = 0; i < 10000; i++) {
            count.getAndIncrement(); // 原子自增,线程安全
        }
    }

    public static void main(String[] args) throws InterruptedException {
        Thread t1 = new Thread(AtomicDemo::add);
        Thread t2 = new Thread(AtomicDemo::add);
        t1.start();
        t2.start();
        t1.join();
        t2.join();
        System.out.println(count); // 固定 20000
    }
}

获取共享变量时,为了保证该变量的可见性,需要使用 volatile 修饰。

它可以用来修饰成员变量和静态成员变量,他可以避免线程从自己的工作缓存中查找变量的值,必须到主存中获取它的值,线程操作 volatile 变量都是直接操作主存。即一个线程对 volatile 变量的修改,对另一个线程可见。

注意 volatile 仅仅保证了共享变量的可见性,让其它线程能够看到最新值,但不能解决指令交错问题(不能保证原子性)

CAS 必须借助 volatile 才能读取到共享变量的最新值来实现【比较并交换】的效果

CAS的的特点

  • CAS 是基于乐观锁的思想:最乐观的估计,不怕别的线程来修改共享变量,就算改了也没关系,我吃亏点再重试呗。
  • synchronized 是基于悲观锁的思想:最悲观的估计,得防着其它线程来修改共享变量,我上了锁你们都别想改,我改完了解开锁,你们才有机会。
  • CAS 体现的是无锁并发、无阻塞并发,请仔细体会这两句话的意思
    • 因为没有使用 synchronized,所以线程不会陷入阻塞,这是效率提升的因素之一
    • 但如果竞争激烈,可以想到重试必然频繁发生,反而效率会受影响

悲观锁默认一定会出现线程竞争,访问共享资源前先上锁,锁住之后其他线程只能阻塞等待,直到持有锁的线程执行完毕、释放锁,别的线程才可以争抢。

原子整数

J.U.C 并发包提供了:

  • AtomicBoolean
  • AtomicInteger
  • AtomicLong

以 AtomicInteger 为例

AtomicInteger i = new AtomicInteger(0);

// 获取并自增(i = 0,结果 i = 1,返回 0),类似于 i++
System.out.println(i.getAndIncrement());

// 自增并获取(i = 1,结果 i = 2,返回 2),类似于 ++i
System.out.println(i.incrementAndGet());

// 自减并获取(i = 2,结果 i = 1,返回 1),类似于 --i
System.out.println(i.decrementAndGet());

// 获取并自减(i = 1,结果 i = 0,返回 1),类似于 i--
System.out.println(i.getAndDecrement());

// 获取并加值(i = 0,结果 i = 5,返回 0)
System.out.println(i.getAndAdd(5));

底层全部依靠 CAS+volatile 实现线程安全,不会出现 i++ 线程错乱问题。

原子引用类型

  • AtomicReference

底层原理

  1. 内部存放 volatile V value,依靠 volatile 保证引用的可见性
  2. compareAndSet() 使用 CPU‑CAS 指令,无锁完成对象替换;
  3. 失败就自旋重试。

缺陷:ABA 问题

什么是 ABA

  • 线程 1 读到对象 A
  • 线程 2 先把 A 改成 B,又改回 A
  • 线程 1 执行 CAS,发现现在还是 A,更新成功,但中间已经被修改过
// 1. 创建,初始化存放对象
AtomicReference<User> ref = new AtomicReference<>(new User("张三"));

// 2. 获取当前对象
User oldUser = ref.get();

// 3. 直接设置新对象
ref.set(new User("李四"));

// 4. CAS交换:期望值等于当前值,才更新新对象
// compareAndSet(预期值, 新值)
ref.compareAndSet(oldUser,new User("王五"));
  • AtomicMarkableReference

    对象 + boolean 标记,只能标记是否改动过

  • AtomicStampedReference

对象 + int 版本号,每次修改版本 + 1,彻底解决 ABA

import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.atomic.AtomicStampedReference;

@Slf4j
public class Test36 {
    // 参数:初始引用A,初始版本号0
    static AtomicStampedReference<String> ref = new AtomicStampedReference<>("A", 0);

    public static void main(String[] args) throws InterruptedException {
        log.debug("main start...");
        // 获取值 A
        String prev = ref.getReference();
        // 获取版本号
        int stamp = ref.getStamp();
        log.debug("读取到的版本号:{}", stamp);

        other();
        Thread.sleep(1);

        log.debug("main线程持有的旧版本号:{}", stamp);
        // 尝试 A → C,版本号+1
        boolean result = ref.compareAndSet(prev, "C", stamp, stamp + 1);
        log.debug("change A->C 是否成功:{}", result);
    }

    private static void other() throws InterruptedException {
        new Thread(() -> {
            int stamp = ref.getStamp();
            log.debug("t1 当前版本号:{}", stamp);
            // A -> B,版本 +1
            ref.compareAndSet(ref.getReference(), "B", stamp, stamp + 1);
        }, "t1").start();

        Thread.sleep(500);

        new Thread(() -> {
            int stamp = ref.getStamp();
            log.debug("t2 当前版本号:{}", stamp);
            // B -> A,版本 +1
            ref.compareAndSet(ref.getReference(), "A", stamp, stamp + 1);
        }, "t2").start();
    }
}

原子数组

Java 并发包 java.util.concurrent.atomic 提供 4 种常用原子数组:

  1. AtomicIntegerArray:int‑类型原子数组
  2. AtomicLongArray:long‑类型原子数组
  3. AtomicReferenceArray<E>:引用类型原子数组
  4. AtomicBooleanArray:boolean 原子数组

常用 API(以 AtomicIntegerArray 举例)

// 1.初始化,传入数组长度
AtomicIntegerArray arr = new AtomicIntegerArray(5);

// 2.获取下标值
int get(int index)

// 3.设置值
void set(int index,int value)

// 4.原子累加
int getAndAdd(int index,int delta)

// 5.先返回旧值再自增
int getAndIncrement(int index)

// 6.自增后返回新值
int incrementAndGet(int index)

// 7.CAS:期望值相等才更新
boolean compareAndSet(int index,int expect,int update)

import java.util.concurrent.atomic.AtomicIntegerArray;

public class Demo {
    public static void main(String[] args) {
        // 创建长度为3的原子int数组
        AtomicIntegerArray atomicArr = new AtomicIntegerArray(3);

        atomicArr.set(0,10);
        System.out.println(atomicArr.get(0));

        // 下标0元素 +1
        atomicArr.incrementAndGet(0);
        System.out.println(atomicArr.get(0));
    }
}

优点

  1. 无锁 CAS,性能高于 synchronized 加锁数组
  2. 数组元素独立竞争,并发粒度细
  3. 简单易用,不需要手动加锁

缺点

  1. 仅单个下标操作原子;多下标复合操作依旧需要加锁
  2. CAS 高并发下会出现自旋空转,消耗 CPU
  3. 不支持复杂业务逻辑,只适合数值原子更新

和普通数组的对比:

字段更新器

三类更新器

  • AtomicReferenceFieldUpdater // 引用类型域字段更新器
  • AtomicIntegerFieldUpdater //int 类型字段更新器
  • AtomicLongFieldUpdater //long 类型字段更新器

字段更新器用于对普通对象里面指定成员字段实现 CAS‑原子操作,不需要把字段封装成完整 Atomic 对象,节省内存开销。

只能配合volatile使用,否则会出现异常

public class Test40 {
    public static void main(String[] args) {
        Student stu = new Student();

        // 带上泛型,代码更规范
        AtomicReferenceFieldUpdater<Student, String> updater
                = AtomicReferenceFieldUpdater.newUpdater(Student.class, String.class, "name");

        // CAS:name为null则赋值张三
        boolean success = updater.compareAndSet(stu, null, "张三");
        System.out.println("CAS是否成功:" + success);
        System.out.println(stu);
    }
}

class Student {
    // 强制要求 volatile
    volatile String name;

    @Override
    public String toString() {
        return "Student{" +
                "name='" + name + '\'' +
                '}';
    }
}

原子累加器

 java在jdk1.8之后提供了原子累加器,性能极高

public class AtomicAddPerformanceTest {

    // 线程数量
    private static final int THREAD_COUNT = 50;
    // 每个线程循环累加次数
    private static final int CYCLE = 100000;

    public static void main(String[] args) throws InterruptedException {
        // 1.测试 AtomicInteger
        AtomicInteger atomicInteger = new AtomicInteger(0);
        long start1 = System.currentTimeMillis();
        Thread[] threads1 = new Thread[THREAD_COUNT];
        for (int i = 0; i < THREAD_COUNT; i++) {
            threads1[i] = new Thread(() -> {
                for (int j = 0; j < CYCLE; j++) {
                    atomicInteger.getAndIncrement();
                }
            });
            threads1[i].start();
        }
        // 等待所有线程结束
        for (Thread t : threads1) {
            t.join();
        }
        long time1 = System.currentTimeMillis() - start1;
        System.out.println("AtomicInteger 耗时:" + time1 + " ms");
        System.out.println("AtomicInteger 最终数值 = " + atomicInteger.get());


        // 2.测试 LongAdder 分段累加器
        LongAdder longAdder = new LongAdder();
        long start2 = System.currentTimeMillis();
        Thread[] threads2 = new Thread[THREAD_COUNT];
        for (int i = 0; i < THREAD_COUNT; i++) {
            threads2[i] = new Thread(() -> {
                for (int j = 0; j < CYCLE; j++) {
                    longAdder.increment();
                }
            });
            threads2[i].start();
        }
        for (Thread t : threads2) {
            t.join();
        }
        long time2 = System.currentTimeMillis() - start2;
        System.out.println("\nLongAdder 耗时:" + time2 + " ms");
        System.out.println("LongAdder 最终数值 = " + longAdder.sum());
    }
}

性能排序:普通 i++ > LongAdder > AtomicInteger

  • i++:最快,但线程不安全
  • LongAdder:高并发下高性能‑线程安全累加
  • AtomicInteger:CAS 自旋竞争损耗,速度最慢

 AtomicInteger.getAndIncrement()

依靠 CAS 自旋 + volatile 内存屏障;多线程争抢同一个变量,CAS 失败就循环重试,高并发下大量空转消耗 CPU。

LongAdder.increment()

分段 Cell 数组,线程哈希分散到不同单元格竞争,降低冲突; 存在哈希寻址、cell 数组运算开销,就是有竞争的时候利用多个线程去自增变量,最后将所有的变量汇总,因此慢于原生 i++,但是远超 AtomicInteger。

什么是 Unsafe

sun.misc.Unsafe 是 JDK 底层后门类,可以直接操作操作系统内存、CAS、线程调度、对象内存布局,所有原子类、LongAdder、字段更新器底层全部依靠 Unsafe。

获取unsafe

import sun.misc.Unsafe;
import java.lang.reflect.Field;

public class UnsafeDemo {
    public static void main(String[] args) throws Exception{
        //1.通过反射获取私有静态字段theUnsafe
        Field field = Unsafe.class.getDeclaredField("theUnsafe");
        field.setAccessible(true);
        Unsafe unsafe = (Unsafe) field.get(null);
    }
}

  1. Unsafe 为什么不安全?
  • 可以操作堆外内存,内存泄漏
  • 跳过构造方法、绕过访问权限
  • 手动操作内存,容易 JVM 崩溃

AbstractQueuedSynchronizer(AQS)

是 Java 所有阻塞锁、同步工具底层的顶层框架,ReentrantLock、ReentrantReadWriteLock、CountDownLatch、Semaphore 全部基于 AQS 实现。

核心成员‑state 同步状态

state 是一个 volatile int,用来标记锁资源占用状态:

  1. 独占模式(排他锁):例如 ReentrantLock
    • state = 0 → 锁空闲
    • state > 0 → 已经有线程持有锁,可支持可重入,重入一次 state +1
  2. 共享模式:读写锁、信号量、倒计时门闩
    • state 代表可用资源数量,多个线程可以同时获取资源

操作 state 的 3 个底层方法(基于 Unsafe‑CAS)

  1. getState():获取当前同步状态
  2. setState(int newState):设置同步状态
  3. compareAndSetState(int expect,int update):CAS 乐观锁原子修改 state,底层调用 Unsafe 的 CAS

  1. 独占模式(排他):同一时刻只允许一条线程拿到锁资源 对应方法:tryAcquire() 获取锁、tryRelease() 释放锁 实现类:ReentrantLock
  2. 共享模式:允许多个线程同时抢占资源 对应方法:tryAcquireShared() 获取、tryReleaseShared() 释放 实现类:ReentrantReadWriteLock读锁、Semaphore、CountDownLatch

  • 提供了基于 FIFO 的等待队列,类似于 Monitor 的 EntryList
  • 条件变量来实现等待、唤醒机制,支持多个条件变量,类似于 Monitor 的 WaitSet

子类主要实现这样一些方法(默认抛出 UnsupportedOperationException)

  • tryAcquire
  • tryRelease
  • tryAcquireShared

获取锁的姿势

// 如果获取锁失败
if (!tryAcquire(arg)) {
    // 入队,可以选择阻塞当前线程
}

释放锁的姿势

// 如果释放锁成功
if (tryRelease(arg)) {
    // 让阻塞线程恢复运行
}

使用AQS自定义不可重入锁

import java.util.concurrent.locks.AbstractQueuedSynchronizer;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;

/**
 * 自定义 不可重入独占锁
 */
public class NonReentrantLock implements Lock {

    // 内部同步器,继承AQS
    private static class Sync extends AbstractQueuedSynchronizer {

        /**
         * 尝试获取独占锁
         * state=0 空闲 → CAS改成1,上锁成功
         * state=1 被占用 → 获取失败,进入AQS队列阻塞
         */
        @Override
        protected boolean tryAcquire(int arg) {
            // CAS尝试把state从0修改为1
            return compareAndSetState(0, 1);
        }

        /**
         * 释放锁
         * state置回0,唤醒后继线程
         */
        @Override
        protected boolean tryRelease(int arg) {
            // 锁已经释放,抛出异常
            if (getState() == 0) {
                throw new IllegalMonitorStateException();
            }
            setState(0);
            return true;
        }

        // 是否独占锁被当前线程持有
        @Override
        protected boolean isHeldExclusively() {
            return getState() == 1;
        }

        Condition newCondition() {
            return new ConditionObject();
        }
    }

    private final Sync sync = new Sync();

    @Override
    public void lock() {
        sync.acquire(1);
    }

    @Override
    public void lockInterruptibly() throws InterruptedException {
        sync.acquireInterruptibly(1);
    }

    @Override
    public boolean tryLock() {
        return sync.tryAcquire(1);
    }

    @Override
    public boolean tryLock(long timeout, java.util.concurrent.TimeUnit unit) throws InterruptedException {
        return sync.tryAcquireNanos(1, unit.toNanos(timeout));
    }

    @Override
    public void unlock() {
        sync.release(1);
    }

    @Override
    public Condition newCondition() {
        return sync.newCondition();
    }
}

ReentrantLock原理

ReentrantLock 实现了 Lock 接口与序列化接口,内部依靠继承 AQS 的同步器 Sync 作为底层,Sync 衍生出 NonfairSync 非公平锁、FairSync 公平锁两个子类;AQS 内置存放阻塞线程双向链表的 Node 结点以及实现等待通知机制的 ConditionObject,依靠 state 同步状态标记锁的重入次数,非公平锁允许线程直接插队抢锁,公平锁会先检查等待队列是否存在前置线程、遵循先来后到,依靠模板方法模式重写 tryAcquire、tryRelease 完成加锁、释放锁,从而实现可重入的显式锁功能。

非公平锁实现原理

没有竞争时

第一个竞争出现时

复制

Thread-1 执行了

  1. CAS 尝试将 state 由 0 改为 1,结果失败
  2. 进入 tryAcquire 逻辑,这时 state 已经是 1,结果仍然失败
  3. 接下来进入 addWaiter 逻辑,构造 Node 队列
  • 图中黄色三角表示该 Node 的 waitStatus 状态,其中 0 为默认正常状态
  • Node 的创建是懒惰的
  • 其中第一个 Node 称为 Dummy(哑元)或哨兵,用来占位,并不关联线程

Thread-0 释放锁,进入 tryRelease 流程,如果成功

  • 设置 exclusiveOwnerThread 为 null
  • state = 0

如果加锁成功(没有竞争),会设置

  • exclusiveOwnerThread 为 Thread‑1,state = 1
  • head 指向刚刚 Thread‑1 所在的 Node,该 Node 清空 Thread
  • 原本的 head 因为从链表断开,而可被垃圾回收 如果这时候有其它线程来竞争(非公平的体现),例如这时有 Thread‑4 来了

可重入原理

// Sync 继承过来的方法,方便阅读,放在此处
final boolean nonfairTryAcquire(int acquires) {
    final Thread current = Thread.currentThread();
    int c = getState();
    if (c == 0) {
        if (compareAndSetState(0, acquires)) {
            setExclusiveOwnerThread(current);
            return true;
        }
    }
    // 如果已经获得了锁,线程还是当前线程,表示发生了锁重入
    else if (current == getExclusiveOwnerThread()) {
        // state++
        int nextc = c + acquires;
        if (nextc < 0) // overflow
            throw new Error("Maximum lock count exceeded");
        setState(nextc);
        return true;
    }
    return false;
}


protected final boolean tryRelease(int releases) {
    // state--
    int c = getState() - releases;
    if (Thread.currentThread() != getExclusiveOwnerThread())
        throw new IllegalMonitorStateException();
    boolean free = false;
    // 支持锁重入,只有 state 减为 0,才释放成功
    if (c == 0) {
        free = true;
        setExclusiveOwnerThread(null);
    }
    setState(c);
    return free;
}

图解

不可打断原理

private final boolean parkAndCheckInterrupt() {
    // 如果打断标记已经是 true, 则 park 会失效
    LockSupport.park(this);
    // interrupted 会清除打断标记
    return Thread.interrupted();
}


final boolean acquireQueued(final Node node, int arg) {
    boolean failed = true;
    try {
        boolean interrupted = false;
        for (;;) {
            final Node p = node.predecessor();
            if (p == head && tryAcquire(arg)) {
                setHead(node);
                p.next = null;
                failed = false;
                // 还是需要获得锁后,才能返回打断状态
                return interrupted;
            }
            if (shouldParkAfterFailedAcquire(p, node) &&
                    parkAndCheckInterrupt())
            {
                // 如果是因为 interrupt 被唤醒,返回打断状态为 true
                interrupted = true;
            }
        }
    } finally {
        if (failed)
            cancelAcquire(node);
    }
}


public final void acquire(int arg) {
    if (!tryAcquire(arg) &&
        acquireQueued(addWaiter(Node.EXCLUSIVE), arg))
    {
        // 如果打断状态为 true
        selfInterrupt();
    }

}


static void selfInterrupt() {
    // 重新产生一次中断
    Thread.currentThread().interrupt();
}

流程图

可打断模式的原理

private void doAcquireInterruptibly(int arg) throws InterruptedException {
    final Node node = addWaiter(Node.EXCLUSIVE);
    boolean failed = true;
    try {
        for (;;) {
            final Node p = node.predecessor();
            if (p == head && tryAcquire(arg)) {
                setHead(node);
                p.next = null; // help GC
                failed = false;
                return;
            }
            if (shouldParkAfterFailedAcquire(p, node) &&
                    parkAndCheckInterrupt()) {
                // 在 park 过程中如果被 interrupt 会进入此
                // 这时候抛出异常, 而不会再次进入 for (;;)
                throw new InterruptedException();
            }
        }
    } finally {
        if (failed)
            cancelAcquire(node);
    }
}

流程图

Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐