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

crossbeam_epoch::Atomic是否真具备原子性?多线程测试结果不符预期

问题:crossbeam_epoch::Atomic多线程递增操作结果不符合预期

我需要全局存储一个以读操作为主、写操作为辅的复杂结构体作为配置,因此尝试用crossbeam_epoch::Atomic实现多线程共享,同时测试其读写的原子性。

我创建了包含version: i32属性的Config结构体,启动5个线程,每个线程循环100次递增配置的version值,预期最终结果为500,但实际结果随机。测试代码如下:

use crossbeam_epoch as epoch;
use std::sync::atomic::Ordering;
use std::{sync::Arc, thread};

#[derive(Debug)]
struct Config {
    version: i32,
}

fn run() -> () {
    let atomic = Arc::new(epoch::Atomic::new(Config { version: 0 }));

    let mut handles = vec![];

    let per_thread_itr = 100;
    let num_thread = 5;

    for _ in 0..num_thread {
        let thread_handle = thread::spawn({
            let atomic = atomic.clone();
            move || {
                let guard = epoch::pin();
                for _ in 1..per_thread_itr + 1 {
                    let shared = atomic.load(Ordering::SeqCst, &guard);
                    if let Some(v) = unsafe { shared.as_ref() } {
                        atomic.store(
                            epoch::Owned::new(Config {
                                version: v.version + 1,
                            }),
                            Ordering::SeqCst,
                        );
                    }
                }
            }
        });

        handles.push(thread_handle);
    }

    for handle in handles {
        handle.join().unwrap();
    }

    // Main thread accesses the updated value
    let guard = epoch::pin();
    let result = unsafe { atomic.load(Ordering::SeqCst, &guard).as_ref() }.unwrap();

    println!("Result: {}", result.version);
}

fn main() {
    for _ in 0..5 {
        run();
    }
}

测试结果:

Result: 297
Result: 398
Result: 349
Result: 235
Result: 365

使用std::sync::atomic::AtomicI32的对比测试

采用相同逻辑用AtomicI32测试时,每次结果均为500。测试代码如下:

use std::sync::atomic::AtomicI32;
use std::sync::atomic::Ordering;
use std::sync::Arc;
use std::thread;

fn run() {
    // Use AtomicI32 with proper atomic operations
    let atomic = Arc::new(AtomicI32::new(0));

    let mut handles = vec![];

    let per_thread_itr = 100;
    let num_thread = 5;

    for _ in 0..num_thread {
        let a = atomic.clone();
        let thread_handle = thread::spawn(move || {
            for _ in 0..per_thread_itr {
                // Safely increment the atomic value
                a.fetch_add(1, Ordering::SeqCst);
            }
        });

        handles.push(thread_handle);
    }

    for handle in handles {
        handle.join().unwrap();
    }

    println!("Result: {}", atomic.load(Ordering::SeqCst));
}

fn main() {
    for _ in 0..5 {
        run();
    }
}

测试结果:

Result: 500
Result: 500
Result: 500
Result: 500
Result: 500

我想知道crossbeam_epoch::Atomic与std::sync::atomic::AtomicI32的工作机制是否存在差异?还是我在使用crossbeam_epoch::Atomic时存在疏漏?


解答

两者工作机制的核心差异

  1. 针对的对象类型不同:

    • std::sync::atomic::AtomicI32是专门针对单个i32值的原子类型,提供了fetch_add这类原子性的读-改-写操作,底层直接利用CPU的原子指令保证操作的不可分割性。
    • crossbeam_epoch::Atomic是用于管理堆分配对象的原子指针,核心能力是支持无锁的内存回收(epoch-based reclamation),而非原生支持复杂对象内部字段的原子操作。它的load和store仅保证指针层面的原子性,不负责对象内部数据的原子更新。
  2. 操作语义不同:

    • AtomicI32的fetch_add是单一原子操作,从读取当前值、加1到写回新值的全过程不会被其他线程打断。
    • 你的crossbeam_epoch::Atomic代码中,load读取指针、修改对象字段、再store新对象是三个独立操作,中间无原子性保障——多个线程可能同时读取到同一个旧版本的Config,各自加1后写回,导致更新丢失。

代码中的疏漏

你错误地将crossbeam_epoch::Atomic的指针原子性,等同于对象内部数据的原子更新能力。当前代码存在典型竞态条件:多个线程同时读取同一个Config实例,各自计算出version+1后覆盖写入新实例,导致部分线程的更新被覆盖。

修正方案:使用原子比较并交换(CAS)操作

要实现crossbeam_epoch::Atomic下的原子递增,需利用compare_and_set方法(对应CAS操作),确保只有当当前指针指向的实例与读取的旧实例一致时,才替换为新实例。若CAS失败(说明其他线程已更新指针),则重试操作。

修正后的代码示例:

use crossbeam_epoch as epoch;
use std::sync::atomic::Ordering;
use std::{sync::Arc, thread};

#[derive(Debug, Clone)]
struct Config {
    version: i32,
}

fn run() -> () {
    let atomic = Arc::new(epoch::Atomic::new(Config { version: 0 }));

    let mut handles = vec![];

    let per_thread_itr = 100;
    let num_thread = 5;

    for _ in 0..num_thread {
        let thread_handle = thread::spawn({
            let atomic = atomic.clone();
            move || {
                let guard = epoch::pin();
                for _ in 1..per_thread_itr + 1 {
                    loop {
                        // 加载当前的共享实例
                        let shared = atomic.load(Ordering::SeqCst, &guard);
                        let old_config = unsafe { shared.as_ref() }.unwrap();
                        
                        // 创建新的配置实例
                        let new_config = Config {
                            version: old_config.version + 1,
                        };
                        
                        // 尝试CAS替换:只有当前指针仍指向old_config时才成功
                        match atomic.compare_and_set(
                            shared,
                            epoch::Owned::new(new_config),
                            Ordering::SeqCst,
                            &guard,
                        ) {
                            Ok(_) => break, // CAS成功,退出循环进行下一次递增
                            Err(_) => continue, // CAS失败,重试操作
                        }
                    }
                }
            }
        });

        handles.push(thread_handle);
    }

    for handle in handles {
        handle.join().unwrap();
    }

    // 主线程读取结果
    let guard = epoch::pin();
    let result = unsafe { atomic.load(Ordering::SeqCst, &guard).as_ref() }.unwrap();

    println!("Result: {}", result.version);
}

fn main() {
    for _ in 0..5 {
        run();
    }
}

这段代码通过循环执行CAS操作,确保每次递增都是原子性的,最终结果会稳定为500。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 08:28:11