предисловие
В прошлой статье я говорил о соответствующей реализации HashMap в исходном коде, и мы знаем, что она небезопасна для потоков.При использовании в параллельной среде HashMap может генерировать круговой связанный список при расширении, что приводит к бесконечному циклу get , тайм-аут. В этой статье давайте представим HashMap, используемый в параллельной среде — ConcurrentHashMap, Ниже приведена диаграмма его классов.
Реализация в JDK1.7
Сегмент — это реентерабельная блокировка, которая действует как блокировка в ConcurrentHashMap; HashEntry используется для хранения данных пары ключ-значение.
ConcurrentHashMap содержит массив сегментов. Структура сегмента аналогична структуре HashMap, которая представляет собой массив и структуру связанного списка. Сегмент содержит массив HashEntry, каждый HashEntry является элементом структуры связанного списка, и каждый сегмент защищает элемент в массиве HashEntry.При изменении данных массива HashEntry вы должны сначала получить соответствующую блокировку сегмента.
ConcurrentHashMap разделяет данные на сегменты с помощью технологии блокировки сегментов, а затем назначает блокировку каждому сегменту данных. Когда поток занимает блокировку для доступа к одному сегменту данных, другие потоки также могут получить доступ к другим сегментам данных. добиться истинного одновременного доступа.
static final class Segment extends ReentrantLock implements Serializable {
private static final long serialVersionUID = 2249069246763182397L;
static final int MAX_SCAN_RETRIES =
Runtime.getRuntime().availableProcessors() > 1 ? 64 : 1;
transient volatile HashEntry[] table;
transient int count;
transient int modCount;
transient int threshold;
final float loadFactor;
... ...
}
1. Структура хранения
static final class HashEntry<K,V> {
final int hash;
final K key;
volatile V value;
volatile HashEntry<K,V> next;
}
static final class Segment<K,V> extends ReentrantLock implements Serializable {
private static final long serialVersionUID = 2249069246763182397L;
static final int MAX_SCAN_RETRIES =
Runtime.getRuntime().availableProcessors() > 1 ? 64 : 1;
transient volatile HashEntry<K,V>[] table;
transient int count;
transient int modCount;
transient int threshold;
final float loadFactor;
}
final Segment<K,V>[] segments;
public ConcurrentHashMap(int initialCapacity,float loadFactor, int concurrencyLevel)
- Начальная емкость. Начальная емкость представляет собой все массивы сегментов с общим количеством HashenTry. Если INITIALCAPACITY не равна мощности 2, потребуется мощность, превышающая INITIALCAPACITY.
- Коэффициент нагрузки: по умолчанию 0,75.
- Уровень параллелизма: сколько одновременных потоков может быть разрешено одновременно. Сегментов столько же, сколько и concurrencyLevel, и, конечно, будет взята степень двойки, превышающая или равная этому значению.
static final int DEFAULT_CONCURRENCY_LEVEL = 16;
Далее давайте рассмотрим несколько ключевых функций в ConcurrentHashMap, методы get, put, rehash (расширение) и size, чтобы увидеть, как он достигает параллелизма.
2. получить операцию
- Вычислить hashCode по ключу;
- Найдите сегмент в соответствии с хэш-кодом, рассчитанным на шаге 1. Если сегмент не равен нулю && segment.table не равен нулю, перейдите к шагу 3, в противном случае верните значение null, а значение, соответствующее ключу, не существует;
- Найдите соответствующий hashEntry в таблице в соответствии с hashCode, пройдите по hashEntry, если ключ существует, верните значение, соответствующее ключу;
- После шага 3 значение, соответствующее ключу, по-прежнему не найдено, и возвращается null, а значение, соответствующее ключу, не существует.
3. поставить операцию
- Проверка параметра, значение не может быть нулевым, и если оно равно нулю, возникает исключение нулевого указателя;
- Вычислить hashCode ключа;
- Найдите сегмент, если сегмент не существует, создайте новый сегмент;
- Вызовите метод put сегмента, чтобы выполнить операцию вставки в соответствующий сегмент.
Реализация метода put сегмента
2. Найдите конкретный HashEntry в массиве HashEntry;
3. Перейдите по связанному списку HashEntry, если вставляемый ключ уже существует:
- Вам нужно обновить значение, соответствующее ключу (!onlyIfAbsent), обновить oldValue=newValue и перейти к шагу 5;
- В противном случае перейдите сразу к шагу 5;
5. Снимите блокировку и верните oldValue.
- Первое: расширение массива HashEntry;
- Второй: найдите позицию, соответствующую добавленному элементу, а затем поместите его в массив HashEntry.
4. размер операции
/**
* The number of elements. Accessed only either within locks
* or among other volatile reads that maintain visibility.
*/
transient int count;
static final int RETRIES_BEFORE_LOCK = 2;
public int size() {
// Try a few times to get accurate count. On failure due to
// continuous async changes in table, resort to locking.
final Segment<K,V>[] segments = this.segments;
int size;
boolean overflow; // true if size overflows 32 bits
long sum; // sum of modCounts
long last = 0L; // previous sum
int retries = -1; // first iteration isn't retry
try {
for (;;) {
// 超过尝试次数,则对每个 Segment 加锁
if (retries++ == RETRIES_BEFORE_LOCK) {
for (int j = 0; j < segments.length; ++j)
ensureSegment(j).lock(); // force creation
}
sum = 0L;
size = 0;
overflow = false;
for (int j = 0; j < segments.length; ++j) {
Segment<K,V> seg = segmentAt(segments, j);
if (seg != null) {
sum += seg.modCount;
int c = seg.count;
if (c < 0 || (size += c) < 0)
overflow = true;
}
}
// 连续两次得到的结果一致,则认为这个结果是正确的
if (sum == last)
break;
last = sum;
}
} finally {
if (retries > RETRIES_BEFORE_LOCK) {
for (int j = 0; j < segments.length; ++j)
segmentAt(segments, j).unlock();
}
}
return overflow ? Integer.MAX_VALUE : size;
}
Как ConcurrentHashMap определяет, что cout сегмента изменился во время статистического процесса?
Изменения в JDK 1.8
Структура данных реализована в виде массива + связанный список + красное черное дерево. Когда количество узлов в цепочке превышает 8, она преобразуется в хранилище структуры данных типа red-hahead, так что целью разработки является повышение эффективности чтения того же связанного списка.
В Java 8 были сделаны следующие оптимизации:
- Откажитесь от сегмента и напрямую используйте Node (унаследованный от Map.Entry) в качестве элемента таблицы.
- При модификации ReentrantLock больше не используется для блокировки, а встроенный синхронизированный используется для блокировки.Встроенная блокировка Java8 гораздо более оптимизирована, чем предыдущая версия.По сравнению с ReentrantLock производительность неплохая.
- Оптимизирован метод size и добавлен внутренний класс CounterCell для параллельного вычисления количества элементов в каждом сегменте.
- Отрицательное число означает инициализацию или расширение, - 1 --1 означает инициализацию, и -n означает N - 1 потоки расширяются.
- 0 положительное число указывает, что не было инициализировано. Другое положительное число указывает на меньший размер расширения.
static class Node<K,V> implements Map.Entry<K,V> {
final int hash;
final K key;
volatile V val;
volatile Node<K,V> next;
}
CAS-операция
В ConcurrentHashMap есть три основные операции CAS.
- Табат: получение массива узла в положении I
- Castabat: установка узла в позиции массива i
- setTabAt: используйте volatile, чтобы установить узел в позицию i.
//获取索引i处Node
static final <K,V> Node<K,V> tabAt(Node<K,V>[] tab, int i) {
return (Node<K,V>)U.getObjectVolatile(tab, ((long)i << ASHIFT) + ABASE);
}
//利用CAS算法设置i位置上的Node节点(将c和table[i]比较,相同则插入v)
static final <K,V> boolean casTabAt(Node<K,V>[] tab, int i,
Node<K,V> c, Node<K,V> v) {
return U.compareAndSwapObject(tab, ((long)i << ASHIFT) + ABASE, c, v);
}
//利用volatile设置节点位置i的值,仅在上锁区被调用
static final <K,V> void setTabAt(Node<K,V>[] tab, int i, Node<K,V> v) {
U.putObjectVolatile(tab, ((long)i << ASHIFT) + ABASE, v);
}
метод initTable()
private final Node<K,V>[] initTable() {
Node<K,V>[] tab; int sc;
while ((tab = table) == null || tab.length == 0) {
//如果一个线程发现sizeCtl<0,意味着另外的线程
//执行CAS操作成功,当前线程只需要让出cpu时间片
if ((sc = sizeCtl) < 0)
Thread.yield();
else if (U.compareAndSwapInt(this, SIZECTL, sc, -1)) {
//CAS方法把sizectl置为-1,表示本线程正在进行初始化
try {
if ((tab = table) == null || tab.length == 0) {
//DEFAULT_CAPACITY 默认初始容量是 16
int n = (sc > 0) ? sc : DEFAULT_CAPACITY;
@SuppressWarnings("unchecked")
//初始化数组,长度为 16 或初始化时提供的长度
Node<K,V>[] nt = (Node<K,V>[])new Node<?,?>[n];
//将这个数组赋值给 table,table 是 volatile 的
table = tab = nt;
//如果 n 为 16 的话,那么这里 sc = 12
//其实就是 0.75 * n
sc = n - (n >>> 2);
}
} finally {
sizeCtl = sc;
}
break;
}
}
return tab;
}
Вызов initTable будет оценивать значение sizeCtl.Если значение равно -1, это означает, что он инициализируется, и yield() будет вызываться для ожидания.
Если значение равно 0, то сначала вызывается алгоритм CAS, чтобы установить его в -1, а затем инициализируется.
Таким образом, поток, который выполняет первую операцию размещения, выполнит метод Unsafe.compareAndSwapInt, чтобы изменить sizeCtl на -1, и только один поток может быть успешно изменен. дождитесь завершения инициализации таблицы.
Из вышеизложенного видно, что инициализация является однопоточной операцией.
метод put()
public V put(K key, V value) {
return putVal(key, value, false);
}
/** Implementation for put and putIfAbsent */
final V putVal(K key, V value, boolean onlyIfAbsent) {
//不允许key、value为空
if (key == null || value == null) throw new NullPointerException();
//返回 (h ^ (h >>> 16)) & HASH_BITS;
int hash = spread(key.hashCode());
int binCount = 0;
//循环,直到插入成功
for (Node[] tab = table;;) {
Node f; int n, i, fh;
if (tab == null || (n = tab.length) == 0)
//table为空,初始化table
tab = initTable();
else if ((f = tabAt(tab, i = (n - 1) & hash)) == null) {
//索引处无值
if (casTabAt(tab, i, null,
new Node(hash, key, value, null)))
break; // no lock when adding to empty bin
}
else if ((fh = f.hash) == MOVED)// MOVED=-1;
//检测到正在扩容,则帮助其扩容
tab = helpTransfer(tab, f);
else {
V oldVal = null;
//上锁(hash值相同的链表的头节点)
synchronized (f) {
if (tabAt(tab, i) == f) {
if (fh >= 0) {
//遍历链表节点
binCount = 1;
for (Node e = f;; ++binCount) {
K ek;
// hash和key相同,则修改value
if (e.hash == hash &&
((ek = e.key) == key ||
(ek != null && key.equals(ek)))) {
oldVal = e.val;
//仅putIfAbsent()方法中onlyIfAbsent为true
if (!onlyIfAbsent)
//putIfAbsent()包含key则返回get,否则put并返回
e.val = value;
break;
}
Node pred = e;
//已遍历到链表尾部,直接插入
if ((e = e.next) == null) {
pred.next = new Node(hash, key,
value, null);
break;
}
}
}
else if (f instanceof TreeBin) {// 树节点
Node p;
binCount = 2;
if ((p = ((TreeBin)f).putTreeVal(hash, key,
value)) != null) {
oldVal = p.val;
if (!onlyIfAbsent)
p.val = value;
}
}
}
}
if (binCount != 0) {
//判断是否要将链表转换为红黑树,临界值和HashMap一样也是8
if (binCount >= TREEIFY_THRESHOLD)
//若length<64,直接tryPresize,两倍table.length;不转树
treeifyBin(tab, i);
if (oldVal != null)
return oldVal;
break;
}
}
}
addCount(1L, binCount);
return null;
}
static final int spread(int h) {
return (h ^ (h >>> 16)) & HASH_BITS;
}
int index = (n - 1) & hash
4. Если f равно null, это означает, что элемент впервые вставляется в эту позицию в таблице, а узел Node вставляется методом Unsafe.compareAndSwapObject.
- Если CAS прошел успешно, значит, узел Node был вставлен, и разрыв выскакивает, тогда метод addCount(1L, binCount) проверит, нужно ли расширять текущую емкость.
- Если CAS дает сбой, это означает, что другой поток вставил узел заранее и снова пытается вставить узел в эту позицию.
6. В других случаях новый узел Node вставляется в соответствующую позицию в виде связанного списка или красно-черного дерева.Этот процесс использует синхронные встроенные блокировки для достижения параллелизма, а код такой же, как и выше.
Синхронизировать на узле f. Прежде чем узел будет вставлен, снова используйте tabAt(tab, i) == f, чтобы предотвратить его изменение другими потоками.
- Если f.Hash> = 0, это означает, что F - это узел головного узла связанного списка структуры, пройти связанный список, если обнаружен соответствующий узел узла, измените значение, в противном случае добавьте узел в конце связанного списка. Отказ
- Если f является узлом типа TreeBin, указывающим, что f является корневым узлом красно-черного дерева, то необходимо пройтись по элементам древовидной структуры, обновить или добавить узлы.
- Если количество узлов в связанном списке binCount >= TREEIFY_THRESHOLD (по умолчанию 8), то преобразовать связанный список в красно-черную древовидную структуру.
Связанный список с красно-черным деревом: treeifyBin()
private final void treeifyBin(Node[] tab, int index) {
Node b; int n, sc;
if (tab != null) {
// MIN_TREEIFY_CAPACITY 为 64
// 所以,如果数组长度小于 64 的时候,其实也就是 32 或者 16 或者更小的时候,会进行数组扩容
if ((n = tab.length) < MIN_TREEIFY_CAPACITY)
// 后面我们再详细分析这个方法
tryPresize(n << 1);
// b 是头结点
else if ((b = tabAt(tab, index)) != null && b.hash >= 0) {
// 加锁
synchronized (b) {
if (tabAt(tab, index) == b) {
// 下面就是遍历链表,建立一颗红黑树
TreeNode hd = null, tl = null;
for (Node e = b; e != null; e = e.next) {
TreeNode p =
new TreeNode(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(hd));
}
}
}
}
}
Расширение: tryPresize()
Расширение здесь тоже двойное, и емкость массива удваивается после расширения.
// 首先要说明的是,方法参数 size 传进来的时候就已经翻了倍了
private final void tryPresize(int size) {
// c:size 的 1.5 倍,再加 1,再往上取最近的 2 的 n 次方。
int c = (size >= (MAXIMUM_CAPACITY >>> 1)) ? MAXIMUM_CAPACITY :
tableSizeFor(size + (size >>> 1) + 1);
int sc;
while ((sc = sizeCtl) >= 0) {
Node<K,V>[] tab = table; int n;
// 这个 if 分支和之前说的初始化数组的代码基本上是一样的
if (tab == null || (n = tab.length) == 0) {
n = (sc > c) ? sc : c;
if (U.compareAndSwapInt(this, SIZECTL, sc, -1)) {
try {
if (table == tab) {
@SuppressWarnings("unchecked")
Node<K,V>[] nt = (Node<K,V>[])new Node<?,?>[n];
table = nt;
sc = n - (n >>> 2); // 0.75 * n
}
} finally {
sizeCtl = sc;
}
}
}
else if (c <= sc || n >= MAXIMUM_CAPACITY)
break;
else if (tab == table) {
int rs = resizeStamp(n);
if (sc < 0) {
Node<K,V>[] nt;
if ((sc >>> RESIZE_STAMP_SHIFT) != rs || sc == rs + 1 ||
sc == rs + MAX_RESIZERS || (nt = nextTable) == null ||
transferIndex <= 0)
break;
// 2. 用 CAS 将 sizeCtl 加 1,然后执行 transfer 方法
// 此时 nextTab 不为 null
if (U.compareAndSwapInt(this, SIZECTL, sc, sc + 1))
transfer(tab, nt);
}
// 1. 将 sizeCtl 设置为 (rs << RESIZE_STAMP_SHIFT) + 2)
// 调用 transfer 方法,此时 nextTab 参数为 null
else if (U.compareAndSwapInt(this, SIZECTL, sc,
(rs << RESIZE_STAMP_SHIFT) + 2))
transfer(tab, null);
}
}
}
Что касается исходного кода метода transfer(), то я его здесь разбирать не буду, его примерная функция такова:Перенесите элементы исходного массива вкладок в новый массив nextTab.
получить () метод
public V get(Object key) {
Node<K,V>[] tab; Node<K,V> e, p; int n, eh; K ek;
int h = spread(key.hashCode());
if ((tab = table) != null && (n = tab.length) > 0 &&
(e = tabAt(tab, (n - 1) & h)) != null) {//tabAt(i),获取索引i处Node
// 判断头结点是否就是我们需要的节点
if ((eh = e.hash) == h) {
if ((ek = e.key) == key || (ek != null && key.equals(ek)))
return e.val;
}
// 如果头结点的 hash<0,说明正在扩容,或者该位置是红黑树
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;
}
Node<K,V> find(int h, Object k) {
Node<K,V> e = this;
if (k != null) {
do {
K ek;
if (e.hash == h &&
((ek = e.key) == k || (ek != null && k.equals(ek))))
return e;
} while ((e = e.next) != null);
}
return null;
}
- Если позиция равна нулю, просто верните значение null напрямую
- Если узел в этой позиции именно то, что нам нужно, просто верните значение узла
- Если хеш-значение узла в этом месте меньше 0, это означает, что емкость расширяется, или это красно-черное дерево.
- Если не удовлетворены три вышеуказанных, то есть список цепочек, сравнение обхода может
резюме
До сих пор я в основном рассмотрел реализацию ConcurrentHashMap в JDK 1.7 и 1.8 и подробно проанализировал несколько важных реализаций методов: инициализация, установка и получение. В JDK1.8 ConcurrentHashMap претерпел серьезные изменения, заменив механизм блокировки сегмента сегмента в исходной версии 1.7, используя реализацию CAS+synchronized, тем самым поддерживая более высокий параллелизм.
Это всего лишь мое второе исследование ConcurrentHashMap. Если я хочу лучше понять и понять тонкости реализации ConcurrentHashMap, я лично чувствую, что мне нужно прочитать его еще несколько раз в будущем. Я верю, что каждый раз будут новые достижения. время.