48 template <
class T,
class Alloc = std::allocator<T>>
65 typename std::allocator_traits<Alloc>::template rebind_alloc<Node>;
75 Node *node = NodeTraits::allocate(node_alloc, 1);
76 NodeTraits::construct(node_alloc, node, data);
83 NodeTraits::destroy(node_alloc, node);
84 NodeTraits::deallocate(node_alloc, node, 1);
104 Alloc(std::move(alloc)),
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();
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);
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);
167 template <
class Iterator>
168 void push(Iterator begin,
size_t count) {
169 if (count == 0)
return;
172 for (
size_t i = 1; i < count; i++) {
175 node->next_.store(
next, std::memory_order_relaxed);
179 for (
Node *p = head; p != node;
180 p = p->
next_.load(std::memory_order_relaxed))
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);
198 template <
class Iterator>
199 void push(Iterator begin, Iterator end) {
200 push(begin, std::distance(begin, end));
213 while (!(head =
head_.exchange(
nullptr, std::memory_order_acq_rel)))
214 std::this_thread::yield();
216 head_.store(head, std::memory_order_release);
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);
232 std::unique_ptr<T>
pop() {
233 std::unique_ptr<T> result(
new T{});
234 if (
pop(*result))
return result;
244 template <
class Iterator>
245 size_t pop(Iterator output,
size_t count) {
247 while (!(head =
head_.exchange(
nullptr, std::memory_order_acq_rel)))
248 std::this_thread::yield();
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);
254 head_.store(tail, std::memory_order_release);
255 size_t real_count = 0;
256 while (head != tail) {
258 *output++ = std::move(
next->data_);
273 template <
class Iterator>
274 size_t pop(Iterator begin, Iterator end) {
275 return pop(begin, std::distance(begin, end));
284 std::vector<T>
pop(
size_t count) {
285 std::vector<T> result;
286 result.reserve(count);
287 pop(std::back_inserter(result), count);
302 template <
class T,
class Alloc = std::allocator<T>,
size_t RobinTimes = 2>
305 std::vector<std::unique_ptr<console::LFQueue<T, Alloc>>>
queues_;
316 thread_local size_t idx =
nth_.fetch_add(1) %
queues_.size();
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;
335 template <
class RandomAccessIterator>
338 std::random_access_iterator_tag) {
347 template <
class Iterator>
351 typename std::iterator_traits<Iterator>::iterator_category{});
360 for (
size_t i = 0; i < num_queues; ++i)
371 template <
class... Args>
373 ->
decltype(
queues_[
_index()]->push(std::forward<Args>(args)...)) {
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();
399 std::unique_ptr<T>
pop() {
400 std::unique_ptr<T> result(
new T{});
401 if (
pop(*result))
return result;
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);
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);
419 popped += popped_this;
420 count -= popped_this;
421 if (count == 0)
break;
433 template <
class Iterator>
434 size_t pop(Iterator begin, Iterator end) {
435 return pop(begin, std::distance(begin, end));
444 std::vector<T>
pop(
size_t count) {
445 std::vector<T> result;
446 result.reserve(count);
447 pop(std::back_inserter(result), count);
基于自旋的高性能 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 next(typename Generator< Derived, T >::iterator &it)
从迭代器获取下一个值的辅助函数。
Definition gen.h:190
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