Spec-Zone.ru › Ruby 2.2

класс Thread

Родитель:
Объект

Потоки — это реализация Ruby для модели параллельного программирования.

Программы, требующие нескольких потоков выполнения, являются идеальными кандидатами для класса Ruby Thread.

Например, мы можем создать новый поток, отдельный от основного потока выполнения, используя ::new.

thr = Thread.new { puts "Whats the big deal" }

Затем мы можем приостановить выполнение основного потока и дождаться завершения нового потока, используя join:

thr.join #=> "Whats the big deal"

Если мы не вызовем thr.join до завершения основного потока, то все остальные потоки, включая thr , будут убиты.

В качестве альтернативы, вы можете использовать массив для одновременной обработки нескольких потоков, как в следующем примере:

threads = []
threads << Thread.new { puts "Whats the big deal" }
threads << Thread.new { 3.times { puts "Threads are fun!" } }

После создания нескольких потоков мы ждем их всех последовательного завершения.

threads.each { |thr| thr.join }

Поток инициализация

Для создания новых потоков Ruby предоставляет ::new, ::start и ::fork. Блок должен быть предоставлен с каждым из этих методов, в противном случае будет поднято исключение ThreadError.

При наследовании от класса Thread метод initialize вашего подкласса будет проигнорирован методами ::start и ::fork. В противном случае, убедитесь, что вы вызываете super в методе initialize.

Поток завершение

Для завершения потоков Ruby предоставляет различные способы.

Метод класса ::kill предназначен для выхода из данного потока:

thr = Thread.new { ... }
Thread.kill(thr) # sends exit() to thr

В качестве альтернативы, вы можете использовать метод экземпляра exit или любой из его псевдонимов kill или terminate.

thr.exit

Поток статус

Ruby предоставляет несколько методов экземпляров для запроса состояния данного потока. Чтобы получить строку с текущим состоянием потока, используйте status

thr = Thread.new { sleep }
thr.status # => "sleep"
thr.exit
thr.status # => false

Вы также можете использовать alive? для определения, запущен ли поток или спит, и stop?, если поток мертв или спит.

Поток переменные и область видимости

Поскольку потоки создаются с блоками, те же правила применяются к другим блокам Ruby для области видимости переменных. Любые локальные переменные, созданные в этом блоке, доступны только этому потоку.

Fiber-локальные против Thread-локальных

У каждого волокна есть свой контейнер для #[] хранения. Когда вы устанавливаете новую fiber-локальную переменную, она доступна только в рамках этого Fiber.

Чтобы проиллюстрировать:

Thread.new {
  Thread.current[:foo] = "bar"
  Fiber.new {
    p Thread.current[:foo] # => nil
  }.resume
}.join

В этом примере используются [] для получения и []= для установки fiber-локальных переменных. Вы также можете использовать keys для перечисления fiber-локальных переменных для данного потока и key? для проверки существования fiber-локальной переменной.

Что касается thread-локальных переменных, они доступны во всей области видимости потока. Учитывая следующий пример:

Thread.new{
  Thread.current.thread_variable_set(:foo, 1)
  p Thread.current.thread_variable_get(:foo) # => 1
  Fiber.new{
    Thread.current.thread_variable_set(:foo, 2)
    p Thread.current.thread_variable_get(:foo) # => 2
  }.resume
  p Thread.current.thread_variable_get(:foo)   # => 2
}.join

Вы можете видеть, что thread-локальная переменная :foo переносилась в волокно и была изменена на 2 к концу потока.

В этом примере используется thread_variable_set для создания новых thread-локальных переменных и thread_variable_get для обращения к ним.

Также есть thread_variables для перечисления всех thread-локальных переменных и thread_variable? для проверки существования данной thread-локальной переменной.

Обработка исключений

Любой поток может вызвать исключение, используя метод экземпляра raise, который работает аналогично Kernel#raise.

Однако важно отметить, что исключение, возникшее в любом потоке, кроме основного, зависит от abort_on_exception. Этот параметр false по умолчанию, что означает, что любое необработанное исключение приведет к молчащему завершению потока при ожидании его с помощью join или value. Вы можете изменить это значение по умолчанию, либо установив abort_on_exception= true , либо установив $DEBUG в true.

С добавлением метода класса ::handle_interrupt вы теперь можете обрабатывать исключения асинхронно с потоками.

Планирование

Ruby предоставляет несколько способов поддержки планирования потоков в вашей программе.

Первый способ — использование метода класса ::stop для приостановки текущего выполняющегося потока и планирования выполнения другого потока.

После того, как поток приостановится, вы можете использовать метод экземпляра wakeup для обозначения вашего потока как подходящего для планирования.

Вы также можете попробовать ::pass, который пытается передать выполнение другому потоку, но зависит от ОС, будет ли текущий поток переключаться или нет. То же самое относится к priority, который позволяет подсказать планировщику потоков, какие потоки вы хотите назначить приоритет при передаче выполнения. Этот метод также зависит от ОС и может быть проигнорирован на некоторых платформах.

Методы публичного класса

DEBUG → num Показать исходный код
static VALUE
rb_thread_s_debug(void)
{
    return INT2NUM(rb_thread_debug_enabled);
}

Возвращает уровень отладки потока. Доступно только при компиляции с THREAD_DEBUG=-1.

DEBUG = num Показать исходный код
static VALUE
rb_thread_s_debug_set(VALUE self, VALUE val)
{
    rb_thread_debug_enabled = RTEST(val) ? NUM2INT(val) : 0;
    return val;
}

Устанавливает уровень отладки потока. Доступно только при компиляции с THREAD_DEBUG=-1.

abort_on_exception → true или false Показать исходный код
static VALUE
rb_thread_s_abort_exc(void)
{
    return GET_THREAD()->vm->thread_abort_on_exception ? Qtrue : Qfalse;
}

Возвращает состояние глобальной опции «прерывание при исключении».

По умолчанию false.

При установке в true, если любой поток прерывается исключением, то возбужденное исключение будет повторно возбуждено в главном потоке.

Также может быть задано глобальным флагом $DEBUG или командной опцией -d.

См. также ::abort_on_exception=.

Существует также метод уровня экземпляра для установки этого значения для конкретного потока, см. abort_on_exception.

abort_on_exception= boolean → true или false Показать исходный код
static VALUE
rb_thread_s_abort_exc_set(VALUE self, VALUE val)
{
    GET_THREAD()->vm->thread_abort_on_exception = RTEST(val);
    return val;
}

При установке в true, если любой поток прерывается исключением, то возбужденное исключение будет повторно возбуждено в главном потоке. Возвращает новое состояние.

Thread.abort_on_exception = true
t1 = Thread.new do
  puts  "In new thread"
  raise "Exception from thread"
end
sleep(1)
puts "not reached"

Это даст:

In new thread
prog.rb:4: Exception from thread (RuntimeError)
 from prog.rb:2:in `initialize'
 from prog.rb:2:in `new'
 from prog.rb:2

См. также ::abort_on_exception.

Существует также метод уровня экземпляра для установки этого значения для конкретного потока, см. abort_on_exception=.

current → поток Показать исходный код
static VALUE
thread_s_current(VALUE klass)
{
    return rb_thread_current();
}

Возвращает текущий выполняемый поток.

Thread.current   #=> #<Thread:0x401bdf4c run>
exclusive { блок } → obj Показать исходный код
# File prelude.rb, line 10
def self.exclusive
  MUTEX_FOR_THREAD_EXCLUSIVE.synchronize{
    yield
  }
end

Оборачивает блок в единый, глобальный для ВМ, Mutex#synchronize, возвращая значение блока. Поток, выполняющий код внутри раздела exclusive, будет блокировать только другие потоки, которые также используют механизм ::exclusive.

exit → поток Показать исходный код
static VALUE
rb_thread_exit(void)
{
    rb_thread_t *th = GET_THREAD();
    return rb_thread_kill(th->self);
}

Завершает текущий поток и планирует выполнение другого потока.

Если этот поток уже помечен для завершения, ::exit возвращает Поток.

Если это главный поток или последний поток, завершает процесс.

fork([args]*) {|args| блок } → поток Показать исходный код
static VALUE
thread_start(VALUE klass, VALUE args)
{
    return thread_create_core(rb_thread_alloc(klass), args, 0);
}

В основном то же, что и ::new. Однако, если класс Thread наследуется, то вызов start в этом подклассе не вызовет метод подкласса initialize.

handle_interrupt(hash) { ... } → результат блока Показать исходный код
static VALUE
rb_thread_s_handle_interrupt(VALUE self, VALUE mask_arg)
{
    VALUE mask;
    rb_thread_t *th = GET_THREAD();
    VALUE r = Qnil;
    int state;

    if (!rb_block_given_p()) {
        rb_raise(rb_eArgError, "block is needed.");
    }

    mask = rb_convert_type(mask_arg, T_HASH, "Hash", "to_hash");
    rb_hash_foreach(mask, handle_interrupt_arg_check_i, 0);
    rb_ary_push(th->pending_interrupt_mask_stack, mask);
    if (!rb_threadptr_pending_interrupt_empty_p(th)) {
        th->pending_interrupt_queue_checked = 0;
        RUBY_VM_SET_INTERRUPT(th);
    }

    TH_PUSH_TAG(th);
    if ((state = EXEC_TAG()) == 0) {
        r = rb_yield(Qnil);
    }
    TH_POP_TAG();

    rb_ary_pop(th->pending_interrupt_mask_stack);
    if (!rb_threadptr_pending_interrupt_empty_p(th)) {
        th->pending_interrupt_queue_checked = 0;
        RUBY_VM_SET_INTERRUPT(th);
    }

    RUBY_VM_CHECK_INTS(th);

    if (state) {
        JUMP_TAG(state);
    }

    return r;
}

Изменяет время асинхронного прерывания.

Прерывание означает асинхронное событие и соответствующую процедуру с помощью #raise, #kill, обработки сигнала (ещё не поддерживается) и завершения главного потока (если главный поток завершается, то все остальные потоки будут убиты).

Указанный hash содержит пары, например, ExceptionClass => :TimingSymbol. Где ExceptionClass — это прерывание, обрабатываемое заданным блоком. TimingSymbol может быть одним из следующих символов:

:immediate

Вызвать прерывания немедленно.

:on_blocking

Вызвать прерывания во время BlockingOperation.

:never

Никогда не вызывать все прерывания.

BlockingOperation означает, что операция заблокирует вызывающий поток, например, чтение и запись. В реализации CRuby BlockingOperation — это любая операция, выполняемая без GVL.

Маскированные асинхронные прерывания откладываются до их включения. Этот метод аналогичен sigprocmask(3).

ПРИМЕЧАНИЕ

Асинхронные прерывания трудно использовать.

Если вам нужно взаимодействовать между потоками, рассмотрите другой способ, например, Queue.

Или используйте их с глубоким пониманием этого метода.

Использование

В этом примере мы можем защититься от исключений #raise.

Используя :never TimingSymbol, исключение RuntimeError всегда будет игнорироваться в первом блоке главного потока. Во втором блоке ::handle_interrupt мы можем намеренно обработать исключения RuntimeError.

th = Thread.new do
  Thread.handle_interrupt(RuntimeError => :never) {
    begin
      # You can write resource allocation code safely.
      Thread.handle_interrupt(RuntimeError => :immediate) {
        # ...
      }
    ensure
      # You can write resource deallocation code safely.
    end
  }
end
Thread.pass
# ...
th.raise "stop"

В то время как мы игнорируем исключение RuntimeError, безопасно написать код выделения ресурсов. Затем в блоке ensure мы можем безопасно освободить ресурсы.

Защита от Timeout::Error

В следующем примере мы защитимся от исключения Timeout::Error. Это поможет предотвратить утечку ресурсов при возникновении исключений Timeout::Error во время обычного блока ensure. Для этого примера мы воспользуемся стандартной библиотекой Timeout из lib/timeout.rb

require 'timeout'
Thread.handle_interrupt(Timeout::Error => :never) {
  timeout(10){
    # Timeout::Error doesn't occur here
    Thread.handle_interrupt(Timeout::Error => :on_blocking) {
      # possible to be killed by Timeout::Error
      # while blocking operation
    }
    # Timeout::Error doesn't occur here
  }
}

В первой части блока timeout мы можем полагаться на то, что Timeout::Error будет проигнорирован. Затем в блоке Timeout::Error => :on_blocking любая операция, которая заблокирует вызывающий поток, подвержена риску возникновения исключения Timeout::Error.

Настройки управления стеком

Можно укладывать несколько уровней блоков ::handle_interrupt, чтобы контролировать более одного ExceptionClass и TimingSymbol одновременно.

Thread.handle_interrupt(FooError => :never) {
  Thread.handle_interrupt(BarError => :never) {
     # FooError and BarError are prohibited.
  }
}

Наследование с ExceptionClass

Все исключения, унаследованные от параметра ExceptionClass, будут рассматриваться.

Thread.handle_interrupt(Exception => :never) {
  # all exceptions inherited from Exception are prohibited.
}
kill(поток) → поток Показать исходный код
static VALUE
rb_thread_s_kill(VALUE obj, VALUE th)
{
    return rb_thread_kill(th);
}

Принудительно завершает указанный thread, см. также ::exit.

count = 0
a = Thread.new { loop { count += 1 } }
sleep(0.1)       #=> 0
Thread.kill(a)   #=> #<Thread:0x401b3d30 dead>
count            #=> 93947
a.alive?         #=> false
list → массив Показать исходный код
VALUE
rb_thread_list(void)
{
    VALUE ary = rb_ary_new();
    rb_vm_t *vm = GET_THREAD()->vm;
    rb_thread_t *th = 0;

    list_for_each(&vm->living_threads, th, vmlt_node) {
        switch (th->status) {
          case THREAD_RUNNABLE:
          case THREAD_STOPPED:
          case THREAD_STOPPED_FOREVER:
            rb_ary_push(ary, th->self);
          default:
            break;
        }
    }
    return ary;
}

Возвращает массив объектов Thread для всех потоков, которые находятся в состоянии готовности или остановлены.

Thread.new { sleep(200) }
Thread.new { 1000000.times {|i| i*i } }
Thread.new { Thread.stop }
Thread.list.each {|t| p t}

Это даст:

#<Thread:0x401b3e84 sleep>
#<Thread:0x401b3f38 run>
#<Thread:0x401b3fb0 sleep>
#<Thread:0x401bdf4c run>
main → поток Показать исходный код
static VALUE
rb_thread_s_main(VALUE klass)
{
    return rb_thread_main();
}

Возвращает главный поток.

new { ... } → поток Показать исходный код
new(*args, &proc) → поток
new(*args) { |args| ... } → поток
static VALUE
thread_s_new(int argc, VALUE *argv, VALUE klass)
{
    rb_thread_t *th;
    VALUE thread = rb_thread_alloc(klass);

    if (GET_VM()->main_thread->status == THREAD_KILLED)
        rb_raise(rb_eThreadError, "can't alloc thread");

    rb_obj_call_init(thread, argc, argv);
    GetThreadPtr(thread, th);
    if (!th->first_args) {
        rb_raise(rb_eThreadError, "uninitialized thread - check `%s#initialize'",
                 rb_class2name(klass));
    }
    return thread;
}

Создаёт новый поток, выполняющий заданный блок.

Любые args, переданные в ::new, будут переданы в блок:

arr = []
a, b, c = 1, 2, 3
Thread.new(a,b,c) { |d,e,f| arr << d << e << f }.join
arr #=> [1, 2, 3]

Исключение ThreadError возбуждается, если ::new вызывается без блока.

Если вы собираетесь создавать подкласс Thread, убедитесь, что вызываете super в методе initialize , в противном случае будет возбуждено исключение ThreadError.

pass → nil Показать исходный код
static VALUE
thread_s_pass(VALUE klass)
{
    rb_thread_schedule();
    return Qnil;
}

Даёт подсказку планировщику потоков для передачи выполнения другому потоку. Текущий поток может или не может переключиться, это зависит от ОС и процессора.

pending_interrupt?(error = nil) → true/false Показать исходный код
static VALUE
rb_thread_s_pending_interrupt_p(int argc, VALUE *argv, VALUE self)
{
    return rb_thread_pending_interrupt_p(argc, argv, GET_THREAD()->self);
}

Возвращает значение, указывающее, пуста ли асинхронная очередь.

Так как ::handle_interrupt может использоваться для отложенного выполнения асинхронных событий, этот метод можно использовать для определения наличия отложенных событий.

Если этот метод возвращает true, то вы можете завершить :never блоки.

Например, следующий метод обрабатывает отложенные асинхронные события немедленно.

def Thread.kick_interrupt_immediately
  Thread.handle_interrupt(Object => :immediate) {
    Thread.pass
  }
end

Если error задано, то проверяются только отложенные события типа error.

Использование

th = Thread.new{
  Thread.handle_interrupt(RuntimeError => :on_blocking){
    while true
      ...
      # reach safe point to invoke interrupt
      if Thread.pending_interrupt?
        Thread.handle_interrupt(Object => :immediate){}
      end
      ...
    end
  }
}
...
th.raise # stop thread

Этот пример также можно переписать следующим образом, что следует использовать для предотвращения асинхронных прерываний.

flag = true
th = Thread.new{
  Thread.handle_interrupt(RuntimeError => :on_blocking){
    while true
      ...
      # reach safe point to invoke interrupt
      break if flag == false
      ...
    end
  }
}
...
flag = false # stop thread
start([args]*) {|args| block } → thread Показать исходный код
static VALUE
thread_start(VALUE klass, VALUE args)
{
    return thread_create_core(rb_thread_alloc(klass), args, 0);
}

В основном то же самое, что и ::new. Однако, если класс Thread является подклассом, то вызов start в этом подклассе не вызовет метод initialize подкласса.

stop → nil Показать исходный код
VALUE
rb_thread_stop(void)
{
    if (rb_thread_alone()) {
        rb_raise(rb_eThreadError,
                 "stopping only thread\n\tnote: use sleep to stop forever");
    }
    rb_thread_sleep_deadly();
    return Qnil;
}

Останавливает выполнение текущей нити, переводя её в состояние «ожидания», и планирует выполнение другой нити.

a = Thread.new { print "a"; Thread.stop; print "c" }
sleep 0.1 while a.status!='sleep'
print "b"
a.run
a.join
#=> "abc"

Публичные методы экземпляра

thr[sym] → obj or nil Show source
static VALUE
rb_thread_aref(VALUE thread, VALUE key)
{
    ID id = rb_check_id(&key);
    if (!id) return Qnil;
    return rb_thread_local_aref(thread, id);
}

Ссылка на атрибут — возвращает значение локальной для волокна переменной (корневое волокно текущего потока, если явно не находится внутри Fiber), используя либо символ, либо строковое имя. Если указанная переменная не существует, возвращает nil.

[
  Thread.new { Thread.current["name"] = "A" },
  Thread.new { Thread.current[:name]  = "B" },
  Thread.new { Thread.current["name"] = "C" }
].each do |th|
  th.join
  puts "#{th.inspect}: #{th[:name]}"
end

Это даст:

#<Thread:0x00000002a54220 dead>: A
#<Thread:0x00000002a541a8 dead>: B
#<Thread:0x00000002a54130 dead>: C

#[] и #[]= не являются локальными для потока, а локальными для волокна. Этого недоразумения не существовало в Ruby 1.8, поскольку волокна доступны только начиная с Ruby 1.9. В Ruby 1.9 выбрано, чтобы методы вели себя локально для волокна, чтобы сохранить следующую идиому для динамической области видимости.

def meth(newvalue)
  begin
    oldvalue = Thread.current[:name]
    Thread.current[:name] = newvalue
    yield
  ensure
    Thread.current[:name] = oldvalue
  end
end

Идиома может не работать как динамическая область видимости, если методы являются локальными для потока, и данный блок переключает волокна.

f = Fiber.new {
  meth(1) {
    Fiber.yield
  }
}
meth(2) {
  f.resume
}
f.resume
p Thread.current[:name]
#=> nil if fiber-local
#=> 2 if thread-local (The value 2 is leaked to outside of meth method.)

Для локальных для потока переменных см. thread_variable_get и thread_variable_set.

thr[sym] = obj → obj Show source
static VALUE
rb_thread_aset(VALUE self, VALUE id, VALUE val)
{
    return rb_thread_local_aset(self, rb_to_id(id), val);
}

Присваивание атрибута — устанавливает или создает значение локальной для волокна переменной, используя либо символ, либо строку.

См. также #[].

Для локальных для потока переменных см. thread_variable_set и thread_variable_get.

abort_on_exception → true or false Show source
static VALUE
rb_thread_abort_exc(VALUE thread)
{
    rb_thread_t *th;
    GetThreadPtr(thread, th);
    return th->abort_on_exception ? Qtrue : Qfalse;
}

Возвращает состояние локального для потока условия «прервать при исключении» для этого thr.

Значение по умолчанию — false.

См. также abort_on_exception=.

Существует также метод уровня класса для установки этого для всех потоков, см. ::abort_on_exception.

abort_on_exception= boolean → true or false Show source
static VALUE
rb_thread_abort_exc_set(VALUE thread, VALUE val)
{
    rb_thread_t *th;

    GetThreadPtr(thread, th);
    th->abort_on_exception = RTEST(val);
    return val;
}

Если установлено в true, если этот thr прерывается исключением, сгенерированное исключение будет снова вызвано в основном потоке.

См. также abort_on_exception.

Существует также метод уровня класса для установки этого для всех потоков, см. ::abort_on_exception=.

add_trace_func(proc) → proc Show source
static VALUE
thread_add_trace_func_m(VALUE obj, VALUE trace)
{
    rb_thread_t *th;

    GetThreadPtr(obj, th);
    thread_add_trace_func(th, trace);
    return trace;
}

Добавляет proc в качестве обработчика трассировки.

См. #set_trace_func и Kernel#set_trace_func.

alive? → true or false Show source
static VALUE
rb_thread_alive_p(VALUE thread)
{
    rb_thread_t *th;
    GetThreadPtr(thread, th);

    if (rb_threadptr_dead(th))
        return Qfalse;
    return Qtrue;
}

Возвращает true, если thr выполняется или находится в состоянии ожидания.

thr = Thread.new { }
thr.join                #=> #<Thread:0x401b3fb0 dead>
Thread.current.alive?   #=> true
thr.alive?              #=> false

См. также stop? и status.

backtrace → array Show source
static VALUE
rb_thread_backtrace_m(int argc, VALUE *argv, VALUE thval)
{
    return rb_vm_thread_backtrace(argc, argv, thval);
}

Возвращает текущий стек вызовов целевого потока.

backtrace_locations(*args) → array or nil Show source
static VALUE
rb_thread_backtrace_locations_m(int argc, VALUE *argv, VALUE thval)
{
    return rb_vm_thread_backtrace_locations(argc, argv, thval);
}

Возвращает стек выполнения для целевого потока — массив, содержащий объекты местоположения стека вызовов.

См. Thread::Backtrace::Location для получения дополнительной информации.

Этот метод ведет себя аналогично Kernel#caller_locations, за исключением того, что он применяется к конкретному потоку.

exit → thr or nil Show source
kill → thr or nil
terminate → thr or nil
VALUE
rb_thread_kill(VALUE thread)
{
    rb_thread_t *th;

    GetThreadPtr(thread, th);

    if (th->to_kill || th->status == THREAD_KILLED) {
        return thread;
    }
    if (th == th->vm->main_thread) {
        rb_exit(EXIT_SUCCESS);
    }

    thread_debug("rb_thread_kill: %p (%"PRI_THREAD_ID")\n", (void *)th, thread_id_str(th));

    if (th == GET_THREAD()) {
        /* kill myself immediately */
        rb_threadptr_to_kill(th);
    }
    else {
        threadptr_check_pending_interrupt_queue(th);
        rb_threadptr_pending_interrupt_enque(th, eKillSignal);
        rb_threadptr_interrupt(th);
    }
    return thread;
}

Завершает thr и планирует выполнение другого потока.

Если этот поток уже помечен как подлежащий уничтожению, exit возвращает Thread.

Если это главный поток или последний поток, завершает процесс.

group → thgrp or nil Show source
VALUE
rb_thread_group(VALUE thread)
{
    rb_thread_t *th;
    VALUE group;
    GetThreadPtr(thread, th);
    group = th->thgroup;

    if (!group) {
        group = Qnil;
    }
    return group;
}

Возвращает ThreadGroup, содержащий данный поток, или возвращает nil, если thr не является членом ни одной группы.

Thread.main.group   #=> #<ThreadGroup:0x4029d914>
inspect → string Show source
static VALUE
rb_thread_inspect(VALUE thread)
{
    return rb_thread_inspect_msg(thread, 1, 1, 1);
}

Выводит имя, идентификатор и состояние thr в строку.

join → thr Show source
join(limit) → thr
static VALUE
thread_join_m(int argc, VALUE *argv, VALUE self)
{
    rb_thread_t *target_th;
    double delay = DELAY_INFTY;
    VALUE limit;

    GetThreadPtr(self, target_th);

    rb_scan_args(argc, argv, "01", &limit);
    if (!NIL_P(limit)) {
        delay = rb_num2dbl(limit);
    }

    return thread_join(target_th, delay);
}

Вызывающий поток приостановит выполнение и запустит этот thr.

Не возвращается до тех пор, пока thr не завершится или пока не пройдут указанные limit секунд.

Если истекает лимит времени, будет возвращено nil, в противном случае возвращается thr.

Все потоки, которые не были объединены, будут уничтожены при завершении основной программы.

Если thr ранее сгенерировал исключение, и флаги ::abort_on_exception или $DEBUG не установлены (поэтому исключение еще не обработано), оно будет обработано в это время.

a = Thread.new { print "a"; sleep(10); print "b"; print "c" }
x = Thread.new { print "x"; Thread.pass; print "y"; print "z" }
x.join # Let thread x finish, thread a will be killed on exit.
#=> "axyz"

Следующий пример иллюстрирует параметр limit.

y = Thread.new { 4.times { sleep 0.1; puts 'tick... ' }}
puts "Waiting" until y.join(0.15)

Это даст:

tick...
Waiting
tick...
Waiting
tick...
tick...
key?(sym) → true or false Show source
static VALUE
rb_thread_key_p(VALUE self, VALUE key)
{
    rb_thread_t *th;
    ID id = rb_check_id(&key);

    GetThreadPtr(self, th);

    if (!id || !th->local_storage) {
        return Qfalse;
    }
    if (st_lookup(th->local_storage, id, 0)) {
        return Qtrue;
    }
    return Qfalse;
}

Возвращает true, если данная строка (или символ) существует как локальная для волокна переменная.

me = Thread.current
me[:oliver] = "a"
me.key?(:oliver)    #=> true
me.key?(:stanley)   #=> false
keys → array Show source
static VALUE
rb_thread_keys(VALUE self)
{
    rb_thread_t *th;
    VALUE ary = rb_ary_new();
    GetThreadPtr(self, th);

    if (th->local_storage) {
        st_foreach(th->local_storage, thread_keys_i, ary);
    }
    return ary;
}

Возвращает массив имен локальных для волокна переменных (в виде символов).

thr = Thread.new do
  Thread.current[:cat] = 'meow'
  Thread.current["dog"] = 'woof'
end
thr.join   #=> #<Thread:0x401b3f10 dead>
thr.keys   #=> [:dog, :cat]
exit → thr or nil Show source
kill → thr or nil
terminate → thr or nil
VALUE
rb_thread_kill(VALUE thread)
{
    rb_thread_t *th;

    GetThreadPtr(thread, th);

    if (th->to_kill || th->status == THREAD_KILLED) {
        return thread;
    }
    if (th == th->vm->main_thread) {
        rb_exit(EXIT_SUCCESS);
    }

    thread_debug("rb_thread_kill: %p (%"PRI_THREAD_ID")\n", (void *)th, thread_id_str(th));

    if (th == GET_THREAD()) {
        /* kill myself immediately */
        rb_threadptr_to_kill(th);
    }
    else {
        threadptr_check_pending_interrupt_queue(th);
        rb_threadptr_pending_interrupt_enque(th, eKillSignal);
        rb_threadptr_interrupt(th);
    }
    return thread;
}

Завершает thr и планирует выполнение другого потока.

Если этот поток уже помечен как подлежащий уничтожению, exit возвращает Thread.

Если это главный поток или последний поток, завершает процесс.

pending_interrupt?(error = nil) → true/false Show source
static VALUE
rb_thread_pending_interrupt_p(int argc, VALUE *argv, VALUE target_thread)
{
    rb_thread_t *target_th;

    GetThreadPtr(target_thread, target_th);

    if (!target_th->pending_interrupt_queue) {
        return Qfalse;
    }
    if (rb_threadptr_pending_interrupt_empty_p(target_th)) {
        return Qfalse;
    }
    else {
        if (argc == 1) {
            VALUE err;
            rb_scan_args(argc, argv, "01", &err);
            if (!rb_obj_is_kind_of(err, rb_cModule)) {
                rb_raise(rb_eTypeError, "class or module required for rescue clause");
            }
            if (rb_threadptr_pending_interrupt_include_p(target_th, err)) {
                return Qtrue;
            }
            else {
                return Qfalse;
            }
        }
        return Qtrue;
    }
}

Возвращает, пуста ли асинхронная очередь для целевого потока.

Если задан error, то проверяются только отложенные события типа error.

См. ::pending_interrupt? для получения дополнительной информации.

priority → целое Показать исходный код
static VALUE
rb_thread_priority(VALUE thread)
{
    rb_thread_t *th;
    GetThreadPtr(thread, th);
    return INT2NUM(th->priority);
}

Возвращает приоритет потока thr. По умолчанию наследуется от текущего потока, создающего новый поток, или равен нулю для начального главного потока; поток с более высоким приоритетом будет выполняться чаще, чем поток с более низким приоритетом (но потоки с более низким приоритетом также могут выполняться).

Это просто подсказка для планировщика потоков Ruby. Она может быть проигнорирована на некоторых платформах.

Thread.current.priority   #=> 0
priority= целое → thr Показать исходный код
static VALUE
rb_thread_priority_set(VALUE thread, VALUE prio)
{
    rb_thread_t *th;
    int priority;
    GetThreadPtr(thread, th);


#if USE_NATIVE_THREAD_PRIORITY
    th->priority = NUM2INT(prio);
    native_thread_apply_priority(th);
#else
    priority = NUM2INT(prio);
    if (priority > RUBY_THREAD_PRIORITY_MAX) {
        priority = RUBY_THREAD_PRIORITY_MAX;
    }
    else if (priority < RUBY_THREAD_PRIORITY_MIN) {
        priority = RUBY_THREAD_PRIORITY_MIN;
    }
    th->priority = priority;
#endif
    return INT2NUM(th->priority);
}

Устанавливает приоритет потока thr на целое. Потоки с более высоким приоритетом будут выполняться чаще, чем потоки с более низким приоритетом (но потоки с более низким приоритетом также могут выполняться).

Это просто подсказка для планировщика потоков Ruby. Она может быть проигнорирована на некоторых платформах.

count1 = count2 = 0
a = Thread.new do
      loop { count1 += 1 }
    end
a.priority = -1

b = Thread.new do
      loop { count2 += 1 }
    end
b.priority = -2
sleep 1   #=> 1
count1    #=> 622504
count2    #=> 5832
raise Показать исходный код
raise(строка)
raise(исключение [, строка [, массив]])
static VALUE
thread_raise_m(int argc, VALUE *argv, VALUE self)
{
    rb_thread_t *target_th;
    rb_thread_t *th = GET_THREAD();
    GetThreadPtr(self, target_th);
    threadptr_check_pending_interrupt_queue(target_th);
    rb_threadptr_raise(target_th, argc, argv);

    /* To perform Thread.current.raise as Kernel.raise */
    if (th == target_th) {
        RUBY_VM_CHECK_INTS(th);
    }
    return Qnil;
}

Выбрасывает исключение из данного потока. Вызывающий поток не обязан быть thr. Подробнее см. Kernel#raise.

Thread.abort_on_exception = true
a = Thread.new { sleep(200) }
a.raise("Gotcha")

Это приведет к:

prog.rb:3: Gotcha (RuntimeError)
 from prog.rb:2:in `initialize'
 from prog.rb:2:in `new'
 from prog.rb:2
run → thr Показать исходный код
VALUE
rb_thread_run(VALUE thread)
{
    rb_thread_wakeup(thread);
    rb_thread_schedule();
    return thread;
}

Разбуждает thr, делая его подходящим для планирования.

a = Thread.new { puts "a"; Thread.stop; puts "c" }
sleep 0.1 while a.status!='sleep'
puts "Got here"
a.run
a.join

Это приведет к:

a
Got here
c

См. также метод экземпляра wakeup.

safe_level → целое Показать исходный код
static VALUE
rb_thread_safe_level(VALUE thread)
{
    rb_thread_t *th;
    GetThreadPtr(thread, th);

    return INT2NUM(th->safe_level);
}

Возвращает уровень безопасности, действующий для потока thr. Установка локальных уровней безопасности потоков может быть полезна при реализации песочниц, в которых выполняется небезопасный код.

thr = Thread.new { $SAFE = 3; sleep }
Thread.current.safe_level   #=> 0
thr.safe_level              #=> 3
set_trace_func(proc) → proc Показать исходный код
set_trace_func(nil) → nil
static VALUE
thread_set_trace_func_m(VALUE obj, VALUE trace)
{
    rb_thread_t *th;

    GetThreadPtr(obj, th);
    rb_threadptr_remove_event_hook(th, call_trace_func, Qundef);

    if (NIL_P(trace)) {
        return Qnil;
    }

    thread_add_trace_func(th, trace);
    return trace;
}

Устанавливает proc для потока thr как обработчик отслеживания или отключает отслеживание, если параметр равен nil.

См. Kernel#set_trace_func.

status → строка, false или nil Показать исходный код
static VALUE
rb_thread_status(VALUE thread)
{
    rb_thread_t *th;
    GetThreadPtr(thread, th);

    if (rb_threadptr_dead(th)) {
        if (!NIL_P(th->errinfo) && !FIXNUM_P(th->errinfo)
            /* TODO */ ) {
            return Qnil;
        }
        return Qfalse;
    }
    return rb_str_new2(thread_status_name(th));
}

Возвращает состояние thr.

"sleep"

Возвращается, если этот поток спит или ожидает ввода-вывода

"run"

Когда этот поток выполняется

"aborting"

Если этот поток прерывается

false

Когда этот поток завершается нормально

nil

Если завершился с исключением.

a = Thread.new { raise("die now") }
b = Thread.new { Thread.stop }
c = Thread.new { Thread.exit }
d = Thread.new { sleep }
d.kill                  #=> #<Thread:0x401b3678 aborting>
a.status                #=> nil
b.status                #=> "sleep"
c.status                #=> false
d.status                #=> "aborting"
Thread.current.status   #=> "run"

См. также методы экземпляра alive? и stop?

stop? → true или false Показать исходный код
static VALUE
rb_thread_stop_p(VALUE thread)
{
    rb_thread_t *th;
    GetThreadPtr(thread, th);

    if (rb_threadptr_dead(th))
        return Qtrue;
    if (th->status == THREAD_STOPPED || th->status == THREAD_STOPPED_FOREVER)
        return Qtrue;
    return Qfalse;
}

Возвращает true если thr мертв или спит.

a = Thread.new { Thread.stop }
b = Thread.current
a.stop?   #=> true
b.stop?   #=> false

См. также alive? и status.

terminate → thr или nil Показать исходный код
VALUE
rb_thread_kill(VALUE thread)
{
    rb_thread_t *th;

    GetThreadPtr(thread, th);

    if (th->to_kill || th->status == THREAD_KILLED) {
        return thread;
    }
    if (th == th->vm->main_thread) {
        rb_exit(EXIT_SUCCESS);
    }

    thread_debug("rb_thread_kill: %p (%"PRI_THREAD_ID")\n", (void *)th, thread_id_str(th));

    if (th == GET_THREAD()) {
        /* kill myself immediately */
        rb_threadptr_to_kill(th);
    }
    else {
        threadptr_check_pending_interrupt_queue(th);
        rb_threadptr_pending_interrupt_enque(th, eKillSignal);
        rb_threadptr_interrupt(th);
    }
    return thread;
}

Прерывает thr и планирует запуск другого потока.

Если этот поток уже помечен для убийства, exit возвращает Thread.

Если это главный поток или последний поток, завершает процесс.

thread_variable?(ключ) → true или false Показать исходный код
static VALUE
rb_thread_variable_p(VALUE thread, VALUE key)
{
    VALUE locals;
    ID id = rb_check_id(&key);

    if (!id) return Qfalse;

    locals = rb_ivar_get(thread, id_locals);

    if (!RHASH(locals)->ntbl)
        return Qfalse;

    if (st_lookup(RHASH(locals)->ntbl, ID2SYM(id), 0)) {
        return Qtrue;
    }

    return Qfalse;
}

Возвращает true если заданная строка (или символ) существует в качестве локальной переменной потока.

me = Thread.current
me.thread_variable_set(:oliver, "a")
me.thread_variable?(:oliver)    #=> true
me.thread_variable?(:stanley)   #=> false

Обратите внимание, что это не локальные переменные волокна. Подробнее см. #[] и #thread_variable_get.

thread_variable_get(ключ) → obj или nil Показать исходный код
static VALUE
rb_thread_variable_get(VALUE thread, VALUE key)
{
    VALUE locals;

    locals = rb_ivar_get(thread, id_locals);
    return rb_hash_aref(locals, rb_to_symbol(key));
}

Возвращает значение локальной переменной потока, которое было установлено. Обратите внимание, что они отличаются от локальных значений волокна. Для локальных значений волокна см. #[] и #[]=.

Thread локальные значения передаются вместе с потоками и не учитывают волокна. Например:

Thread.new {
  Thread.current.thread_variable_set("foo", "bar") # set a thread local
  Thread.current["foo"] = "bar"                    # set a fiber local

  Fiber.new {
    Fiber.yield [
      Thread.current.thread_variable_get("foo"), # get the thread local
      Thread.current["foo"],                     # get the fiber local
    ]
  }.resume
}.join.value # => ['bar', nil]

Значение «bar» возвращается для локальной переменной потока, где nil возвращается для локальной переменной волокна. Волокно выполняется в том же потоке, поэтому доступны значения локальных переменных потока.

thread_variable_set(ключ, значение) Показать исходный код
static VALUE
rb_thread_variable_set(VALUE thread, VALUE id, VALUE val)
{
    VALUE locals;

    if (OBJ_FROZEN(thread)) {
        rb_error_frozen("thread locals");
    }

    locals = rb_ivar_get(thread, id_locals);
    return rb_hash_aset(locals, rb_to_symbol(id), val);
}

Устанавливает локальную переменную потока со значением key в value. Обратите внимание, что они локальны для потоков, а не для волокон. Подробнее см. #thread_variable_get и #[].

thread_variables → массив Показать исходный код
static VALUE
rb_thread_variables(VALUE thread)
{
    VALUE locals;
    VALUE ary;

    locals = rb_ivar_get(thread, id_locals);
    ary = rb_ary_new();
    rb_hash_foreach(locals, keys_i, ary);

    return ary;
}

Возвращает массив имен локальных переменных потока (в виде символов).

thr = Thread.new do
  Thread.current.thread_variable_set(:cat, 'meow')
  Thread.current.thread_variable_set("dog", 'woof')
end
thr.join               #=> #<Thread:0x401b3f10 dead>
thr.thread_variables   #=> [:dog, :cat]

Обратите внимание, что это не локальные переменные волокна. Подробнее см. #[] и #thread_variable_get.

value → obj Показать исходный код
static VALUE
thread_value(VALUE self)
{
    rb_thread_t *th;
    GetThreadPtr(self, th);
    thread_join(th, DELAY_INFTY);
    return th->value;
}

Ожидает завершения thr, используя join, и возвращает его значение или вызывает исключение, которое завершило поток.

a = Thread.new { 2 + 2 }
a.value   #=> 4

b = Thread.new { raise 'something went wrong' }
b.value   #=> RuntimeError: something went wrong
wakeup → thr Показать исходный код
VALUE
rb_thread_wakeup(VALUE thread)
{
    if (!RTEST(rb_thread_wakeup_alive(thread))) {
        rb_raise(rb_eThreadError, "killed thread");
    }
    return thread;
}

Помечает данный поток как подходящий для планирования, однако он все еще может быть заблокирован ввода-выводом.

Примечание: Это не вызывает планировщик, см. run для получения дополнительной информации.

c = Thread.new { Thread.stop; puts "hey!" }
sleep 0.1 while c.status!='sleep'
c.wakeup
c.join
#=> "hey!"

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

Spec-Zone.ru

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