阻塞队列实现
代码搬运地址
#ifndef BLOCKQUEUE_HPP
#define BLOCKQUEUE_HPP
#include
#include
#include
#include
#include
//阻塞队列
template
class BlockDeque
{
public:
explicit BlockDeque(size_t MaxCapacity = 1000);
~BlockDeque();
void Close();
void clear();
bool empty();
bool full();
size_t size();
size_t capacity();
T front();
T back();
void push_back(const T& item);
void push_front(const T& item);
bool pop(T& item);
bool pop(T& item, int timeout);
void flush();
private:
std::deque _deq;
size_t _capacity;
bool _isClose;
std::mutex _mtx;
std::condition_variable _consumer;
std::condition_variable _producer;
};
template
BlockDeque::BlockDeque(size_t MaxCapacity)
: _capacity(MaxCapacity)
{
assert(MaxCapacity > 0);
_isClose = false;
}
template
BlockDeque::~BlockDeque()
{
Close();
}
template
void
BlockDeque::Close()
{
{
std::lock_guard locker(_mtx);
_deq.clear();
_isClose = true;
}
_producer.notify_all();
_consumer.notify_all();
}
template
void
BlockDeque::flush()
{
_consumer.notify_one();
}
template
void
BlockDeque::clear()
{
std::lock_guard locker(_mtx);
_deq.clear();
}
template
bool
BlockDeque::empty()
{
std::lock_guard locker(_mtx);
return _deq.empty();
}
template
bool
BlockDeque::full()
{
std::lock_guard locker(_mtx);
return _deq.size() >= _capacity;
}
template
size_t
BlockDeque::size()
{
std::lock_guard locker(_mtx);
return _deq.size();
}
template
size_t
BlockDeque::capacity()
{
std::lock_guard locker(_mtx);
return _capacity;
}
template
T
BlockDeque::front()
{
std::lock_guard locker(_mtx);
return _deq.front();
}
template
T
BlockDeque::back()
{
std::lock_guard locker(_mtx);
return _deq.back();
}
template
void
BlockDeque::push_back(const T& item)
{
std::unique_lock locker(_mtx);
while (_deq.size() >= _capacity) //生产满了,等待消费
{
_producer.wait(locker);
}
_deq.push_back(item);
_consumer.notify_one(); //通知一个消费者
}
template
void
BlockDeque::push_front(const T& item)
{
std::unique_lock locker(_mtx);
while (_deq.size() >= _capacity) {
_producer.wait(locker);
}
_deq.push_front(item);
_consumer.notify_one();
}
template
bool
BlockDeque::pop(T& item)
{
std::unique_lock locker(_mtx);
while (_deq.empty()) {
//线程阻塞,解锁,等待唤醒; 唤醒后不断尝试加锁,直至成功
_consumer.wait(locker);
if (_isClose) {
return false;
}
}
item = _deq.front();
_deq.pop_front();
_producer.notify_one();
return true;
}
template
bool
BlockDeque::pop(T& item, int timeout)
{
std::unique_lock locker(_mtx);
while (_deq.empty()) {
if (_consumer.wait_for(locker, std::chrono::seconds(timeout)) ==
std::cv_status::timeout) {
return false;
}
if (_isClose) {
return false;
}
}
item = _deq.front();
_deq.pop_front();
_producer.notify_one();
return true;
}
#endif // BLOCKQUEUE_HPP
/*vim 替换
% 当前文件(省略表示当前行)
s 替换
%s/a/b 当前文件 /a替换成/b
*/