ITADN
wvwwvwwv/scalable-concurrent-containers
wvwwvwwv/scalable-concurrent-containers · 文件 下载 ZIP
文件最后提交记录最后更新时间
README.md
以下内容由 AI 翻译,如有问题请点此提交 issue 反馈

已迁移至 codeberg

此仓库已迁移至 codeberg,将不再在此处维护。

Scalable Concurrent Containers

Cargo Crates.io GitHub Workflow Status

一组高性能容器,同时提供异步和同步接口。

特性

  • 同时提供异步和同步接口。
  • SIMD 查找以并行扫描多个条目:在 x86_64 上需要 RUSTFLAGS='-C target_feature=+avx2'
  • EquivalentLoomSerde 支持:features = ["equivalent", "loom", "serde"]

并发容器

  • HashMap 是一个并发哈希映射。
  • HashSet 是一个并发哈希集合。
  • HashIndex 是一个读优化的并发哈希映射。
  • HashCache 是一个由 HashMap 支持的 32 路关联并发缓存。
  • TreeIndex 是一个读优化的并发 B+ 树。

HashMap

HashMap 是一种针对高并行写密集型工作负载优化的并发哈希映射。HashMap 的结构是一个无锁的条目桶数组栈。条目桶数组由 sdd 管理,从而实现对它的无锁访问和非阻塞容器调整大小。每个桶是一个固定大小的条目数组,由一个读写锁保护,该锁同时提供阻塞和异步方法。

锁定行为

条目访问:细粒度锁定

对条目的读/写访问由包含该条目的桶中的读写锁进行串行化。没有容器级别的锁;因此,容器越大,桶级别锁发生竞争的可能性就越低。

调整大小:无锁

HashMap 的调整大小是完全非阻塞和无锁的;调整大小不会阻塞对容器的任何其他读/写访问或调整大小尝试。调整大小类似于向无锁栈中压入一个新的桶数组。旧桶数组中的每个条目将在未来对容器的访问时增量迁移到新的桶数组中,并且旧桶数组在变为空后最终会被丢弃。

示例

插入的条目可以同步或异步地更新、读取和删除。

use scc::HashMap;

let hashmap: HashMap<u64, u32> = HashMap::default();

assert!(hashmap.insert_sync(1, 0).is_ok());
assert!(hashmap.insert_sync(1, 1).is_err());
assert_eq!(hashmap.upsert_sync(1, 1).unwrap(), 0);
assert_eq!(hashmap.update_sync(&1, |_, v| { *v = 3; *v }).unwrap(), 3);
assert_eq!(hashmap.read_sync(&1, |_, v| *v).unwrap(), 3);
assert_eq!(hashmap.remove_sync(&1).unwrap(), (1, 3));

hashmap.entry_sync(7).or_insert(17);
assert_eq!(hashmap.read_sync(&7, |_, v| *v).unwrap(), 17);

let future_insert = hashmap.insert_async(2, 1);
let future_remove = hashmap.remove_async(&1);

HashMapEntry API 在工作流复杂时很有帮助。

use scc::HashMap;

let hashmap: HashMap<u64, u32> = HashMap::default();

hashmap.entry_sync(3).or_insert(7);
assert_eq!(hashmap.read_sync(&3, |_, v| *v), Some(7));

hashmap.entry_sync(4).and_modify(|v| { *v += 1 }).or_insert(5);
assert_eq!(hashmap.read_sync(&4, |_, v| *v), Some(5));

HashMap 不提供 Iterator,因为无法将 Iterator::Item 的生命周期限制在 Iterator 内。可以通过依赖内部可变性来规避此限制,例如让返回的引用持有锁。然而,如果未正确使用,可能会导致死锁,并且频繁获取锁可能会影响性能。因此,未实现 Iterator;相反,HashMap 提供了若干方法以同步或异步方式遍历条目:iter_{async|sync}iter_mut_{async|sync}retain_{async|sync}begin_{async|sync}OccupiedEntry::next_{async|sync}OccupiedEntry::remove_and_{async|sync}

use scc::HashMap;

let hashmap: HashMap<u64, u32> = HashMap::default();

assert!(hashmap.insert_sync(1, 0).is_ok());
assert!(hashmap.insert_sync(2, 1).is_ok());

// Entries can be modified or removed via `retain_sync`.
let mut acc = 0;
hashmap.retain_sync(|k, v_mut| { acc += *k; *v_mut = 2; true });
assert_eq!(acc, 3);
assert_eq!(hashmap.read_sync(&1, |_, v| *v).unwrap(), 2);
assert_eq!(hashmap.read_sync(&2, |_, v| *v).unwrap(), 2);

// `iter_sync` returns `true` when all the entries satisfy the predicate.
assert!(hashmap.insert_sync(3, 2).is_ok());
assert!(!hashmap.iter_sync(|k, _| *k == 3));

// Multiple entries can be removed through `retain_sync`.
hashmap.retain_sync(|k, v| *k == 1 && *v == 2);

// `hash_map::OccupiedEntry` also can return the next closest occupied entry.
let first_entry = hashmap.begin_sync();
assert_eq!(*first_entry.as_ref().unwrap().key(), 1);
let second_entry = first_entry.and_then(|e| e.next_sync());
assert!(second_entry.is_none());

fn is_send<T: Send>(f: &T) -> bool {
    true
}

// Asynchronous iteration over entries using `iter_async`.
let future_scan = hashmap.iter_async(|k, v| { println!("{k} {v}"); true });
assert!(is_send(&future_scan));

// Asynchronous iteration over entries using the `Entry` API.
let future_iter = async {
    let mut iter = hashmap.begin_async().await;
    while let Some(entry) = iter {
        // `OccupiedEntry` can be sent across awaits and threads.
        assert!(is_send(&entry));
        assert_eq!(*entry.key(), 1);
        iter = entry.next_async().await;
    }
};
assert!(is_send(&future_iter));

HashSet

HashSetHashMap 的一个特殊版本,其值类型为 ()

示例

大多数 HashSet 方法与 HashMap 的方法相同,只是它们不接收值参数,并且一些用于值修改的 HashMap 方法未在 HashSet 中实现。

use scc::HashSet;

let hashset: HashSet<u64> = HashSet::default();

assert!(hashset.read_sync(&1, |_| true).is_none());
assert!(hashset.insert_sync(1).is_ok());
assert!(hashset.read_sync(&1, |_| true).unwrap());

let future_insert = hashset.insert_async(2);
let future_remove = hashset.remove_async(&1);

HashIndex

HashIndexHashMap 的读优化版本。在 HashIndex 中,桶数组的内存由 sdd 管理,条目桶的内存也由 sdd 保护,从而实现对单个条目的无锁读取访问。

条目生命周期

HashIndex 不会立即丢弃已删除的条目;相反,只有当 sdd 机制确保这些条目没有潜在读者后,桶再次被访问时才会丢弃它们。这意味着,只要存在潜在读者或桶未被访问,已删除的条目就可能一直存在。因此,HashIndex 不是写密集型且涉及大条目大小的工作负载的最佳选择。

示例

peekpeek_with 方法是完全无锁的。

use scc::HashIndex;

let hashindex: HashIndex<u64, u32> = HashIndex::default();

assert!(hashindex.insert_sync(1, 0).is_ok());

// `peek` and `peek_with` are lock-free.
assert_eq!(hashindex.peek_with(&1, |_, v| *v).unwrap(), 0);

let future_insert = hashindex.insert_async(2, 1);
let future_remove = hashindex.remove_if_async(&1, |_| true);

HashIndexEntry API 可以更新现有条目。

use scc::HashIndex;

let hashindex: HashIndex<u64, u32> = HashIndex::default();
assert!(hashindex.insert_sync(1, 1).is_ok());

if let Some(mut o) = hashindex.get_sync(&1) {
    // Create a new version of the entry.
    o.update(2);
};

if let Some(mut o) = hashindex.get_sync(&1) {
    // Update the entry in place.
    unsafe { *o.get_mut() = 3; }
};

HashIndex 实现了 Iterator

use scc::HashIndex;

use sdd::Guard;

let hashindex: HashIndex<u64, u32> = HashIndex::default();

assert!(hashindex.insert_sync(1, 0).is_ok());

// Existing values can be replaced with a new one.
hashindex.get_sync(&1).unwrap().update(1);

let guard = Guard::new();

// An `Guard` has to be supplied to `iter`.
let mut iter = hashindex.iter(&guard);

let entry_ref = iter.next().unwrap();
assert_eq!(iter.next(), None);

HashCache

HashCache 是一种基于 HashMap 实现的 32 路组相联并发缓存。HashCache 不跟踪整个缓存中最近最少使用的条目。相反,每个桶维护一个已占用条目的双向链表,该链表在条目访问时更新。

示例

当插入新条目且桶已满时,桶中的 LRU 条目会被驱逐。

use scc::HashCache;

let hashcache: HashCache<u64, u32> = HashCache::with_capacity(100, 2000);

/// The capacity cannot exceed the maximum capacity.
assert_eq!(hashcache.capacity_range(), 128..=2048);

/// If the bucket corresponding to `1` or `2` is full, the LRU entry will be evicted.
assert!(hashcache.put_sync(1, 0).is_ok());
assert!(hashcache.put_sync(2, 0).is_ok());

/// `1` becomes the most recently accessed entry in the bucket.
assert!(hashcache.get_sync(&1).is_some());

/// An entry can be normally removed.
assert_eq!(hashcache.remove_sync(&2).unwrap(), (2, 0));

TreeIndex

TreeIndex 是一种针对读操作优化的 B+ 树变体。sdd 保护单个条目所使用的内存,从而实现对它们的无锁读取访问。

锁定行为

读取访问始终是无锁且非阻塞的。对条目的写入访问在无锁且非阻塞,只要不需要进行结构变更。然而,当写操作导致节点分裂或合并时,受影响范围内的其他写操作将被阻塞。

条目生命周期

TreeIndex 不会立即删除已移除的条目。相反,这些条目会在叶节点被清空或分裂时被删除,这使得 TreeIndex 在写密集型工作负载下成为次优选择。

示例

当内部节点分裂或合并时,会获取或等待锁,但阻塞操作不会影响读取操作。

use scc::TreeIndex;

let treeindex: TreeIndex<u64, u32> = TreeIndex::new();

assert!(treeindex.insert_sync(1, 2).is_ok());

// `peek` and `peek_with` are lock-free.
assert_eq!(treeindex.peek_with(&1, |_, v| *v).unwrap(), 2);
assert!(treeindex.remove_sync(&1));

let future_insert = treeindex.insert_async(2, 3);
let future_remove = treeindex.remove_if_async(&1, |v| *v == 2);

条目可以在不获取任何锁的情况下进行扫描。

use scc::TreeIndex;

use sdd::Guard;

let treeindex: TreeIndex<u64, u32> = TreeIndex::new();

assert!(treeindex.insert_sync(1, 10).is_ok());
assert!(treeindex.insert_sync(2, 11).is_ok());
assert!(treeindex.insert_sync(3, 13).is_ok());

let guard = Guard::new();

// `iter` iterates over entries without acquiring a lock.
let mut iter = treeindex.iter(&guard);
assert_eq!(iter.next().unwrap(), (&1, &10));
assert_eq!(iter.next().unwrap(), (&2, &11));
assert_eq!(iter.next().unwrap(), (&3, &13));
assert!(iter.next().is_none());

可以扫描特定范围的键。

use scc::TreeIndex;

use sdd::Guard;

let treeindex: TreeIndex<u64, u32> = TreeIndex::new();

for i in 0..10 {
    assert!(treeindex.insert_sync(i, 10).is_ok());
}

let guard = Guard::new();

assert_eq!(treeindex.range(1..1, &guard).count(), 0);
assert_eq!(treeindex.range(4..8, &guard).count(), 4);
assert_eq!(treeindex.range(4..=8, &guard).count(), 5);

性能

SIMD 支持

HashMap 针对 256 位 SIMD 指令进行了优化。因此,建议在 x86-64 目标上使用 avx2 或等效选项进行编译,或在其他平台上使用相应的特性。

  • 请注意,Apple M 系列 CPU 不支持实现最佳性能所需的 256 位 SIMD 指令。

HashMap 尾部延迟

在 Apple M4 Pro 上,1048576 次插入操作(K = u64, V = u64)的延迟分布的预期尾部延迟小于 50 微秒。

HashMapHashIndex 吞吐量

Changelog