商城首页欢迎来到中国正版软件门户

您的位置: 首页 > 文章列表 > 编程开发 > Java并发编程之锁、并发容器、阻塞队列与异步编程实战代码

Java并发编程之锁、并发容器、阻塞队列与异步编程实战代码

  发布于2026-06-30 阅读(0)

扫一扫,手机访问

在多核CPU早已成为标配的今天,并发编程已经不是什么进阶技能,而是后端开发者必须掌握的看家本领了。这篇文章我们系统梳理一下Ja va并发编程中几个最核心的模块:锁机制、并发容器、阻塞队列以及异步编程。不光讲原理,也会配上实战代码,帮大家在高并发场景下做出更明智的技术选型,少踩一些坑。

Ja va并发编程之锁、并发容器、阻塞队列与异步编程实战代码

一、显式锁 ReentrantLock 深度解析

1.1 ReentrantLock 与 synchronized 核心区别

synchronized是Ja va原生的隐式锁,基于JVM实现,加锁解锁全自动;而ReentrantLock是JDK提供的显式锁,基于AQS(AbstractQueuedSynchronizer)框架,需要手动控制。

对比维度

synchronized

ReentrantLock

实现层面

JVM 层面

JDK API 层面

锁释放

自动释放

必须手动 unlock ()

可中断

不支持

支持 lockInterruptibly ()

可超时

不支持

支持 tryLock (timeout)

公平锁

仅非公平

支持公平 / 非公平

条件变量

仅 1 个等待队列

支持多个 Condition

当然,synchronized在JDK6之后做了大量优化,像偏向锁、轻量级锁、自旋锁这些,性能已经大幅提升了。不过,在需要灵活控制锁行为、或者需要更精细的并发控制时,ReentrantLock依然是绕不开的选择。

1.2 可重入性实现原理

可重入锁的意思是:同一个线程可以多次获取同一把锁而不会把自己给卡死。ReentrantLock通过AQS的state状态变量和exclusiveOwnerThread来实现这一点。

具体来说,线程首次获取锁时,state从0变成1,同时记录下当前持有线程是谁。同一线程再次获取锁时,state继续累加;释放时则递减,直到state归零,锁才真正释放。

// ReentrantLock.NonfairSync.tryAcquire()核心逻辑
final boolean nonfairTryAcquire(int acquires) {
    final Thread current = Thread.currentThread();
    int c = getState();
    if (c == 0) {
        // 锁空闲,CAS尝试获取
        if (compareAndSetState(0, acquires)) {
            setExclusiveOwnerThread(current);
            return true;
        }
    }
    else if (current == getExclusiveOwnerThread()) {
        // 同一线程重入,state累加
        int nextc = c + acquires;
        if (nextc < 0) throw new Error("Maximum lock count exceeded");
        setState(nextc);
        return true;
    }
    return false;
}

可重入性避免了递归调用或同一线程反复获取锁导致的死锁,这是它最核心的价值所在。

1.3 可中断锁实现原理

synchronized获取锁时是不可中断的——线程会一直傻等,直到拿到锁为止。ReentrantLock则通过lockInterruptibly()支持中断响应,给阻塞中的线程一个“反悔”的机会。

当调用lockInterruptibly()时,如果线程在等待队列中被中断了,它会直接抛出InterruptedException,不再继续傻等。这为实现超时获取锁、取消阻塞操作提供了基础。

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

// AQS.acquireInterruptibly()
public final void acquireInterruptibly(int arg) throws InterruptedException {
    if (Thread.interrupted())
        throw new InterruptedException();
    if (!tryAcquire(arg))
        doAcquireInterruptibly(arg); // 可中断的等待逻辑
}

1.4 可超时获取锁实现原理

tryLock(long timeout, TimeUnit unit)允许在指定时间内尝试获取锁,超时直接返回false。它的底层基于LockSupport.parkNanos()进行限时等待,同时会检测中断和超时两种信号。

// 超时获取锁使用示例
public boolean tryLockWithTimeout(ReentrantLock lock, long timeoutMs) {
    try {
        return lock.tryLock(timeoutMs, TimeUnit.MILLISECONDS);
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
        return false;
    }
}

这个特性在分布式系统中尤其重要——能有效避免因网络波动导致的线程永久阻塞。

1.5 公平锁与非公平锁实现原理

ReentrantLock默认采用非公平锁,但可以通过构造函数参数指定为公平锁:

ReentrantLock fairLock = new ReentrantLock(true);    // 公平锁
ReentrantLock unfairLock = new ReentrantLock(false); // 非公平锁(默认)

公平锁严格按照线程请求的顺序来分配锁,新线程来了必须乖乖排到队尾;非公平锁则允许新线程在锁释放的那一刻直接入场“抢跑”,不管队列里还有没有人在等。

公平锁的tryAcquire()会额外检查hasQueuedPredecessors(),确保只有等待队列中没有前驱节点时才尝试获取锁。

从性能角度看,非公平锁通常更高(吞吐量能高出约30%),但代价是有可能导致线程饥饿。公平锁保证了公平,但增加了上下文切换的开销。

1.6 实战代码示例

/**
 * ReentrantLock完整使用示例
 */
public class ReentrantLockDemo {
    private final ReentrantLock lock = new ReentrantLock(true);
    private final Condition notFull = lock.newCondition();
    private final Condition notEmpty = lock.newCondition();
    private final Queue queue = new LinkedList<>();
    private static final int CAPACITY = 10;

    public void produce(int data) throws InterruptedException {
        lock.lock();
        try {
            while (queue.size() == CAPACITY) {
                notFull.await(); // 队列满,生产者等待
            }
            queue.offer(data);
            notEmpty.signal(); // 唤醒消费者
        } finally {
            lock.unlock(); // 必须在finally中释放锁
        }
    }

    public int consume() throws InterruptedException {
        lock.lockInterruptibly(); // 可中断方式获取锁
        try {
            while (queue.isEmpty()) {
                notEmpty.await();
            }
            int data = queue.poll();
            notFull.signal();
            return data;
        } finally {
            lock.unlock();
        }
    }
}

几个关键注意点

  • unlock()必须放在finally块里,否则一旦出异常,锁就容易“泄漏”

  • 多个Condition可以精确控制唤醒条件,这是synchronized做不到的

二、ConcurrentHashMap 底层原理与演进

2.1 JDK7 分段锁设计原理

JDK7的ConcurrentHashMap采用\ 分段锁(Segment)\ 设计,整个哈希表被拆成16个Segment数组,每个Segment独立加锁。

核心数据结构长这样:

ConcurrentHashMap
└── Segment[] (默认16个)
    └── HashEntry[] (每个Segment独立的哈希表)
        └── HashEntry链表

Segment继承了ReentrantLock,每次put操作只锁定对应的Segment,其他Segment的读写不受影响。理论上,最大并发度等于Segment的数量(默认16)。

但分段锁有几个明显的缺陷:

  1. 并发度固定,数组扩容后并发能力并不会跟着提升

  2. 每个Segment都需要独立的锁和数据结构,浪费内存

  3. 跨段操作(比如size())需要锁定所有Segment,性能一下就下来了

2.2 JDK8 架构演进:CAS + synchronized

JDK8彻底放弃了分段锁,改用数组+链表+红黑树的结构,和HashMap保持了一致。锁的粒度也从Segment级别细化到了每个哈希桶(Node)级别。

核心改进点:

  • synchronized替换了ReentrantLock,因为JVM对synchronized的优化更成熟

  • 无锁竞争时,直接用CAS做无锁化更新

  • 只有在哈希桶发生冲突时,才对桶的首节点加锁

  • 链表长度超过8时自动转成红黑树,解决哈希碰撞攻击

// JDK8 ConcurrentHashMap.putVal()核心逻辑
final V putVal(K key, V value, boolean onlyIfAbsent) {
    int hash = spread(key.hashCode());
    for (Node[] tab = table;;) {
        int i = (n - 1) & hash;
        Node f = tabAt(tab, i);
        
        if (f == null) {
            // 桶为空,CAS直接插入
            if (casTabAt(tab, i, null, new Node(hash, key, value)))
                break;
        } else {
            synchronized (f) { // 仅锁定桶的首节点
                if (tabAt(tab, i) == f) { // 双重检查
                    // 链表或红黑树插入逻辑
                }
            }
        }
    }
}

2.3 放弃分段锁的深层原因

JDK8放弃分段锁,背后有四个关键考量:

1. 锁粒度更细

  • 分段锁最小粒度是Segment,JDK8最小粒度是哈希桶

  • 并发度随数组容量动态提升,理论最大并发度等于数组长度

2. synchronized 性能优化

  • JDK6之后synchronized引入了偏向锁、轻量级锁、自适应自旋

  • 无锁竞争时,偏向锁的性能甚至优于ReentrantLock

  • JVM可以对synchronized做逃逸分析等深度优化

3. 减少内存开销

  • 消除了Segment对象的内存占用

  • 数据结构与HashMap统一,代码复用性更高

4. 红黑树引入

  • 解决了哈希碰撞导致链表过长的问题

  • 极端情况下,时间复杂度从O(n)降到了O(logn)

2.4 扩容机制深度解析

ConcurrentHashMap的扩容是并发协作式的,支持多线程一起参与搬砖,这是它的一个亮点。

扩容触发条件

  • 元素数量达到阈值(容量×加载因子)

  • 单桶链表长度超过8但数组容量小于64

扩容核心流程

  1. 扩容准备:创建nextTable,容量是原数组的2倍

  2. 扩容标记:sizeCtl设为负数,表明正在扩容中

  3. 任务分配:每个线程负责连续的16个桶的迁移

  4. 并发迁移

    1. 处理完的桶会被标记为ForwardingNode

    2. 遇到ForwardingNode自动跳过,或者主动协助扩容

  5. 扩容完成:table指向nextTable,重置sizeCtl

// 扩容时的ForwardingNode标记
static final class ForwardingNode extends Node {
    final Node[] nextTable;
    ForwardingNode(Node[] tab) {
        super(MOVED, null, null, null);
        this.nextTable = tab;
    }
}

并发扩容的巧妙之处

  • 读操作遇到ForwardingNode会直接转发到新数组,不阻塞

  • 写操作遇到ForwardingNode会主动帮忙扩容,实现了“多线程搭把手”的效果

  • 迁移采用“复制+清除”方式,保证数据一致性

2.5 完整工作流程梳理

以put操作为例,整个流程是这样的:

  1. 计算key的哈希值,定位哈希桶位置

  2. 如果数组没初始化,用CAS触发初始化

  3. 如果目标桶是空的,CAS直接插入新节点

  4. 如果遇到ForwardingNode,帮忙扩容后重试

  5. 否则,对桶的首节点加synchronized锁

  6. 遍历链表或红黑树,key存在就更新,不存在就追加

  7. 链表长度超过8时,触发树化或扩容

  8. 释放锁,CAS更新元素计数,检查是否需要扩容

整个过程中,锁的持有时间非常短,绝大多数操作都是无锁的CAS——这正是ConcurrentHashMap高性能的根本原因。

三、BlockingQueue 阻塞队列实战

3.1 阻塞队列核心作用

BlockingQueue是Ja va并发包里最重要的数据结构之一,专门用来解决生产者-消费者模式下的线程协作问题。

核心特性:

  • 队列满时,生产者线程自动阻塞,直到有消费者来消费

  • 队列空时,消费者线程自动阻塞,直到有生产者生产

  • 所有操作都是线程安全的,内部通过锁和条件变量实现

四种核心操作模式:

操作方式

抛出异常

返回特殊值

阻塞等待

超时等待

插入

add(e)

offer(e)

put(e)

offer(e, time, unit)

移除

remove()

poll()

take()

poll(time, unit)

检查

element()

peek()

-

-

可以说,阻塞队列不仅是线程池的核心组件,也是解耦生产者和消费者速率不匹配的标准方案。

3.2 ArrayBlockingQueue 详解

ArrayBlockingQueue是基于数组实现的有界阻塞队列,创建时必须指定容量。

核心特点:

  • 有界队列,容量固定,不能扩容

  • 单锁双Condition机制(notEmpty + notFull)

  • 支持公平/非公平模式

  • 读写共用同一把锁,并发度相对较低

// ArrayBlockingQueue核心结构
public class ArrayBlockingQueue {
    final Object[] items;      // 存储数组
    int takeIndex;             // 取元素位置
    int putIndex;              // 放元素位置
    int count;                 // 元素数量
    final ReentrantLock lock;  // 单锁
    private final Condition notEmpty;
    private final Condition notFull;
}

适用场景:队列大小可以预估、对内存占用比较敏感的场景。

3.3 LinkedBlockingQueue 详解

LinkedBlockingQueue是基于链表实现的阻塞队列,默认容量为Integer.MAX_VALUE。

核心特点:

  • 双锁分离设计(takeLock + putLock),读写互不阻塞

  • 默认无界(实际最大2^31-1),也可以指定容量

  • 吞吐量高于ArrayBlockingQueue

  • 内存占用相对较高

// LinkedBlockingQueue双锁设计
private final ReentrantLock takeLock = new ReentrantLock();
private final Condition notEmpty = takeLock.newCondition();
private final ReentrantLock putLock = new ReentrantLock();
private final Condition notFull = putLock.newCondition();

需要留意:在无界模式下,如果生产速度远大于消费速度,OOM可不是开玩笑的。

适用场景:生产消费速率差异较大、追求高吞吐量的场景。

3.4 SynchronousQueue 详解

SynchronousQueue是一个不存储元素的阻塞队列——每个插入操作必须等待另一个线程的移除操作。

核心特点:

  • 内部没有缓冲区,队列容量始终为0

  • 支持公平/非公平模式

  • 直接传递,不存储任何元素

  • 是Executors.newCachedThreadPool()的默认队列

// SynchronousQueue典型用法
SynchronousQueue queue = new SynchronousQueue<>();
// 生产者线程
new Thread(() -> {
    queue.put(1); // 会阻塞直到有消费者take
}).start();
// 消费者线程
new Thread(() -> {
    queue.take(); // 会阻塞直到有生产者put
}).start();

适用场景:任务必须立即处理、不允许排队的场景,比如CachedThreadPool。

3.5 DelayQueue 原理与延时任务应用

DelayQueue是支持延时获取元素的无界阻塞队列,元素必须实现Delayed接口。

核心原理:

  • 内部基于PriorityQueue实现,按过期时间排序

  • 只有元素延迟时间到期后才能被取出

  • 队首永远是最早过期的元素

// 延时任务元素定义
public class DelayedTask implements Delayed {
    private final long executeTime;
    private final Runnable task;

    public DelayedTask(long delayMs, Runnable task) {
        this.executeTime = System.currentTimeMillis() + delayMs;
        this.task = task;
    }

    @Override
    public long getDelay(TimeUnit unit) {
        return unit.convert(executeTime - System.currentTimeMillis(), TimeUnit.MILLISECONDS);
    }

    @Override
    public int compareTo(Delayed other) {
        return Long.compare(this.executeTime, ((DelayedTask)other).executeTime);
    }
}

典型应用场景

  1. 订单超时自动取消

  2. 会话超时清理

  3. 定时任务调度

  4. 缓存过期失效

// 延时队列使用示例
DelayQueue delayQueue = new DelayQueue<>();
delayQueue.put(new DelayedTask(30000, () -> System.out.println("30秒后执行")));
delayQueue.put(new DelayedTask(60000, () -> System.out.println("60秒后执行")));

// 消费者线程
while (true) {
    DelayedTask task = delayQueue.take(); // 阻塞直到任务到期
    task.run();
}

3.6 四者对比与选型建议

特性

ArrayBlockingQueue

LinkedBlockingQueue

SynchronousQueue

DelayQueue

容量

有界

可选有界 / 无界

0

无界

数据结构

数组

链表

直接传递

优先级队列

锁机制

单锁

双锁分离

CAS

单锁

公平性

支持

不支持

支持

不支持

吞吐量

中等

极高

典型应用

固定大小线程池

固定大小线程池

缓存线程池

延时任务

选型建议

  • 需要控制队列大小 → ArrayBlockingQueue

  • 追求高吞吐量 → LinkedBlockingQueue

  • 任务必须立即执行 → SynchronousQueue

  • 需要延时执行 → DelayQueue

四、CompletableFuture 异步编程

4.1 异步编程背景

传统Future接口的局限性其实很明显:

  • 无法手动完成

  • 不支持链式调用

  • 不支持异常处理

  • 无法组合多个Future

CompletableFuture在Ja va8中被引入,实现了CompletionStage接口,提供了丰富的异步编程能力,而且支持函数式编程风格,用起来非常灵活。

4.2 核心 API 分类详解

1. 创建类 API

// 使用默认线程池
CompletableFuture future1 = CompletableFuture.runAsync(() -> System.out.println("异步任务"));
CompletableFuture future2 = CompletableFuture.supplyAsync(() -> "返回结果");

// 使用自定义线程池(推荐)
ExecutorService executor = Executors.newFixedThreadPool(10);
CompletableFuture future3 = CompletableFuture.supplyAsync(() -> "自定义线程池", executor);

// 手动完成
CompletableFuture future4 = new CompletableFuture<>();
future4.complete("手动设置结果");
future4.completeExceptionally(new RuntimeException("手动异常"));

需要提醒的是:默认使用ForkJoinPool.commonPool(),所有CompletableFuture共享这个线程池。如果是CPU密集型任务,强烈建议用自定义线程池。

2. 链式转换类 API

CompletableFuture future = CompletableFuture.supplyAsync(() -> "Hello")
    .thenApply(s -> s + " World")           // 同步转换
    .thenApplyAsync(s -> s.toUpperCase())   // 异步转换
    .thenAccept(s -> System.out.println(s)) // 消费结果
    .thenRun(() -> System.out.println("执行完成")); // 仅执行,不消费结果
  • thenApply:输入T,输出钱,类似map操作

  • thenAccept:输入T,无输出,消费型

  • thenRun:无输入无输出,只执行动作

3. 组合类 API

// AND组合:两个都完成才执行
CompletableFuture f1 = CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture f2 = CompletableFuture.supplyAsync(() -> "World");

f1.thenCombine(f2, (s1, s2) -> s1 + " " + s2)
  .thenAccept(System.out::println); // 输出 Hello World

// OR组合:任意一个完成就执行
CompletableFuture fast = f1.applyToEither(f2, s -> s + " faster");

4. 异常处理 API

CompletableFuture.supplyAsync(() -> {
    if (true) throw new RuntimeException("出错了");
    return "正常";
})
.exceptionally(ex -> {
    System.out.println("捕获异常: " + ex.getMessage());
    return "默认值"; // 异常时返回默认值
})
.handle((result, ex) -> {
    if (ex != null) {
        return "处理异常";
    }
    return result;
})
.whenComplete((result, ex) -> {
    // 无论成功失败都会执行,不改变结果
    System.out.println("执行完成");
});

5. 多任务组合

// 所有任务都完成
CompletableFuture all = CompletableFuture.allOf(f1, f2, f3);

// 任意一个任务完成
CompletableFuture any = CompletableFuture.anyOf(f1, f2, f3);

4.3 实战示例:并行调用优化

/**
 * 并行调用多个服务,聚合结果
 */
public class CompletableFutureDemo {
    public UserInfo getUserInfo(Long userId) {
        // 并行调用三个接口
        CompletableFuture userFuture = CompletableFuture.supplyAsync(() -> userService.getUser(userId));
        CompletableFuture> orderFuture = CompletableFuture.supplyAsync(() -> orderService.getOrders(userId));
        CompletableFuture> addrFuture = CompletableFuture.supplyAsync(() -> addressService.getAddresses(userId));

        // 等待所有完成,聚合结果
        return CompletableFuture.allOf(userFuture, orderFuture, addrFuture)
            .thenApply(v -> {
                UserInfo info = new UserInfo();
                info.setUser(userFuture.join());
                info.setOrders(orderFuture.join());
                info.setAddresses(addrFuture.join());
                return info;
            })
            .exceptionally(ex -> {
                log.error("获取用户信息失败", ex);
                return null;
            })
            .join();
    }
}

通过CompletableFuture,原本串行需要3秒的调用,可以优化到1秒内完成——这在微服务架构下是非常常用的优化手段。

结语

这篇文章系统梳理了Ja va并发编程的四大核心模块:ReentrantLock的灵活锁机制、ConcurrentHashMap的高性能并发设计、BlockingQueue的生产者-消费者模式,以及CompletableFuture的异步编程能力。

这些技术是构建高并发系统的基石。建议结合具体的项目去深入实践,重点关注各组件的适用场景和性能特性,避免在生产环境中间出现并发安全问题。如果想继续深入,可以去研究AQS框架、JMM内存模型以及无锁算法,这些才是真正让你对并发编程“豁然开朗”的东西。

本文转载于:https://www.jb51.net/program/3649795y4.htm 如有侵犯,请联系zhengruancom@outlook.com删除。
免责声明:正软商城发布此文仅为传递信息,不代表正软商城认同其观点或证实其描述。