класс Thread::Queue
Класс Thread::Queue реализует очереди с множественным производством и множественным потреблением. Он особенно полезен в многопоточной программировании, когда информация должна безопасно обмениваться между несколькими потоками. Класс Thread::Queue реализует все необходимые семантики блокировок.
Класс реализует очередь типа FIFO (первым вошел, первым вышел). В очереди FIFO первые добавленные задачи — первые извлеченные.
Пример:
queue = Thread::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(int argc, VALUE *argv, VALUE self)
{
VALUE initial;
struct rb_queue *q = queue_ptr(self);
if ((argc = rb_scan_args(argc, argv, "01", &initial)) == 1) {
initial = rb_to_array(initial);
}
RB_OBJ_WRITE(self, queue_list(q), ary_buf_new());
ccan_list_head_init(queue_waitq(q));
if (argc == 1) {
rb_ary_concat(q->que, initial);
}
return self;
} Создаёт новый экземпляр очереди, необязательно используя содержимое enumerable для его начального состояния.
Пример:
q = Thread::Queue.new #=> #<Thread::Queue:0x00007ff7501110d0> q.empty? #=> true q = Thread::Queue.new([1, 2, 3]) #=> #<Thread::Queue:0x00007ff7500ec500> q.empty? #=> false q.pop #=> 1
Публичные методы экземпляра
Исходный код
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, чтобы вы могли прервать цикл.
Пример:
q = Thread::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 RBOOL(queue_closed_p(self));
} Возвращает true если очередь закрыта.
Исходный код
static VALUE
rb_queue_empty_p(VALUE self)
{
return RBOOL(queue_length(self, queue_ptr(self)) == 0);
} Возвращает true если очередь пуста.
Исходный код
static VALUE
rb_queue_freeze(VALUE self)
{
rb_raise(rb_eTypeError, "cannot freeze " "%+"PRIsVALUE, self);
UNREACHABLE_RETURN(self);
} Очередь нельзя заморозить, поэтому этот метод вызывает исключение:
Thread::Queue.new.freeze # Raises TypeError (cannot freeze #<Thread::Queue:0x...>)
Исходный код
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);
} Возвращает количество потоков, ожидающих в очереди.
Исходный код
# File thread_sync.rb, line 14
def pop(non_block = false, timeout: nil)
if non_block && timeout
raise ArgumentError, "can't set a timeout if non_block is enabled"
end
Primitive.rb_queue_pop(non_block, timeout)
end Извлекает данные из очереди.
Если очередь пуста, вызывающий поток приостанавливается до тех пор, пока данные не будут добавлены в очередь. Если non_block равно true, поток не приостанавливается, и генерируется исключение ThreadError.
Если прошло timeout секунд и данные недоступны, возвращается nil. Если timeout равно 0, возвращается сразу.
Исходный код
static VALUE
rb_queue_push(VALUE self, VALUE obj)
{
return queue_do_push(self, queue_ptr(self), obj);
} Добавляет указанный object в очередь.
Ruby Core © 1993–2024 Yukihiro Matsumoto
Licensed under the Ruby License.
Ruby Standard Library © contributors
Licensed under their own licenses.