Console Library 8.0.0
A header-only library that makes C++ simple
Loading...
Searching...
No Matches
queue.h
Go to the documentation of this file.
1
9
10/*
11Copyright (c) 2026 MrXie1109
12
13Permission is hereby granted, free of charge, to any person obtaining a copy
14of this software and associated documentation files (the "Software"), to deal
15in the Software without restriction, including without limitation the rights
16to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
17copies of the Software, and to permit persons to whom the Software is
18furnished to do so, subject to the following conditions:
19
20The above copyright notice and this permission notice shall be included in all
21copies or substantial portions of the Software.
22
23THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
24IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
25FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
26AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
27LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
28OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
29SOFTWARE.
30*/
31
32#pragma once
33#include <atomic>
34#include <memory>
35#include <thread>
36#include <utility>
37#include <vector>
38
39namespace console {
48 template <class T, class Alloc = std::allocator<T>>
49 class LFQueue : private Alloc {
50 struct Node {
52 std::atomic<Node *> next_;
53
55 Node(const T &data) : data_(data), next_(nullptr) {}
57 Node(T &&data) : data_(std::move(data)), next_(nullptr) {}
58 };
59
60 std::atomic<Node *> head_;
61 std::atomic<Node *> tail_;
62
64 using NodeAlloc =
65 typename std::allocator_traits<Alloc>::template rebind_alloc<Node>;
67 using NodeTraits = std::allocator_traits<NodeAlloc>;
68
70 Alloc &get_alloc() { return *this; }
71
73 Node *create_node(const T &data) {
74 NodeAlloc node_alloc(get_alloc());
75 Node *node = NodeTraits::allocate(node_alloc, 1);
76 NodeTraits::construct(node_alloc, node, data);
77 return node;
78 }
79
81 void destroy_node(Node *node) {
82 NodeAlloc node_alloc(get_alloc());
83 NodeTraits::destroy(node_alloc, node);
84 NodeTraits::deallocate(node_alloc, node, 1);
85 }
86
87 public:
90 Alloc(Alloc{}), head_(create_node(T())), tail_(head_.load()) {}
91
96 LFQueue(const Alloc &alloc) :
97 Alloc(alloc), head_(create_node(T())), tail_(head_.load()) {}
98
103 LFQueue(Alloc &&alloc) :
104 Alloc(std::move(alloc)), //
105 head_(create_node(T())), tail_(head_.load()) {}
106
114 Node *head;
115 Node *tail;
116 while (!(head = head_.exchange(nullptr, std::memory_order_acq_rel)))
117 std::this_thread::yield();
118 while (!(tail = tail_.exchange(nullptr, std::memory_order_acq_rel)))
119 std::this_thread::yield();
120 while (head) {
121 Node *next = head->next_.load(std::memory_order_relaxed);
122 destroy_node(head);
123 head = next;
124 }
125 (void)tail;
126 }
127
134 void push(const T &data) {
135 Node *new_node = create_node(data);
136 Node *tail;
137 while (!(tail = tail_.exchange(nullptr, std::memory_order_acq_rel)))
138 std::this_thread::yield();
139 tail->next_.store(new_node, std::memory_order_relaxed);
140 tail_.store(new_node, std::memory_order_release);
141 }
142
149 void push(T &&data) {
150 Node *new_node = create_node(std::move(data));
151 Node *tail;
152 while (!(tail = tail_.exchange(nullptr, std::memory_order_acq_rel)))
153 std::this_thread::yield();
154 tail->next_.store(new_node, std::memory_order_relaxed);
155 tail_.store(new_node, std::memory_order_release);
156 }
157
167 template <class Iterator>
168 void push(Iterator begin, size_t count) {
169 if (count == 0) return;
170 Node *head = create_node(*begin), *node = head;
171 try {
172 for (size_t i = 1; i < count; i++) {
173 ++begin;
174 Node *next = create_node(*begin);
175 node->next_.store(next, std::memory_order_relaxed);
176 node = next;
177 }
178 } catch (...) {
179 for (Node *p = head; p != node;
180 /* */ p = p->next_.load(std::memory_order_relaxed))
181 destroy_node(p);
182 destroy_node(node);
183 throw;
184 }
185 Node *tail;
186 while (!(tail = tail_.exchange(nullptr, std::memory_order_acq_rel)))
187 std::this_thread::yield();
188 tail->next_.store(head, std::memory_order_relaxed);
189 tail_.store(node, std::memory_order_release);
190 }
191
198 template <class Iterator>
199 void push(Iterator begin, Iterator end) {
200 push(begin, std::distance(begin, end));
201 }
202
211 bool pop(T &output) {
212 Node *head;
213 while (!(head = head_.exchange(nullptr, std::memory_order_acq_rel)))
214 std::this_thread::yield();
215 if (!head->next_) {
216 head_.store(head, std::memory_order_release);
217 return false;
218 }
219 output
220 = std::move(head->next_.load(std::memory_order_relaxed)->data_);
221 Node *new_head = head->next_.load(std::memory_order_relaxed);
222 head_.store(new_head, std::memory_order_release);
223 destroy_node(head);
224 return true;
225 }
226
232 std::unique_ptr<T> pop() {
233 std::unique_ptr<T> result(new T{});
234 if (pop(*result)) return result;
235 return nullptr;
236 }
237
244 template <class Iterator>
245 size_t pop(Iterator output, size_t count) {
246 Node *head, *tail;
247 while (!(head = head_.exchange(nullptr, std::memory_order_acq_rel)))
248 std::this_thread::yield();
249 tail = head;
250 for (size_t i = 0; i < count; ++i) {
251 if (!tail->next_.load(std::memory_order_relaxed)) break;
252 tail = tail->next_.load(std::memory_order_relaxed);
253 }
254 head_.store(tail, std::memory_order_release);
255 size_t real_count = 0;
256 while (head != tail) {
257 Node *next = head->next_.load(std::memory_order_relaxed);
258 *output++ = std::move(next->data_);
259 destroy_node(head);
260 head = next;
261 ++real_count;
262 }
263 return real_count;
264 }
265
273 template <class Iterator>
274 size_t pop(Iterator begin, Iterator end) {
275 return pop(begin, std::distance(begin, end));
276 }
277
284 std::vector<T> pop(size_t count) {
285 std::vector<T> result;
286 result.reserve(count);
287 pop(std::back_inserter(result), count);
288 return result;
289 }
290 };
291
302 template <class T, class Alloc = std::allocator<T>, size_t RobinTimes = 2>
305 std::vector<std::unique_ptr<console::LFQueue<T, Alloc>>> queues_;
307 std::atomic<size_t> round_robin_{0};
309 std::atomic<size_t> nth_{0};
310
315 size_t _index() {
316 thread_local size_t idx = nth_.fetch_add(1) % queues_.size();
317 return idx;
318 }
319
325 template <class OtherIterator, class IteratorTag>
326 void _advance(OtherIterator &it, size_t n, IteratorTag) {
327 for (size_t i = 0; i < n; ++i) ++it;
328 }
329
335 template <class RandomAccessIterator>
336 void _advance(RandomAccessIterator &it,
337 size_t n,
338 std::random_access_iterator_tag) {
339 it += n;
340 }
341
347 template <class Iterator>
348 void _advance(Iterator &it, size_t n) {
349 _advance(it,
350 n,
351 typename std::iterator_traits<Iterator>::iterator_category{});
352 }
353
354 public:
359 MultiLFQueue(size_t num_queues = 1) {
360 for (size_t i = 0; i < num_queues; ++i)
361 queues_.emplace_back(new console::LFQueue<T, Alloc>());
362 }
363
371 template <class... Args>
372 auto push(Args &&...args)
373 -> decltype(queues_[_index()]->push(std::forward<Args>(args)...)) {
374 return queues_[_index()]->push(std::forward<Args>(args)...);
375 }
376
385 bool pop(T &output) {
386 size_t start = round_robin_.fetch_add(1, std::memory_order_relaxed);
387 for (size_t i = 0; i < queues_.size() * RobinTimes; ++i) {
388 size_t idx = (start + i) % queues_.size();
389 if (queues_[idx]->pop(output)) return true;
390 }
391 return false;
392 }
393
399 std::unique_ptr<T> pop() {
400 std::unique_ptr<T> result(new T{});
401 if (pop(*result)) return result;
402 return nullptr;
403 }
404
411 template <class Iterator>
412 size_t pop(Iterator output, size_t count) {
413 size_t start = round_robin_.fetch_add(1, std::memory_order_relaxed);
414 size_t popped = 0;
415 for (size_t i = 0; i < queues_.size() * RobinTimes; ++i) {
416 size_t idx = (start + i) % queues_.size();
417 size_t popped_this = queues_[idx]->pop(output, count);
418 _advance(output, popped_this);
419 popped += popped_this;
420 count -= popped_this;
421 if (count == 0) break;
422 }
423 return popped;
424 }
425
433 template <class Iterator>
434 size_t pop(Iterator begin, Iterator end) {
435 return pop(begin, std::distance(begin, end));
436 }
437
444 std::vector<T> pop(size_t count) {
445 std::vector<T> result;
446 result.reserve(count);
447 pop(std::back_inserter(result), count);
448 return result;
449 }
450 };
451}
基于自旋的高性能 FIFO 队列,支持多生产者/多消费者并发。
Definition queue.h:49
std::allocator_traits< NodeAlloc > NodeTraits
节点分配器的特征类型。
Definition queue.h:67
void destroy_node(Node *node)
销毁节点。
Definition queue.h:81
std::atomic< Node * > tail_
指向队列尾部的原子指针
Definition queue.h:61
std::atomic< Node * > head_
指向队列头部的原子指针
Definition queue.h:60
void push(Iterator begin, size_t count)
将大批数据推送到队列中。
Definition queue.h:168
Alloc & get_alloc()
获取分配器的引用。
Definition queue.h:70
typename std::allocator_traits< Alloc >::template rebind_alloc< Node > NodeAlloc
节点分配器类型。
Definition queue.h:64
void push(T &&data)
将数据推送到队列中。
Definition queue.h:149
void push(Iterator begin, Iterator end)
将大批数据推送到队列中。
Definition queue.h:199
void push(const T &data)
将数据推送到队列中。
Definition queue.h:134
LFQueue()
默认构造函数,使用默认分配器初始化队列。
Definition queue.h:89
LFQueue(const Alloc &alloc)
使用指定分配器初始化队列。
Definition queue.h:96
std::vector< T > pop(size_t count)
从队列中弹出数据大批数据,返回包含弹出结果的 vector。
Definition queue.h:284
LFQueue(Alloc &&alloc)
使用移动分配器初始化队列。
Definition queue.h:103
size_t pop(Iterator begin, Iterator end)
从队列中弹出数据大批数据,写入输出迭代器。
Definition queue.h:274
Node * create_node(const T &data)
创建一个新的节点。
Definition queue.h:73
std::unique_ptr< T > pop()
从队列中弹出数据,并返回一个独占所有权的 std::unique_ptr。
Definition queue.h:232
~LFQueue()
析构函数,销毁队列中的所有节点。
Definition queue.h:113
size_t pop(Iterator output, size_t count)
从队列中弹出数据大批数据,写入输出迭代器。
Definition queue.h:245
bool pop(T &output)
从队列中弹出数据。
Definition queue.h:211
size_t pop(Iterator output, size_t count)
从队列中弹出大批数据,写入输入迭代器。
Definition queue.h:412
size_t _index()
获取当前线程的索引,用于推送时的负载均衡。
Definition queue.h:315
std::atomic< size_t > nth_
序号计数器,用于推送时的线程索引分配。
Definition queue.h:309
void _advance(OtherIterator &it, size_t n, IteratorTag)
类似标准库的 advance。
Definition queue.h:326
bool pop(T &output)
从队列中弹出数据。
Definition queue.h:385
MultiLFQueue(size_t num_queues=1)
构造函数,初始化指定数量的子队列。
Definition queue.h:359
auto push(Args &&...args) -> decltype(queues_[_index()]->push(std::forward< Args >(args)...))
将数据推送到队列中。
Definition queue.h:372
std::atomic< size_t > round_robin_
轮询计数器,用于弹出时的负载均衡。
Definition queue.h:307
std::vector< T > pop(size_t count)
从队列中弹出数据大批数据,返回包含弹出结果的 vector。
Definition queue.h:444
std::unique_ptr< T > pop()
从队列中弹出数据,并返回一个独占所有权的 std::unique_ptr。
Definition queue.h:399
void _advance(RandomAccessIterator &it, size_t n, std::random_access_iterator_tag)
类似标准库的 advance。
Definition queue.h:336
void _advance(Iterator &it, size_t n)
类似标准库的 advance。
Definition queue.h:348
std::vector< std::unique_ptr< console::LFQueue< T, Alloc > > > queues_
子队列集合,每个子队列用于存储不同线程的队列数据。
Definition queue.h:305
size_t pop(Iterator begin, Iterator end)
从队列中弹出数据大批数据,写入输出迭代器。
Definition queue.h:434
本库所有组件所在的顶层命名空间。
@ T
Definition kb.h:99
T next(typename Generator< Derived, T >::iterator &it)
从迭代器获取下一个值的辅助函数。
Definition gen.h:190
Definition queue.h:50
T data_
节点数据
Definition queue.h:51
Node(const T &data)
构造函数,使用数据初始化节点。
Definition queue.h:55
std::atomic< Node * > next_
指向下一个节点的原子指针
Definition queue.h:52
Node(T &&data)
构造函数,使用移动数据初始化节点。
Definition queue.h:57