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

《C++ Concurrency in Action》无锁队列元素丢失问题排查求助

《C++ Concurrency in Action》无锁队列元素丢失问题排查求助

我在阅读《C++ Concurrency in Action》(第二版)第7章时,测试书中提供的无锁队列实现,偶尔会出现入队元素丢失的情况。测试中发现,入队完成后本应执行的set_new_tail函数if块内的原子计数,偶尔会比实际入队操作数少1-2次,无法定位元素丢失的根本原因,特此求助。

无锁队列实现代码

template <typename T>
class lock_free_queue
{
private:
    struct node;
    struct counted_node_ptr
    {
        int external_count;
        node *ptr;
        uint8_t padding[sizeof(void *) - sizeof(int)];
    };

    std::atomic<counted_node_ptr> _head;
    std::atomic<counted_node_ptr> _tail;

    struct node_counter
    {
        unsigned internal_count : 30;
        unsigned external_counters : 2;
    };

    struct node
    {
        std::atomic<T *> data;
        std::atomic<node_counter> count;
        std::atomic<counted_node_ptr> next;

        node()
        {
            data.store(nullptr);
            node_counter new_count;
            new_count.internal_count = 0;
            new_count.external_counters = 2;
            count.store(new_count);
            counted_node_ptr new_next{0, nullptr, {0}};
            next.store(new_next);
        }

        void release_ref()
        {
            node_counter old_counter = count.load(std::memory_order_relaxed);
            node_counter new_counter;
            do
            {
                new_counter = old_counter;
                --new_counter.internal_count;
            } while (!count.compare_exchange_strong(old_counter, new_counter,
                                                    std::memory_order_acquire,
                                                    std::memory_order_relaxed));
            if (!new_counter.internal_count && !new_counter.external_counters)
            {
                delete this;
            }
        }
    };

    static void increase_external_count(std::atomic<counted_node_ptr> &counter,
                                        counted_node_ptr &old_counter)
    {
        counted_node_ptr new_counter;
        do
        {
            new_counter = old_counter;
            ++new_counter.external_count;
        } while (!counter.compare_exchange_strong(old_counter, new_counter,
                                                  std::memory_order_acquire,
                                                  std::memory_order_relaxed));
        old_counter.external_count = new_counter.external_count;
    }

    static void free_external_counter(counted_node_ptr &old_node_ptr)
    {
        node *const ptr = old_node_ptr.ptr;
        int const count_increase = old_node_ptr.external_count - 2;
        node_counter old_counter = ptr->count.load(std::memory_order_relaxed);
        node_counter new_counter;
        do
        {
            new_counter = old_counter;
            --new_counter.external_counters;
            new_counter.internal_count += count_increase;
        } while (!ptr->count.compare_exchange_strong(old_counter, new_counter,
                                                     std::memory_order_acquire,
                                                     std::memory_order_relaxed));
        if (!new_counter.internal_count && !new_counter.external_counters)
        {
            delete ptr;
        }
    }

    void set_new_tail(counted_node_ptr &old_tail,
                      counted_node_ptr const &new_tail)
    {
        node *const current_tail_ptr = old_tail.ptr;
        while (!_tail.compare_exchange_weak(old_tail, new_tail) && old_tail.ptr == current_tail_ptr)
            ;
        if (old_tail.ptr == current_tail_ptr)
        {
            free_external_counter(old_tail);
        }
        else
        {
            current_tail_ptr->release_ref();
        }
    }

public:
    lock_free_queue()
    {
        counted_node_ptr dummy{1, new node, {0}};
        _head.store(dummy);
        _tail.store(dummy);
    }

    void push(T new_value)
    {
        std::unique_ptr<T> new_data(new T(new_value));
        counted_node_ptr new_next;
        new_next.ptr = new node;
        new_next.external_count = 1;
        counted_node_ptr old_tail = _tail.load();
        for (;;)
        {
            increase_external_count(_tail, old_tail);
            T *old_data = nullptr;
            if (old_tail.ptr->data.compare_exchange_strong(
                    old_data, new_data.get()))
            {
                counted_node_ptr old_next = {0, nullptr, {0}};
                if (!old_tail.ptr->next.compare_exchange_strong(
                        old_next, new_next))
                {
                    delete new_next.ptr;
                    new_next = old_next;
                }
                set_new_tail(old_tail, new_next);
                new_data.release();
                break;
            }
            else
            {
                counted_node_ptr old_next = {0, nullptr, {0}};
                if (old_tail.ptr->next.compare_exchange_strong(
                        old_next, new_next))
                {
                    old_next = new_next;
                    new_next.ptr = new node;
                }
                set_new_tail(old_tail, old_next);
            }
        }
    }

    std::unique_ptr<T> pop()
    {
        counted_node_ptr old_head = _head.load(std::memory_order_relaxed);
        for (;;)
        {
            increase_external_count(_head, old_head);
            node *const ptr = old_head.ptr;
            if (ptr == _tail.load().ptr)
            {
                ptr->release_ref(); //
                return std::unique_ptr<T>();
            }
            counted_node_ptr next = ptr->next.load();
            if (_head.compare_exchange_strong(old_head, next))
            {
                T *const res = ptr->data.exchange(nullptr);
                free_external_counter(old_head);
                return std::unique_ptr<T>(res);
            }
            ptr->release_ref();
        }
    }
};

测试代码

void t1_producer(lock_free_queue<int> &queue, int start, int count)
{
    for (int i = start; i < start + count; ++i)
    {
        queue.push(i);
    }
}

void t1_consumer(lock_free_queue<int> &queue, std::atomic<int> &sum, int consume_count)
{
    for (int i = 0; i < consume_count;)
    {
        auto value = queue.pop();
        if (value)
        {
            sum += *value;
            ++i;
        }
    }
}

int main()
{
    for (int i = 0; i < 100; i++)
    {
        lock_free_queue<int> queue;
        std::atomic<int> sum(0);

        int num_producers = 4;
        int num_consumers = 4;
        int items_per_producer = 1000;

        std::vector<std::thread> producers;
        std::vector<std::thread> consumers;

        for (int i = 0; i < num_producers; ++i)
        {
            producers.emplace_back(t1_producer, std::ref(queue), i * items_per_producer, items_per_producer);
        }
        
        for (int i = 0; i < num_consumers; ++i)
        {
            consumers.emplace_back(t1_consumer, std::ref(queue), std::ref(sum), (num_producers * items_per_producer) / num_consumers);
        }

        for (auto &t : producers)
        {
            t.join();
        }

        for (auto &t : consumers)
        {
            t.join();
        }

        int expected_sum = (num_producers * items_per_producer * (items_per_producer * num_producers - 1)) / 2;
        std::cout << "Expected sum: " << expected_sum << "\n";
        std::cout << "Actual   sum: " << sum.load() << "\n";

        if (sum.load() == expected_sum)
        {
            std::cout << "Test passed: Queue is thread-safe.\n";
        }
        else
        {
            std::cout << "Test failed: Queue is not thread-safe.\n";
        }
    }
}

排查过程

根据书中说明,元素完成入队操作后会进入set_new_tail函数并执行if块,因此我在该if块内添加了原子计数,修改后的函数代码如下:

void set_new_tail(counted_node_ptr &old_tail,
                  counted_node_ptr const &new_tail)
{
    node *const current_tail_ptr = old_tail.ptr;
    while (!_tail.compare_exchange_weak(old_tail, new_tail) && old_tail.ptr == current_tail_ptr);
    if (old_tail.ptr == current_tail_ptr)
    {
        // increase counter here...
        free_external_counter(old_tail);
    }
    else
    {
        current_tail_ptr->release_ref();
    }
}

测试中发现,该计数偶尔会比实际入队操作数少1-2次,无法定位元素丢失的根本原因。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 02:44:53