класс Queue
Класс Queue реализует очереди с множественными производителями и потребителями. Он особенно полезен в многопоточной программировании, когда обмен информацией между несколькими потоками должен быть безопасным. Класс Queue реализует всю необходимую семантику блокировки.
Класс реализует очередь типа FIFO. В очереди FIFO первые добавленные задачи — первые извлечённые.
Пример:
queue = Queue.new
producer = Thread.new do
5.times do |i|
sleep rand(i) # simulate expense
queue << i
puts "#{i} produced"
end
end
consumer = Thread.new do
5.times do |i|
value = queue.pop
sleep rand(i/2) # simulate expense
puts "consumed #{value}"
end
end
consumer.join
Публичные методы класса
static VALUE
rb_queue_initialize(VALUE self)
{
struct rb_queue *q = queue_ptr(self);
RB_OBJ_WRITE(self, &q->que, ary_buf_new());
list_head_init(queue_waitq(q));
return self;
} Создаёт новый экземпляр очереди.
Публичные методы экземпляра
static VALUE
rb_queue_clear(VALUE self)
{
struct rb_queue *q = queue_ptr(self);
rb_ary_clear(check_array(self, q->que));
return self;
} Удаляет все объекты из очереди.
static VALUE
rb_queue_close(VALUE self)
{
struct rb_queue *q = queue_ptr(self);
if (!queue_closed_p(self)) {
FL_SET(self, QUEUE_CLOSED);
wakeup_all(queue_waitq(q));
}
return self;
} Закрывает очередь. Закрытую очередь нельзя повторно открыть.
После выполнения close следующие условия выполняются:
-
closed?вернёт true -
closeбудут проигнорированы. -
вызов enq/push/<< вызовет
ClosedQueueError. -
когда
empty?равно false, вызов deq/pop/shift вернёт объект из очереди как обычно. -
когда
empty?равно true, deq(false) не приостановит поток и вернёт nil. deq(true) вызоветThreadError.
ClosedQueueError унаследована от StopIteration, поэтому вы можете прервать цикл.
Example:
q = Queue.new
Thread.new{
while e = q.deq # wait for nil to break loop
# ...
end
}
q.close static VALUE
rb_queue_closed_p(VALUE self)
{
return queue_closed_p(self) ? Qtrue : Qfalse;
} Возвращает true, если очередь закрыта.
static VALUE
rb_queue_empty_p(VALUE self)
{
return queue_length(self, queue_ptr(self)) == 0 ? Qtrue : Qfalse;
} Возвращает true, если очередь пуста.
static VALUE
rb_queue_length(VALUE self)
{
return LONG2NUM(queue_length(self, queue_ptr(self)));
} Возвращает длину очереди.
static VALUE
rb_queue_num_waiting(VALUE self)
{
struct rb_queue *q = queue_ptr(self);
return INT2NUM(q->num_waiting);
} Возвращает количество потоков, ожидающих в очереди.
static VALUE
rb_queue_pop(int argc, VALUE *argv, VALUE self)
{
int should_block = queue_pop_should_block(argc, argv);
return queue_do_pop(self, queue_ptr(self), should_block);
} Извлекает данные из очереди.
Если очередь пуста, вызывающий поток приостанавливается, пока данные не будут добавлены в очередь. Если non_block равно true, поток не приостанавливается и возникает ThreadError.
static VALUE
rb_queue_push(VALUE self, VALUE obj)
{
return queue_do_push(self, queue_ptr(self), obj);
} Добавляет указанный object в очередь.
Ruby Core © 1993–2017 Yukihiro Matsumoto
Licensed under the Ruby License.
Ruby Standard Library © contributors
Licensed under their own licenses.