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

Rust中如何实现跨线程无锁更新内存存储?

解决方案

针对你遇到的Mutex阻塞问题,这里提供几种无锁/低阻塞的实现方案,适配不同的场景:

1. 用Tokio Watch通道实现全量状态更新

如果你的Store需要每次全量替换(比如每次API拉取的是完整数据集),tokio::sync::watch是最安全简单的选择。它支持单生产者多消费者,读操作完全无阻塞,写操作仅需发送新的状态副本,不会阻塞服务线程。

代码示例

use tokio::sync::watch;

#[tokio::main]
async fn main() {
    // 创建Watch通道,初始值为默认Store
    let (store_sender, store_receiver) = watch::channel(store::Store::default());

    // 启动服务和数据拉取任务
    let server_handle = start_server(store_receiver.clone());
    let fetch_a = start_fetch_a(store_sender.clone());
    let fetch_b = start_fetch_b(store_sender.clone());

    let _ = tokio::join!(server_handle, fetch_a, fetch_b);
}

// 服务端逻辑:获取最新状态
async fn start_server(mut receiver: watch::Receiver<store::Store>) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        loop {
            // 非阻塞获取当前最新的Store
            let current_store = receiver.borrow();
            // 用current_store处理REST请求
            // ...

            // 可选:等待状态更新(如果需要实时响应变化)
            if receiver.changed().await.is_err() {
                // 发送端已关闭,退出循环
                break;
            }
        }
    })
}

// 数据拉取任务A:更新状态并发送
async fn start_fetch_a(mut sender: watch::Sender<store::Store>) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        loop {
            // 从API拉取数据
            let data_a = fetch_api_a().await;

            // 克隆当前状态并更新
            let mut new_store = sender.borrow().clone();
            new_store.insert_a(data_a);

            // 发送新状态,通知所有接收端
            if sender.send(new_store).is_err() {
                // 所有接收端已关闭,退出循环
                break;
            }

            // 定期拉取
            tokio::time::sleep(tokio::time::Duration::from_secs(60)).await;
        }
    })
}

注意:如果Store体积较大,每次克隆会有一定开销,适合数据量不大的场景。

2. 用无锁哈希表实现增量更新

如果你的Store是键值对结构,且每次仅需更新部分数据,dashmap(无锁哈希表)是最优选择。它的读操作完全无锁,写操作仅锁定单个哈希桶,不会阻塞其他键的读写,性能远优于全局Mutex。

步骤

  1. 在Cargo.toml中添加依赖:
[dependencies]
dashmap = "5.4"
  1. 代码示例
use dashmap::DashMap;
use std::sync::Arc;

#[tokio::main]
async fn main() {
    // 用Arc包裹无锁哈希表作为Store
    let store = Arc::new(DashMap::new());

    let server_handle = start_server(store.clone());
    let fetch_a = start_fetch_a(store.clone());
    let fetch_b = start_fetch_b(store.clone());

    let _ = tokio::join!(server_handle, fetch_a, fetch_b);
}

// 服务端逻辑:无锁读取数据
async fn start_server(store: Arc<DashMap<String, Data>>) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        // 处理请求时直接读取,无阻塞
        if let Some(entry) = store.get("key_a") {
            // 使用entry.value()处理请求
        }
        // ...
    })
}

// 数据拉取任务A:增量更新哈希表
async fn start_fetch_a(store: Arc<DashMap<String, Data>>) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        loop {
            let data_a = fetch_api_a().await;
            // 仅锁定当前键对应的哈希桶,不影响其他操作
            store.insert("key_a".to_string(), data_a);

            tokio::time::sleep(tokio::time::Duration::from_secs(60)).await;
        }
    })
}

3. 用原子指针实现手动无锁替换

如果你需要极致性能且熟悉Rust内存模型,可以用AtomicPtr直接管理Store的指针,实现完全无锁的状态替换。这种方式需要手动处理内存安全,适合对性能要求极高的场景。

代码示例

use std::sync::atomic::{AtomicPtr, Ordering};
use std::sync::Arc;

#[tokio::main]
async fn main() {
    // 初始化默认Store并转为原始指针
    let initial_store = Arc::new(store::Store::default());
    let store_ptr = AtomicPtr::new(Arc::into_raw(initial_store));
    let store = Arc::new(store_ptr);

    let server_handle = start_server(store.clone());
    let fetch_a = start_fetch_a(store.clone());
    let fetch_b = start_fetch_b(store.clone());

    let _ = tokio::join!(server_handle, fetch_a, fetch_b);
}

// 服务端逻辑:无锁读取指针
async fn start_server(store: Arc<AtomicPtr<store::Store>>) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        loop {
            // 原子加载指针
            let ptr = store.load(Ordering::Acquire);
            // 转为Arc确保内存安全
            let current_store = unsafe { Arc::clone(&*Arc::from_raw(ptr)) };
            // 释放临时Arc,避免引用计数泄漏
            drop(unsafe { Arc::from_raw(ptr) });

            // 使用current_store处理请求
            // ...
        }
    })
}

// 数据拉取任务A:原子替换指针
async fn start_fetch_a(store: Arc<AtomicPtr<store::Store>>) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        loop {
            let data_a = fetch_api_a().await;

            // 获取当前Store并克隆
            let old_ptr = store.load(Ordering::Acquire);
            let old_store = unsafe { Arc::from_raw(old_ptr) };
            let mut new_store = (*old_store).clone();
            new_store.insert_a(data_a);
            let new_store = Arc::new(new_store);

            // 原子替换指针
            store.swap(Arc::into_raw(new_store), Ordering::Release);
            // 释放旧Store的Arc,确保内存被正确回收
            drop(old_store);

            tokio::time::sleep(tokio::time::Duration::from_secs(60)).await;
        }
    })
}

注意:这种方式需要严格遵循Rust的内存安全规则,避免悬垂指针和内存泄漏,仅推荐资深开发者使用。

4. 退而求其次:用异步读写锁降低阻塞

如果不想引入额外依赖,tokio::sync::RwLock是比Mutex更好的选择。它支持并发读操作,仅在写操作时阻塞读,阻塞范围远小于全局Mutex。

代码示例

use tokio::sync::RwLock;
use std::sync::Arc;

#[tokio::main]
async fn main() {
    let store = Arc::new(RwLock::new(store::Store::default()));

    let server_handle = start_server(store.clone());
    let fetch_a = start_fetch_a(store.clone());
    let fetch_b = start_fetch_b(store.clone());

    let _ = tokio::join!(server_handle, fetch_a, fetch_b);
}

// 服务端逻辑:获取读锁,允许多个并发读
async fn start_server(store: Arc<RwLock<store::Store>>) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        let current_store = store.read().await;
        // 使用current_store处理请求
        // ...
    })
}

// 数据拉取任务A:获取写锁,独占更新
async fn start_fetch_a(store: Arc<RwLock<store::Store>>) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        loop {
            let data_a = fetch_api_a().await;
            let mut store = store.write().await;
            store.insert_a(data_a);
            drop(store); // 提前释放写锁,减少阻塞时间

            tokio::time::sleep(tokio::time::Duration::from_secs(60)).await;
        }
    })
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 00:07:59