《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
相关产品推荐
相关产品推荐

