阻塞队列实现


代码搬运地址

#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
*/