You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何处理无锁数据结构的原子性变更?并发写入问题求助

问题分析

你的代码在并发场景下出现数据超量写入,核心原因是判断是否需要清理与后续的写入/清理操作不是原子执行:

  • 线程A读取size判断需要清理,释放读锁后,线程B可能已经写入数据并更新了size;等线程A拿到写锁执行清理时,线程B的写入已经生效,后续其他线程还可能在清理完成前继续写入,导致size统计完全混乱。
  • 单独维护的size和DashMap的实际数据总大小无法保证一致性,DashMap的并发修改不会同步更新RwLock保护的size,进一步加剧统计误差。
无锁友好的设计方案:双缓冲(切换式)存储模式

利用无锁结构的优势,通过双缓冲模式实现原子性的轮转切换,彻底避免竞态条件:

  1. 维护两个无锁Map实例(比如DashMap),用原子变量标记当前活跃的存储实例。
  2. 写入时直接写入活跃实例,同时用原子计数器累加该实例的大小。
  3. 当活跃实例的大小超过阈值时,原子切换到另一个空实例,旧实例可在后台异步清理(或直接丢弃)。

这种模式下所有操作都是原子或无锁的,完全规避RwLock带来的阻塞和竞态问题。

改进后的代码示例
use dashmap::DashMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;

struct Store {
    // 两个无锁存储实例
    stores: [Arc<DashMap<String, Vec<u8>>>; 2],
    // 标记当前活跃的存储索引(0或1)
    active_idx: AtomicUsize,
    // 每个存储实例的大小计数器
    sizes: [AtomicUsize; 2],
    // 存储容量阈值
    threshold: usize,
}

impl Store {
    pub fn new(threshold: usize) -> Self {
        Self {
            stores: [
                Arc::new(DashMap::new()),
                Arc::new(DashMap::new()),
            ],
            active_idx: AtomicUsize::new(0),
            sizes: [AtomicUsize::new(0), AtomicUsize::new(0)],
            threshold,
        }
    }

    pub fn put(&self, key: String, buf: Vec<u8>) {
        let buf_size = buf.len();
        loop {
            let idx = self.active_idx.load(Ordering::Acquire);
            let current_size = self.sizes[idx].load(Ordering::Acquire);

            // 检查当前活跃实例是否超过阈值
            if current_size + buf_size > self.threshold {
                // 尝试切换到另一个实例
                let new_idx = 1 - idx;
                if self.active_idx.compare_exchange(
                    idx,
                    new_idx,
                    Ordering::Release,
                    Ordering::Relaxed,
                ).is_ok() {
                    // 切换成功,清空旧实例(可选:也可后台异步清理)
                    self.stores[idx].clear();
                    self.sizes[idx].store(0, Ordering::Release);
                    println!("Store switched to instance {}, old instance cleaned", new_idx);
                }
                // 切换失败说明其他线程已完成切换,重新循环
                continue;
            }

            // 写入当前活跃实例
            self.stores[idx].entry(key.clone())
                .and_modify(|e| e.extend(&buf))
                .or_insert(buf);
            // 原子更新大小
            self.sizes[idx].fetch_add(buf_size, Ordering::Release);
            println!("Buff added: {} to instance {}", buf_size, idx);
            break;
        }
    }
}
方案优势
  • 完全无锁:仅用原子变量控制状态切换,DashMap本身是无锁结构,避免RwLock的阻塞和竞态。
  • 原子切换:通过compare_exchange保证只有一个线程能完成存储实例切换,不会出现多线程同时触发清理的情况。
  • 数据一致性:每个存储实例的大小与实际数据绑定,避免单独维护size带来的不一致问题。

内容的提问来源于stack exchange,提问作者newSqlz

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 18:53:31