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

如何处理无锁可扩容环形缓冲区中的线程竞争问题?

Solution to Lock-Free Resizable Ring Buffer Issues in Rust

Your core issues stem from uncoordinated access to buffer slots and partial state updates that can't be rolled back. Here's a practical approach to fix these problems, using sequence number tracking and atomic slot ownership:

1. Redesign the Queue and Slot Structure

Replace raw indices with monotonic sequence numbers for head/tail, and add atomic sequence markers to each slot to track validity. This eliminates race conditions where multiple threads write to the same slot.

use std::sync::atomic::{AtomicU64, AtomicPtr, Ordering};
use std::mem::MaybeUninit;

struct Slot<T> {
    data: MaybeUninit<T>,
    seq: AtomicU64, // Tracks when the slot was written/freed
}

struct Buffer<T> {
    data: Vec<Slot<T>>,
    capacity: usize,
}

pub struct Queue<T> {
    buffer: AtomicPtr<Buffer<T>>,
    head: AtomicU64, // Next sequence to pop from front
    tail: AtomicU64, // Next sequence to push to back
}

impl<T> Queue<T> {
    pub fn new(initial_capacity: usize) -> Self {
        let mut data = Vec::with_capacity(initial_capacity);
        for i in 0..initial_capacity {
            data.push(Slot {
                data: MaybeUninit::uninit(),
                seq: AtomicU64::new(i as u64),
            });
        }
        let buffer = Box::new(Buffer { data, capacity: initial_capacity });
        Queue {
            buffer: AtomicPtr::new(Box::into_raw(buffer)),
            head: AtomicU64::new(0),
            tail: AtomicU64::new(0),
        }
    }
}

2. Safe Push Back Operation

The push operation atomically claims a slot before writing data, ensuring only one thread can access a slot at a time:

impl<T> Queue<T> {
    pub fn push_back(&self, value: T) {
        loop {
            let buffer_ptr = self.buffer.load(Ordering::Acquire);
            let buffer = unsafe { &*buffer_ptr };
            let tail = self.tail.load(Ordering::Acquire);
            let idx = (tail % buffer.capacity as u64) as usize;
            let slot = &buffer.data[idx];

            // Check if slot is available for writing
            if slot.seq.load(Ordering::Acquire) == tail {
                // Atomically claim the slot
                if slot.seq.compare_exchange_weak(
                    tail,
                    tail + 1,
                    Ordering::Release,
                    Ordering::Relaxed,
                ).is_ok() {
                    // Write data safely
                    unsafe { slot.data.write(value); }
                    // Make the new element visible to readers
                    self.tail.fetch_add(1, Ordering::Release);
                    return;
                }
            } else {
                // Check if queue is full and needs resizing
                let head = self.head.load(Ordering::Acquire);
                if tail - head == buffer.capacity as u64 {
                    self.resize();
                }
                // Retry loop if slot was claimed by another thread
            }
        }
    }
}

3. Safe Pop Front Operation

Pop works similarly, atomically claiming a filled slot before reading:

impl<T> Queue<T> {
    pub fn pop_front(&self) -> Option<T> {
        loop {
            let buffer_ptr = self.buffer.load(Ordering::Acquire);
            let buffer = unsafe { &*buffer_ptr };
            let head = self.head.load(Ordering::Acquire);
            let idx = (head % buffer.capacity as u64) as usize;
            let slot = &buffer.data[idx];
            let slot_seq = slot.seq.load(Ordering::Acquire);

            // Check if slot contains valid data
            if slot_seq == head + 1 {
                // Atomically claim the slot for reading
                if slot.seq.compare_exchange_weak(
                    head + 1,
                    head + buffer.capacity as u64, // Mark as empty for next cycle
                    Ordering::Release,
                    Ordering::Relaxed,
                ).is_ok() {
                    let value = unsafe { slot.data.read() };
                    // Update head to mark the slot as freed
                    self.head.fetch_add(1, Ordering::Release);
                    return Some(value);
                }
            } else {
                // Check if queue is empty
                let tail = self.tail.load(Ordering::Acquire);
                if head == tail {
                    return None;
                }
                // Retry loop if slot was modified by another thread
            }
        }
    }
}

4. Push Front Operation

For pushing to the front, mirror the push back logic but target the slot before the current head:

impl<T> Queue<T> {
    pub fn push_front(&self, value: T) {
        loop {
            let buffer_ptr = self.buffer.load(Ordering::Acquire);
            let buffer = unsafe { &*buffer_ptr };
            let head = self.head.load(Ordering::Acquire);
            // Calculate slot index before current head (handles wrap-around)
            let idx = ((head - 1) % buffer.capacity as u64) as usize;
            let slot = &buffer.data[idx];

            // Check if slot is available
            if slot.seq.load(Ordering::Acquire) == head - 1 {
                if slot.seq.compare_exchange_weak(
                    head - 1,
                    head,
                    Ordering::Release,
                    Ordering::Relaxed,
                ).is_ok() {
                    unsafe { slot.data.write(value); }
                    self.head.fetch_sub(1, Ordering::Release);
                    return;
                }
            } else {
                let tail = self.tail.load(Ordering::Acquire);
                if tail - head == buffer.capacity as u64 {
                    self.resize();
                }
            }
        }
    }
}

5. Resizing the Buffer

Resize safely by creating a new buffer, copying valid elements, and atomically replacing the old buffer. To avoid use-after-free, you'll need hazard pointers or epoch-based reclamation to track when the old buffer is no longer in use:

impl<T> Queue<T> {
    fn resize(&self) {
        let old_buffer_ptr = self.buffer.load(Ordering::Acquire);
        let old_buffer = unsafe { &*old_buffer_ptr };
        let new_capacity = old_buffer.capacity * 2;

        // Initialize new buffer with empty slots
        let mut new_data = Vec::with_capacity(new_capacity);
        let head = self.head.load(Ordering::Acquire);
        for i in 0..new_capacity {
            new_data.push(Slot {
                data: MaybeUninit::uninit(),
                seq: AtomicU64::new(head + i as u64),
            });
        }

        // Copy valid elements from old buffer to new buffer
        let tail = self.tail.load(Ordering::Acquire);
        for seq in head..tail {
            let old_idx = (seq % old_buffer.capacity as u64) as usize;
            let old_slot = &old_buffer.data[old_idx];
            if old_slot.seq.load(Ordering::Acquire) == seq + 1 {
                let new_idx = (seq % new_capacity as u64) as usize;
                let new_slot = &mut new_data[new_idx];
                unsafe { new_slot.data.write(old_slot.data.read()); }
                new_slot.seq.store(seq + 1, Ordering::Release);
            }
        }

        let new_buffer = Box::new(Buffer { data: new_data, capacity: new_capacity });
        let new_buffer_ptr = Box::into_raw(new_buffer);

        // Atomically replace the buffer
        if self.buffer.compare_exchange(
            old_buffer_ptr,
            new_buffer_ptr,
            Ordering::Release,
            Ordering::Relaxed,
        ).is_ok() {
            // SAFETY: Use hazard pointers here to ensure no threads are accessing old_buffer
            // For simplicity, this example skips safe deallocation (add hazard pointers in production)
            // unsafe { Box::from_raw(old_buffer_ptr); }
        } else {
            // Another thread resized first, free the new buffer
            unsafe { Box::from_raw(new_buffer_ptr); }
        }
    }
}

Key Fixes to Your Original Issues

  • No more overwrites: Slots are atomically claimed before writing, so only one thread can access a slot at a time.
  • No partial state rollbacks: Head/tail are only updated after data is safely written/read, so there's no inconsistent state if a CAS fails.
  • Consistent push_front: The push_front operation uses the same atomic slot claiming logic to avoid race conditions.

Critical Notes for Production

  • Hazard Pointers: You must implement hazard pointers (or epoch-based reclamation) to safely free old buffers. Without this, you risk use-after-free bugs when threads access the old buffer after resizing.
  • Memory Ordering: The example uses Acquire/Release orderings to ensure cross-thread visibility of operations. Never use Relaxed unless you're certain no visibility guarantees are needed.
  • Panic Safety: Add guards to handle panics during data writes/reads to avoid leaving slots in an inconsistent state.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 18:47:35