概述

ConcurrentHashmap是一个支持并发检索和并发更新的线程安全的HashMap,它是不支持空key和value的。ConcurrentHashMap在JDK1.7之前使用的Lock和Segment(分段锁)来实现并发安全的,JDK1.8改用CAS和synchronized来实现的。我们这次主要分析JDK1.8中的实现方式。

简单使用

ConcurrentHashMap在使用上和我们平时常用的HashMap差异不是很大。只不过一个支持并发操作,一个不支持并发操作而已。
下面我们就写一个非常非常简单的例子:

1
2
3
4
5
public static void main(String[] args) {
ConcurrentHashMap<Integer,Integer> map=new ConcurrentHashMap<>();
map.put(1,2);
System.out.println(map.get(1));
}`

下面我们就针对这个例子,分析一下其内部实现。

源码分析

数据结构

lbaPOJ.png
从图中我们可以看出有许多种的节点类。那么我们下面就分别解释一波这些节点类。
lba3TI.png

Node<K,V>这个节点类是最基础的,它存储了键值对(值使用了volatile关键字确保可见性),hash值,下一个节点的引用。

TreeNode<K,V>,这个节点类表示红黑树节点,当链表的长度大于等于8,且数组的长度大于64的时候,就会将链表节点转换为红黑树节点,然后将这些红黑树节点放到TreeBin对象中,由TreeBin对象来完成对红黑树的封装。

TreeBin<K,V>,封装了红黑树根节点。

ForwardingNode<K, V>,在节点转移的时候,用于连接两个table的节点类。它的内部包含一个nextTable指针,指向下一个table。这个节点仅仅作为占位节点表示当前节点已经被移动。

ReservationNode<K,V>这个也是一个占位节点,表示当前节点已经被占用。

重要属性

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
/*最大容量*/
private static final int MAXIMUM_CAPACITY = 1 << 30;

/*默认容量*/
private static final int DEFAULT_CAPACITY = 16;

/*数组的最大容量*/
static final int MAX_ARRAY_SIZE = Integer.MAX_VALUE - 8;

/*默认的最大并发数,为了兼容1.7*/
private static final int DEFAULT_CONCURRENCY_LEVEL = 16;

/*负载因子*/
private static final float LOAD_FACTOR = 0.75f;

/*链表转红黑树的阈值*/
static final int TREEIFY_THRESHOLD = 8;

/*红黑树转链表的阈值*/
static final int UNTREEIFY_THRESHOLD = 6;

/*转红黑树时,数组容量的最小要求*/
static final int MIN_TREEIFY_CAPACITY = 64;

/*扩容的最小转移节点数*/
private static final int MIN_TRANSFER_STRIDE = 16;

/*sizeCtl中记录stamp的位数*/
private static int RESIZE_STAMP_BITS = 16;

/*帮助扩容时的最大线程数*/
private static final int MAX_RESIZERS = (1 << (32 - RESIZE_STAMP_BITS)) - 1;

/*size在sizeCtl中的偏移量*/
private static final int RESIZE_STAMP_SHIFT = 32 - RESIZE_STAMP_BITS;

/*存放节点的数组*/
transient volatile Node<K,V>[] table;

/*一个过度表,只会在扩容的时候使用*/
private transient volatile Node<K,V>[] nextTable;

/*基础计数器值*/
private transient volatile long baseCount;

/*控制table初始化和扩容操作*/
private transient volatile int sizeCtl;

/*节点转移时,下一个需要转移的table索引*/
private transient volatile int transferIndex;

/*元素变化时用于控制自旋*/
private transient volatile int cellsBusy;

/*保存table中每个节点的元素个数作*/
private transient volatile CounterCell[] counterCells;

重要方法源码分析

put方法

put方法,实际上是调用了V putVal(K key, V value, boolean onlyIfAbsent)

1
2
3
4
public V put(K key, V value) {
return putVal(key, value, false);
}

V putVal(K key, V value, boolean onlyIfAbsent)的具体实现如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
final V putVal(K key, V value, boolean onlyIfAbsent) {
//从这个地方可以看出是不支持空值和空键的。
if (key == null || value == null) throw new NullPointerException();
//计算hash值
int hash = spread(key.hashCode());
int binCount = 0;
for (Node<K,V>[] tab = table;;) { //自旋
Node<K,V> f; int n, i, fh;
if (tab == null || (n = tab.length) == 0)
//还未初始化,进行初始化
tab = initTable();
else if ((f = tabAt(tab, i = (n - 1) & hash)) == null) {
//索引i的位置为空,直接插入
if (casTabAt(tab, i, null,
new Node<K,V>(hash, key, value, null)))
//CAS插入节点,成功了的话就直接返回
break; // no lock when adding to empty bin
}
else if ((fh = f.hash) == MOVED)
//当前节点正在进行移动状态,帮助移动
tab = helpTransfer(tab, f);
else {//hash冲突的处理
V oldVal = null;
synchronized (f) { //进行加锁操作
if (tabAt(tab, i) == f) {
if (fh >= 0) {
//f.hash>=0说明f是链表的头
//用于记录链表的节点数,后面以此来判断是否转为红黑树
binCount = 1;
for (Node<K,V> e = f;; ++binCount) {
K ek;
if (e.hash == hash &&
((ek = e.key) == key ||
(ek != null && key.equals(ek)))) {
oldVal = e.val;
if (!onlyIfAbsent)
e.val = value;
break;
}
Node<K,V> pred = e;
if ((e = e.next) == null) {
//将当前节点插入到链表尾部
pred.next = new Node<K,V>(hash, key,
value, null);
break;
}
}
}
//当前位置已经采用红黑树
else if (f instanceof TreeBin) {
Node<K,V> p;
binCount = 2;
if ((p = ((TreeBin<K,V>)f).putTreeVal(hash, key,
value)) != null) {
oldVal = p.val;
if (!onlyIfAbsent)
p.val = value;
}
}
}
}
if (binCount != 0) {
if (binCount >= TREEIFY_THRESHOLD)
//链表长度已经达到阈值了,需要转换为红黑树了
treeifyBin(tab, i);
if (oldVal != null)
return oldVal;
break;
}
}
}
//更新元素的数量
addCount(1L, binCount);
return null;
}

下面我们就梳理一下,整个put操作的整个流程:

  • 计算key的hash值,并根据hash值计算索引i
  • 如果但其哈希表还未初始化就调用initTable()进行初始化
  • 如果索引i的位置为空,那么就直接CAS将当前节点放入该位置即可。
  • 如果当前节点的hash值为-1,即当前节点处于移动状态,那么就调用helpTransfer(tab, f)帮助扩容
  • 如果不满足上述两种情况,那么使用synchronized进行加锁后,进行hash冲突处理。
    • 如果位置i是个链表,那么就遍历整个链表,如果在遍历的过程中,发现了某个节点的hash值与当前key的哈希值相同,那么就覆盖该节点,并停止遍历。如果到了链表尾部了,那么就创建一个新的节点加入链表尾部。
    • 如果位置i上为TreeBin。说明这个位置是红黑树,那么就调用putTreeVal方法,要么覆盖某个节点,要么创建一个新节点加入。
  • 插入完毕之后,如果是链表的话会检查链表的长度,如果达到链表转红黑树的阈值的话,会进行链表转红黑树的操作
  • 最后调用addCount方法,更新元素的数量

把整个流程看下来,我们发现在整个流程中,还有几个方法起到了非常重要的作用。
它们分别是initTable(),helpTransfer(tab, f) treeifyBin(tab, i);addCount(1L, binCount);下面我们就仔细分析一波这些方法。

initTable方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
private final Node<K,V>[] initTable() {
Node<K,V>[] tab; int sc;
while ((tab = table) == null || tab.length == 0) {//自旋,直到初始化完成
if ((sc = sizeCtl) < 0)//其它线程在初始化或转移时,让出CPU
Thread.yield(); // lost initialization race; just spin
else if (U.compareAndSwapInt(this, SIZECTL, sc, -1)) {
//设置sizectl为-1,表示当前线程正在初始化
try {
if ((tab = table) == null || tab.length == 0) {
int n = (sc > 0) ? sc : DEFAULT_CAPACITY;
@SuppressWarnings("unchecked")
//初始化
Node<K,V>[] nt = (Node<K,V>[])new Node<?,?>[n];
table = tab = nt;
//设置扩容阈值(0.75*n)
sc = n - (n >>> 2);
}
} finally {
//将sizeCtl设置为0.75*n
sizeCtl = sc;
}
break;
}
}
return tab;
}

在初始化过程中,首先会判断是否有线程正在初始化,如果正在初始化,那么就让出CPU时间片,自旋等待创建成功。如果没有线程正在初始化,那么该线程就会开始初始化。

那么是如和判断是否有其它线程正在初始化,和保障自己在创建过程中其它线程不会创建呢?
这个时候sizeCtl起到了非常重要的作用,一个线程开始初始化就会将sizeCtl设置为-1,这样其它线程就可以以此判断已经有线程在初始化了。初始化完成之后,会将sizeCtl设置为0.75*n。

helpTransfer方法

这个方法的作用是帮助其它线程进行转移操作

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
final Node<K,V>[] helpTransfer(Node<K,V>[] tab, Node<K,V> f) {
Node<K,V>[] nextTab; int sc;
if (tab != null && (f instanceof ForwardingNode) &&
(nextTab = ((ForwardingNode<K,V>)f).nextTable) != null) {
//计算操作栈校验码
int rs = resizeStamp(tab.length);
while (nextTab == nextTable && table == tab &&
(sc = sizeCtl) < 0) {
if ((sc >>> RESIZE_STAMP_SHIFT) != rs || sc == rs + 1 ||
sc == rs + MAX_RESIZERS || transferIndex <= 0)
break;//不需要转移,跳出
if (U.compareAndSwapInt(this, SIZECTL, sc, sc + 1)) {/*CAS
更新帮助转移的线程数*/
//调用transfer进行真正的转移
transfer(tab, nextTab);
break;
}
}
return nextTab;
}
return table;
}

transfer方法

transfer方法的作用,主要是转移或复制节点到新的table

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
   private final void transfer(Node<K,V>[] tab, Node<K,V>[] nextTab) {
int n = tab.length, stride;
if ((stride = (NCPU > 1) ? (n >>> 3) / NCPU : n) < MIN_TRANSFER_STRIDE)
stride = MIN_TRANSFER_STRIDE; // subdivide range
if (nextTab == null) { // initiating
try {
@SuppressWarnings("unchecked")
//建立一个比长两倍的nextTab
Node<K,V>[] nt = (Node<K,V>[])new Node<?,?>[n << 1];
nextTab = nt;
} catch (Throwable ex) { // try to cope with OOME
sizeCtl = Integer.MAX_VALUE;
return;
}
nextTable = nextTab;
transferIndex = n; //初始化table的最后一个索引
}
int nextn = nextTab.length;
/*初始化ForwardingNode节点,它持有nextTab的引用,每处理完后,
用该节点进行占位,表示该位置已经处理*/
ForwardingNode<K,V> fwd = new ForwardingNode<K,V>(nextTab);
boolean advance = true; //节点是否已经处理
boolean finishing = false; // to ensure sweep before committing nextTab
//i:当前处理的node的索引,bound表示需要处理节点的索引边界
for (int i = 0, bound = 0;;) {
Node<K,V> f; int fh;
while (advance) {
/*nextIndex:下一个需要处理的节点索引,
nextBound下一个需要处理节点的边界*/
int nextIndex, nextBound;
if (--i >= bound || finishing)
advance = false;
else if ((nextIndex = transferIndex) <= 0) {
//节点已全部转移
i = -1;
advance = false;
}
/*更新transferIndex(初始值为最后一个节点的索引)表示
从transferIndex开始,后面的节点的转移任务已经被领取了
在这个地方更新transferIndex的值(transferIndex-stride)
同时更新索引的边界
*/
else if (U.compareAndSwapInt
(this, TRANSFERINDEX, nextIndex,
nextBound = (nextIndex > stride ?
nextIndex - stride : 0))) {
bound = nextBound;
i = nextIndex - 1;
advance = false;
}
}
if (i < 0 || i >= n || i + n >= nextn) {
int sc;
if (finishing) {
//已经完成转移,更新相关属性
nextTable = null;
table = nextTab;
//更新扩容阈值
sizeCtl = (n << 1) - (n >>> 1);
return;
}

//当前线程任务完成,sizectl-1.每个线程的任务完成之后都是如此
if (U.compareAndSwapInt(this, SIZECTL, sc = sizeCtl, sc - 1)) {
//判断是否 还有其它线程在执行
if ((sc - 2) != resizeStamp(n) << RESIZE_STAMP_SHIFT)
return;//还有线程在执行,直接返回
//否则,再检查一遍,说明当前线程就是最后一个完成
//任务的线程了,它需要做一次检查
finishing = advance = true;
i = n; // recheck before commit
}
}
else if ((f = tabAt(tab, i)) == null)
//当前位置为null,直接替换为ForwardingNode,表明该位置已经被处理
advance = casTabAt(tab, i, null, fwd);
else if ((fh = f.hash) == MOVED)
//如果这个位置已经被处理过了,那么跳过该位置
advance = true; // already processed
else {
//对当前位置的节点进行真正的转移
synchronized (f) {
if (tabAt(tab, i) == f) {
//处理当前拿到的节点,构建两个node:ln原位置,hn:i+1位置
Node<K,V> ln, hn;
if (fh >= 0) {
//如果当前为链表节点

//fh&n将链表中的元素分为两部分
int runBit = fh & n;
Node<K,V> lastRun = f;
//从索引i查找到最后一个有效节点
for (Node<K,V> p = f.next; p != null; p = p.next) {
int b = p.hash & n;
if (b != runBit) {
runBit = b;
lastRun = p;
}
}
if (runBit == 0) {
ln = lastRun;
hn = null;
}
else {
hn = lastRun;
ln = null;
}
for (Node<K,V> p = f; p != lastRun; p = p.next) {
int ph = p.hash; K pk = p.key; V pv = p.val;
//将链表分解为两部分
if ((ph & n) == 0)
//在原位置
ln = new Node<K,V>(ph, pk, pv, ln);
else
//在i+1位置
hn = new Node<K,V>(ph, pk, pv, hn);
}
//放入nextTable的指定的位置
setTabAt(nextTab, i, ln);
setTabAt(nextTab, i + n, hn);
//在table位置上插入forwardNode节点,表示该位置已经处理
setTabAt(tab, i, fwd);
advance = true;
}
else if (f instanceof TreeBin) {
//如果是红黑树,同样也进行拆分为两部分
TreeBin<K,V> t = (TreeBin<K,V>)f;
TreeNode<K,V> lo = null, loTail = null;
TreeNode<K,V> hi = null, hiTail = null;
int lc = 0, hc = 0;
for (Node<K,V> e = t.first; e != null; e = e.next) {
int h = e.hash;
TreeNode<K,V> p = new TreeNode<K,V>
(h, e.key, e.val, null, null);
if ((h & n) == 0) {
if ((p.prev = loTail) == null)
lo = p;
else
loTail.next = p;
loTail = p;
++lc;
}
else {
if ((p.prev = hiTail) == null)
hi = p;
else
hiTail.next = p;
hiTail = p;
++hc;
}
}
//如果扩容后,已经不需要tree结构了,那么反向转换为链表结构
ln = (lc <= UNTREEIFY_THRESHOLD) ? untreeify(lo) :
(hc != 0) ? new TreeBin<K,V>(lo) : t;
hn = (hc <= UNTREEIFY_THRESHOLD) ? untreeify(hi) :
(lc != 0) ? new TreeBin<K,V>(hi) : t;
setTabAt(nextTab, i, ln);
setTabAt(nextTab, i + n, hn);
setTabAt(tab, i, fwd);
advance = true;
}
}
}
}
}
}

因为ConcurrentHashMap的扩容,实际上是新建了一个table,因此扩容最主要的任务就是将旧table中的节点转移到新的table中。

在以下三种情况下,是需要进行转移的:

  1. 对table进行扩容的时候
  2. 在调用addCount方法更新元素的数量的时候,发现元素的数量已经达到扩容的阈值的时候。
  3. 在进行put操作的时候,发现需要加入的位置的节点正在进行转移的时候,那么当前线程会帮助扩容。

在整个转移的过程中,有两个比较重要的地方,其中一个是transferIndex,它的初始值是最后一个节点,它的含义是:从transferIndex到最后一个节点的转移任务已经被领取。
还有一个是forwardNode节点,它用于标记已经处理过的位置。

addCount方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
private final void addCount(long x, int check) {
CounterCell[] as; long b, s;
if ((as = counterCells) != null ||
!U.compareAndSwapLong(this, BASECOUNT, b = baseCount, s = b + x)) {
CounterCell a; long v; int m;
boolean uncontended = true;
if (as == null || (m = as.length - 1) < 0 ||
(a = as[ThreadLocalRandom.getProbe() & m]) == null ||
!(uncontended =
U.compareAndSwapLong(a, CELLVALUE, v = a.value, v + x))) {
//在线程争用资源时,使用fullAddCount计算更新元素
fullAddCount(x, uncontended);
return;
}
if (check <= 1)
return;
//计算元素的总数,用于之后的扩容
s = sumCount();
}
if (check >= 0) {
//检查扩容
Node<K,V>[] tab, nt; int n, sc;
while (s >= (long)(sc = sizeCtl) && (tab = table) != null &&
(n = tab.length) < MAXIMUM_CAPACITY) {
int rs = resizeStamp(n);
if (sc < 0) {
//其它线程正在进行扩容
if ((sc >>> RESIZE_STAMP_SHIFT) != rs || sc == rs + 1 ||
sc == rs + MAX_RESIZERS || (nt = nextTable) == null ||
transferIndex <= 0)
break;
if (U.compareAndSwapInt(this, SIZECTL, sc, sc + 1))
transfer(tab, nt);
}
else if (U.compareAndSwapInt(this, SIZECTL, sc,
(rs << RESIZE_STAMP_SHIFT) + 2))
transfer(tab, null);
s = sumCount();
}
}
}

treeifyBin

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
private final void treeifyBin(Node<K,V>[] tab, int index) {
Node<K,V> b; int n, sc;
if (tab != null) {
//当数组还未超过64的时候,优先使用数组扩容,否则将链表转为红黑树
if ((n = tab.length) < MIN_TREEIFY_CAPACITY)
//两倍扩容
tryPresize(n << 1);
else if ((b = tabAt(tab, index)) != null && b.hash >= 0) {
synchronized (b) {
if (tabAt(tab, index) == b) {
//hd为头节点
TreeNode<K,V> hd = null, tl = null;
//遍历转换节点
for (Node<K,V> e = b; e != null; e = e.next) {
TreeNode<K,V> p =
new TreeNode<K,V>(e.hash, e.key, e.val,
null, null);
if ((p.prev = tl) == null)
hd = p;
else
tl.next = p;
tl = p;
}
setTabAt(tab, index, new TreeBin<K,V>(hd));
}
}
}
}
}

get方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
public V get(Object key) {
Node<K,V>[] tab; Node<K,V> e, p; int n, eh; K ek;
//计算hash值
int h = spread(key.hashCode());
if ((tab = table) != null && (n = tab.length) > 0 &&
(e = tabAt(tab, (n - 1) & h)) != null) {
if ((eh = e.hash) == h) {
if ((ek = e.key) == key || (ek != null && key.equals(ek)))
return e.val;
}
else if (eh < 0)
return (p = e.find(h, key)) != null ? p.val : null;
while ((e = e.next) != null) {
if (e.hash == h &&
((ek = e.key) == key || (ek != null && key.equals(ek))))
return e.val;
}
}
return null;
}

get方法的实现与hashmap的实现差不多,就不赘述了。