并发容器与数据交接
并发容器主要服务三类操作:按键查找共享对象、遍历一组元素、在线程之间交接任务。它们对“读到哪一版数据”和“容量耗尽怎么办”有不同约定。
共享数据
├─ 映射:ConcurrentHashMap、ConcurrentSkipListMap
├─ 集合:并发 key set、CopyOnWriteArrayList / Set
└─ 交接:BlockingQueue、ConcurrentLinkedQueue、TransferQueue选择时先确定业务操作是否跨多个键,读取是否需要固定快照,以及生产速度超过消费速度时是否允许等待。容器的内部同步只能覆盖其提供的操作,元素内部的可变字段仍要另行保护。
按键创建与更新共享对象
建立一个按需创建的注册表
Registry 以租户标识为键,第一次查找时生成字符串客户端标识,再次查找使用既有值:
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
public final class Registry {
public static void main(String[] args) {
ConcurrentHashMap<String, String> clients = new ConcurrentHashMap<>();
AtomicInteger creates = new AtomicInteger();
String a = clients.computeIfAbsent("tenant-42", key -> {
creates.incrementAndGet();
return "client-for-" + key;
});
String b = clients.computeIfAbsent("tenant-42", key -> "unused");
if (a != b || creates.get() != 1) throw new AssertionError("mapping");
System.out.println("client=client-for-tenant-42 creates=1");
}
}在 Linux 普通用户可写目录保存为 Registry.java。需要完整 JDK,安装与选择见Java 版本基线。设置实际 JAVA_HOME;实验运行于 Temurin 25.0.4+7,兼容 Java 17。
export JAVA_HOME=/opt/jdk-25
unset JAVA_TOOL_OPTIONS JDK_JAVA_OPTIONS _JAVA_OPTIONS
LAB_OUT=$(mktemp -d /tmp/container-learning.XXXXXX)
"$JAVA_HOME/bin/javac" --release 17 -Xlint:all -Werror -d "$LAB_OUT" Registry.java
"$JAVA_HOME/bin/java" -cp "$LAB_OUT" Registry输出 client=client-for-tenant-42 creates=1。这里用字符串避免真正创建连接资源;组件接入真实客户端时,还要负责替换、关闭和失败重试。若编译失败,先检查文件名、完整 JDK 和 release 目标,不能用旧 class 的输出代替本次编译。
将判断与修改放入同一次操作
两个线程都执行“containsKey 为 false 后创建并 put”,可能各创建一次,最终表里却只剩一个值。线程安全的单次 containsKey 与 put 无法保护中间的空隙。
| 操作 | 原子判断与结果 |
|---|---|
| putIfAbsent(key, value) | 键不存在才放入,返回原值或 null |
| remove(key, expected) | 当前值匹配 expected 才删除 |
| replace(key, expected, next) | 当前值匹配才替换 |
| computeIfAbsent | 缺失时计算,非 null 结果才建立映射 |
| compute / computeIfPresent | 根据当前映射重新计算,返回 null 可删除 |
| merge | 缺失时放入给定值,存在时合并;合并为 null 则删除 |
契约见 ConcurrentHashMap API。条件值比较按这些 Map 方法的相等语义执行,不要直接套用 AtomicReference 的引用身份比较。
ConcurrentHashMap 的 computeIfAbsent 对一次调用按其原子契约计算;成功值保留且没有删除时,并发同键调用能够复用它。返回 null、抛异常或后来删除映射都会让后续调用有机会再次计算。不能把它理解为某个 key 在整个进程生命周期中最多执行一次函数。ConcurrentMap 默认实现与不同具体类的函数重试规则也有差别,详见 ConcurrentMap 接口。
下载并发容器实验包,解压进入 concurrent-containers 目录,沿用 LAB_OUT:
"$JAVA_HOME/bin/javac" --release 17 -Xlint:all -Werror -d "$LAB_OUT" *.java
"$JAVA_HOME/bin/java" -cp "$LAB_OUT" ContainersLab创建对照输出:
check-then-put creates=2 size=1
compute creates=1 size=1错误模式用屏障固定“两者都已发现缺失”的时刻。正确模式让两个调用共享 computeIfAbsent,映射函数不放屏障,避免一个计算等待另一个被自己阻挡的计算。
回调失败与值内部的共享
映射函数应保持短小,不递归修改同一张表,也不执行长时间网络操作。实现可以检测部分无法完成的递归更新并抛 IllegalStateException;更复杂的跨键或跨组件等待仍可能形成死锁,不能依赖递归检测代替设计。
实验构造同键递归、返回 null 和抛异常三种情况:
recursive-update=rejected null-not-cached=true failed-load-retried=true异常会传播给调用方,失败键仍无映射,后续成功加载才写入。若值是 CompletableFuture,失败 Future 本身却是一个正常非 null 值,容器会保留它;是否移除失败占位必须由应用决定。
map.computeIfAbsent(key, ignored -> new ArrayList<>()).add(value);这段代码只保证取得列表映射,后面的 ArrayList.add 仍可能被多个线程并发调用。可改为不可变列表在 compute 内整体替换,或给值本身提供并发访问协议。发布对象之前的初始化与读者取得该映射形成可见性关系,发布后的继续修改则需要新的同步。
ConcurrentHashMap 的桶、扩容与视图
一次访问如何找到节点
固定 OpenJDK jdk-25+36 ConcurrentHashMap使用 table 保存桶入口。哈希值定位桶后,读取沿节点或树结构查键;空桶插入可通过 CAS 安装,已有桶的更新需要桶内协调。计算映射时还可能使用 ReservationNode 表示计算占位。
树化用于缓解较长冲突链;表较小时实现可能优先扩容。树化阈值、树退化和迁移细节属于具体版本实现,应用应先确保 key 的 equals/hashCode 稳定、分布合理。把参与 hashCode 的字段在入表后修改,会让查找和删除失去可靠位置。
sizeCtl 在初始化、正常容量阈值和扩容协作中含义不同;nextTable 指向迁移目标,transferIndex 分配迁移区间,旧桶上的 ForwardingNode 引导访问新表。遇到迁移的线程可能协助搬迁,扩容成本由业务访问共同承担。
hash(key) → table 桶
├─ 空桶:尝试安装
├─ 普通节点 / 树:查找或桶内更新
└─ 转发节点:进入新表,可能协助迁移合理 initialCapacity 减少高峰扩容,但不会限制最大条目数。现代 concurrencyLevel 只是构造时的容量提示,不是 JDK 7 Segment 数量。热点 key 和碰撞桶仍可能集中竞争;增大整张表不一定能改善同一个键的慢回调。
遍历与整体快照
ConcurrentHashMap 的迭代器弱一致,可以与更新并行,不抛普通结构修改导致的 ConcurrentModificationException,但也不冻结整张表。size、isEmpty 和批量汇总在并发更新时适合观察趋势,不能据此保证多个键同时满足业务条件。
把正在变化的 map 拷贝到另一个容器,可以得到一个后续不再变化的副本,却不保证这些条目来自原表的同一瞬间。需要一致配置版本时,可在更新侧构造完整不可变配置,再通过原子引用切换;跨键交易则使用共同锁或持久化事务。
bulk 操作可能利用 commonPool 并行执行,阈值决定是否尝试并行。回调应短小且能接受并发视图,不能在汇总函数里等待依赖同一执行资源的远端任务。
按键排序、范围查询可考虑 ConcurrentSkipListMap,但它的范围视图与遍历同样不构成跨键事务。只需并发集合时,可用 ConcurrentHashMap.newKeySet;同步包装集合则要求遍历期间按包装器规定持锁,不能套用并发迭代器的使用方式。包装器约定见 Collections API。
CopyOnWrite 保留哪一版数组
CopyOnWriteArrayList 的写者在锁内复制数组、修改副本并发布新数组引用。迭代器保存创建时的数组,因此后续增删不会改变其遍历内容。它适合元素不多、修改很少、遍历频繁的监听器集合。
ContainersLab 先取得 v1 的迭代器,再把当前列表更新为 v2:
iterator=v1 current=v2 iterator-remove=rejected旧迭代器和当前列表分别持有不同数组。迭代器 remove 不支持,抛 UnsupportedOperationException。这个固定数组快照与 ConcurrentHashMap 的弱一致遍历是两种语义,不能混称。规则见 CopyOnWriteArrayList API。
旧 iterator ──→ 数组版本 1:[v1]
当前 list ──→ 数组版本 2:[v2]快照固定的是元素引用序列,不是元素对象的深拷贝。监听器对象内部继续变化时,仍需自己的同步。删除监听器也不会使已经创建的快照停止调用它;若注销必须等待在途回调结束,需要状态标记、在途计数和退出等待协议。
每次写入复制与元素数量同阶的数组。高频小修改可以合并成批量替换,避免连续复制;长期持有迭代器会延长旧数组及其元素引用的存活。诊断内存增长时,既看写入分配,也看是谁保留旧快照。CopyOnWriteArraySet 使用类似读多写少思路,不能拿来承载高频去重队列。
队列的交接与容量
空和满各有四种处理方式
| 处理方式 | 插入 | 取出 | 调用方承担的动作 |
|---|---|---|---|
| 抛异常 | add | remove | 处理已满或为空的异常 |
| 立即特殊值 | offer | poll | 处理 false 或 null |
| 可中断等待 | put | take | 提供取消和停止协议 |
| 限时等待 | offer(time, unit) | poll(time, unit) | 耗尽预算后返回失败 |
这是 BlockingQueue API定义的操作分组。队列拒绝 null,以便 poll 的 null 表示没有取到元素。成功交接保证放入之前的动作对后续访问或移除该元素的线程可见;生产者交接后继续修改同一对象,会再次产生并发访问风险。不可变任务可以简化所有权。
实验先填满容量为 1 的 ArrayBlockingQueue,再限时 offer 第二个元素;取走第一个后限时 poll 空队列。随后只向无界队列放入 40 个元素并排空:
bounded-full=rejected empty=timeout unbounded-retained=40 drained=true40 是有限演示规模,不通过耗尽内存验证“无界”。无界表示 API 没有设定元素容量上限,实际仍受可用内存约束。
常见队列如何放置压力
| 类型 | 存储与等待特点 | 常见用途与限制 |
|---|---|---|
| ArrayBlockingQueue | 固定数组容量,锁与条件协调 | 明确限额的生产消费 |
| LinkedBlockingQueue / Deque | 链式节点,可指定容量 | 未显式指定时默认容量很大,仍需主动设限 |
| SynchronousQueue | 不存储元素,交接双方配对 | 直接移交,offer 可能立即失败 |
| PriorityBlockingQueue | 优先队列,逻辑无界 | 同优先级顺序需自行编码,低优先任务可能饥饿 |
| DelayQueue | 到期元素才可取,逻辑无界 | 内存延迟任务,不能替代持久调度 |
| LinkedTransferQueue | 无界链式交接,transfer 可等接收 | 已接收不表示任务业务执行完成 |
| ConcurrentLinkedQueue / Deque | 非阻塞节点操作,无容量背压 | 必须在外层约束积压 |
队列类入口与一致性概述见 java.util.concurrent 包说明。ConcurrentLinkedQueue 的 size 需要遍历,并发变化时也不适合作硬上限判断,详见其 API。先 size 再 offer 会让多个生产者同时通过检查,真正的限额应由有界队列或共同许可协议执行。
ArrayBlockingQueue 的一把锁协调读写条件;LinkedBlockingQueue 的实现分离入队与出队锁,并通过计数与信号配合。这影响竞争形态,却无法单凭锁数量决定实际性能。测量应保持元素大小、生产消费比例、CPU 配额和队列容量一致。
队列长度之外还要看任务年龄
排队会消耗请求的剩余时间。任务应保存入队时刻或单调时钟截止预算,消费者取出后先判断是否仍值得执行;过期任务要向调用方完成失败或按业务协议补偿,不能静默消失。
在稳定吞吐条件下,平均在队数量约等于到达率乘平均等待时间。它可以辅助容量估算,但突发流量、服务时间长尾和不同任务大小会使实际峰值更高。容量至少同时考虑允许等待时长与任务保留内存,不能只凭“放一万个”指定缓冲。
把慢消费者前的队列无限拉长,会把过载推迟成延迟和堆增长。使用 offer 失败后,应返回过载、降级或进入持久化流程;请求线程无限 put 则把等待移到入口线程。线程池的具体入队与拒绝顺序见线程池与拒绝策略。
创建、替换与停止时的资源处理
ConcurrentHashMap 没有自动 TTL、淘汰或资源 close。按租户缓存真实客户端时,组件应规定最大规模,明确“创建中、可用、失败、关闭”的状态以及替换后旧客户端的使用者何时退出。
慢加载可以先用 putIfAbsent 安装 CompletableFuture 占位,由获胜者在独立、受控的执行器加载。提交被拒绝或加载失败时,先异常完成 Future,再使用 remove(key, sameFuture) 条件删除,避免误删后继的新值。已经创建但未被采用的资源要关闭;关闭异常不能掩盖最初加载错误。等待者自己的超时不应直接取消所有共享者仍需要的加载,任务取消与 Future 结果语义见异步上下文与 Future。
保留失败占位可以短期抑制重试风暴,但需要明确期限和再次加载规则。需要按时间、权重、刷新等完整缓存功能时,应选择具备相应契约的缓存实现,而非在表上不断补定时扫描。
BlockingQueue 也没有统一 close 方法。停止流程先拒绝新生产,决定已接受任务是排空还是返回失败,再通知或中断消费者,最后等待任务退出。结束消息要考虑消费者数量、队列已满和多生产者顺序;直接 clear 会丢掉任务及其回执,不能作为无条件关停办法。
从异常与积压定位容器问题
| 现象 | 首先核对 | 处理与复验 |
|---|---|---|
| 客户端重复创建,map 里却只有一个 | contains+put、加载失败与删除历史 | 原子安装或占位合并,再统计创建与关闭数量 |
| compute 长期阻挡更新 | 回调栈、同键或碰撞桶、外部等待 | 移出慢加载,复验同负载等待与失败传播 |
| 遍历结果新旧混合 | 容器是否弱一致,是否要求同一版本 | 构造整体快照或共同事务 |
| COW 分配高、旧数组不回收 | 写频率、集合规模、迭代器持有者 | 合并更新、缩短快照生命周期,再看分配与保留对象 |
| 队列长度和最老年龄持续增加 | 生产速率、完成速率、下游耗时 | 限制入口、处理慢消费、丢弃过期任务并完成回执 |
| 元素数量稳定但堆很大 | 单元素捕获对象和资源引用 | 缩小任务载荷、释放完成后引用 |
线程和内存采集入口见并发诊断。map.size、队列容量只能描述容器中的元素;一个 Future 可能间接保留完整请求体,一个监听器快照可能保留已注销组件。根据实际引用链解释内存,而不是只看集合类名。
运行全部模式
在解压目录设置 JAVA_HOME 后执行 bash run.sh,严格编译两份源文件,两个 Java 子进程各有 20 秒外部上限,最后打印 PASS containers-lab。也可以在有 Docker 权限的 Linux 主机上运行:
LAB_SOURCE="$(pwd -P)"
JDK_IMAGE=eclipse-temurin:25.0.4_7-jdk
docker pull "$JDK_IMAGE"
docker run --pull never --rm --network none --read-only \
--user 10001:10001 --cap-drop ALL --security-opt no-new-privileges \
--mount "type=bind,src=$LAB_SOURCE,dst=/lab,readonly" \
--tmpfs /tmp:rw,nosuid,nodev,size=128m,mode=1777 \
"$JDK_IMAGE" bash /lab/run.sh容器普通 UID 读取源码,在 /tmp 编译,不修改宿主示例。离线环境提前导入经校验的可信镜像,Java 17 对照镜像为 eclipse-temurin:17.0.20_8-jdk。断言或超时失败时检查完整输出,不从已经打印的前几行推断全部模式通过。
实验结束后回收手工输出目录:
case "$LAB_OUT" in
/tmp/container-learning.*) rm -r -- "$LAB_OUT" ;;
*) printf '%s\n' '保留未知目录' ;;
esac
unset LAB_OUT权威资料与规范地址
按方法语义查 API,按实现字段查固定版本源码。
| 资料 | 完整地址 |
|---|---|
| ConcurrentHashMap API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ConcurrentHashMap.html |
| ConcurrentMap 接口 | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ConcurrentMap.html |
| OpenJDK jdk-25+36 ConcurrentHashMap | https://github.com/openjdk/jdk/blob/jdk-25%2B36/src/java.base/share/classes/java/util/concurrent/ConcurrentHashMap.java |
| ConcurrentSkipListMap | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ConcurrentSkipListMap.html |
| Collections API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/Collections.html |
| CopyOnWriteArrayList API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/CopyOnWriteArrayList.html |
| BlockingQueue API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/BlockingQueue.html |
| java.util.concurrent 包说明 | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/package-summary.html |
| 其 API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ConcurrentLinkedQueue.html |
