Spec-Zone.ru › Ruby 4.0
  1. Thread::
  2. Queue

class Thread::Queue

Родительский класс:
Object

Класс 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

Методы класса

Thread::Queue.new → empty_queue Показать исходный код
Thread::Queue.new(enumerable) → queue
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

Методы экземпляра

<<(object)

Добавляет указанный object в очередь.

Псевдоним для: push
clear () Показать исходный код
static VALUE
rb_queue_clear(VALUE self)
{
    struct rb_queue *q = queue_ptr(self);

    rb_ary_clear(check_array(self, q->que));
    return self;
}

Удаляет все объекты из очереди.

close Показать исходный код
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
closed? Показать исходный код
static VALUE
rb_queue_closed_p(VALUE self)
{
    return RBOOL(queue_closed_p(self));
}

Возвращает true, если очередь закрыта.

deq
Псевдоним для: pop
empty? Показать исходный код
static VALUE
rb_queue_empty_p(VALUE self)
{
    return RBOOL(queue_length(self, queue_ptr(self)) == 0);
}

Возвращает true, если очередь пуста.

enq(object)

Добавляет указанный object в очередь.

Псевдоним для: push
freeze Показать исходный код
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...>)
length Показать исходный код
static VALUE
rb_queue_length(VALUE self)
{
    return LONG2NUM(queue_length(self, queue_ptr(self)));
}

Возвращает длину очереди.

Также имеет псевдоним: size
num_waiting () Показать исходный код
static VALUE
rb_queue_num_waiting(VALUE self)
{
    struct rb_queue *q = queue_ptr(self);

    return INT2NUM(q->num_waiting);
}

Возвращает количество потоков, ожидающих в очереди.

pop(non_block=false, timeout: nil) Показать исходный код
# File thread_sync.rb, line 16
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, метод возвращает результат немедленно.

Также имеет псевдонимы: deq, shift
push(object) Показать исходный код
static VALUE
rb_queue_push(VALUE self, VALUE obj)
{
    return queue_do_push(self, queue_ptr(self), obj);
}

Добавляет указанный object в очередь.

Также имеет псевдонимы: enq, <<
shift
Псевдоним для: pop
size

Возвращает длину очереди.

Псевдоним для: length

Ruby Core © 1993–2025 Yukihiro Matsumoto
Licensed under the Ruby License.
Ruby Standard Library © contributors
Licensed under their own licenses.

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API