Spec-Zone.ru › Ruby 4.0

класс Thread

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

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

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

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

thr = Thread.new { puts "What's the big deal" }

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

thr.join #=> "What's the big deal"

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

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

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

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

threads.each { |thr| thr.join }

Чтобы получить последнее значение потока, используйте value

thr = Thread.new { sleep 1; "Useful value" }
thr.value #=> "Useful value"

Thread: инициализация

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

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

Thread: завершение

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

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

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

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

thr.exit

Thread: состояние

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

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

Также можно использовать alive?, чтобы узнать, выполняется поток или находится в режиме ожидания, и stop?, чтобы проверить, завершён поток или находится в режиме ожидания.

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

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

Локальные переменные Fiber и потока

У каждого волокна есть собственное хранилище для Thread#[]. Новая локальная переменная волокна доступна только внутри этого Fiber. Рассмотрим пример:

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

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

Локальные переменные потока доступны в пределах всей области видимости потока. Рассмотрим следующий пример:

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

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

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

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

Exception: обработка

Если внутри потока возникает необработанное исключение, поток завершается. По умолчанию это исключение не передаётся другим потокам. Исключение сохраняется, и когда другой поток вызывает value или join, оно возбуждается повторно в этом потоке.

t = Thread.new{ raise 'something went wrong' }
t.value #=> RuntimeError: something went wrong

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

Если установить Thread.abort_on_exception = true, Thread#abort_on_exception = true или $DEBUG = true, то любое последующее необработанное исключение, возникшее в потоке, будет автоматически возбуждено повторно в основном потоке.

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

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

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

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

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

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

Общедоступные методы класса

abort_on_exception → true or false Показать исходный код
static VALUE
rb_thread_s_abort_exc(VALUE _)
{
    return RBOOL(GET_THREAD()->vm->thread_abort_on_exception);
}

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

По умолчанию установлено значение false.

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

Также это условие можно задать с помощью глобального флага $DEBUG или параметра командной строки -d.

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

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

abort_on_exception= boolean → true or 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 → thread Показать исходный код
static VALUE
thread_s_current(VALUE klass)
{
    return rb_thread_current();
}

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

Thread.current   #=> #<Thread:0x401bdf4c run>
each_caller_location(...) { |loc| ... } → nil Показать исходный код
static VALUE
each_caller_location(int argc, VALUE *argv, VALUE _)
{
    rb_execution_context_t *ec = GET_EC();
    long n, lev = ec_backtrace_range(ec, argc, argv, 1, 1, &n);
    if (lev >= 0 && n != 0) {
        rb_ec_partial_backtrace_object(ec, lev, n, NULL, FALSE, TRUE);
    }
    return Qnil;
}

Передаёт в блок каждый кадр текущего стека вызовов в виде объекта расположения в трассировке стека.

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

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

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

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

start([args]*) {|args| block } → thread Показать исходный код
fork([args]*) {|args| block } → thread
static VALUE
thread_start(VALUE klass, VALUE args)
{
    struct thread_create_params params = {
        .type = thread_invoke_type_proc,
        .args = args,
        .proc = rb_block_proc(),
    };
    return thread_create_core(rb_thread_alloc(klass), &params);
}

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

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

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

    mask_arg = rb_to_hash_type(mask_arg);

    if (OBJ_FROZEN(mask_arg) && rb_hash_compare_by_id_p(mask_arg)) {
        mask = Qnil;
    }

    rb_hash_foreach(mask_arg, handle_interrupt_arg_check_i, (VALUE)&mask);

    if (UNDEF_P(mask)) {
        return rb_yield(Qnil);
    }

    if (!RTEST(mask)) {
        mask = mask_arg;
    }
    else if (RB_TYPE_P(mask, T_HASH)) {
        OBJ_FREEZE(mask);
    }

    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->ec);
    }

    EC_PUSH_TAG(th->ec);
    if ((state = EC_EXEC_TAG()) == TAG_NONE) {
        r = rb_yield(Qnil);
    }
    EC_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->ec);
    }

    RUBY_VM_CHECK_INTS(th->ec);

    if (state) {
        EC_JUMP_TAG(th->ec, state);
    }

    return r;
}

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

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

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

:immediate

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

:on_blocking

Обрабатывать прерывания во время блокирующей операции.

:never

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

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

Замаскированные асинхронные прерывания откладываются до тех пор, пока их обработка не будет разрешена. Этот метод аналогичен sigprocmask(3).

Примечание

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

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

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

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

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

При использовании TimingSymbol :never исключение 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 позволяет безопасно освободить ресурсы.

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

Можно вкладывать несколько уровней блоков ::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.
}

Для обработки всех прерываний используйте Object, а не Exception в качестве ExceptionClass, поскольку прерывания kill/terminate не обрабатываются классом Exception.

ignore_deadlock → true or false Показать исходный код
static VALUE
rb_thread_s_ignore_deadlock(VALUE _)
{
    return RBOOL(GET_THREAD()->vm->thread_ignore_deadlock);
}

Возвращает состояние глобального условия «игнорировать взаимоблокировку». По умолчанию установлено значение false, поэтому условия взаимоблокировки не игнорируются.

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

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

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

Thread.ignore_deadlock = true
queue = Thread::Queue.new

trap(:SIGUSR1){queue.push "Received signal"}

# raises fatal error unless ignoring deadlock
puts queue.pop

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

kill(thread) → thread Показать исходный код
static VALUE
rb_thread_s_kill(VALUE obj, VALUE th)
{
    return rb_thread_kill(th);
}

Завершает заданный thread. См. также 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 → array Показать исходный код
static VALUE
thread_list(VALUE _)
{
    return rb_thread_list();
}

Возвращает массив объектов 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 → thread Показать исходный код
static VALUE
rb_thread_s_main(VALUE klass)
{
    return rb_thread_main();
}

Возвращает основной поток.

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

    if (GET_RACTOR()->threads.main->status == THREAD_KILLED) {
        rb_raise(rb_eThreadError, "can't alloc thread");
    }

    rb_obj_call_init_kw(thread, argc, argv, RB_PASS_CALLED_KEYWORDS);
    th = rb_thread_ptr(thread);
    if (!threadptr_initialized(th)) {
        rb_raise(rb_eThreadError, "uninitialized thread - check '%"PRIsVALUE"#initialize'",
                 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]

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

Если вы собираетесь создать подкласс 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);
}

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

Поскольку Thread::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
report_on_exception → true or false Показать исходный код
static VALUE
rb_thread_s_report_exc(VALUE _)
{
    return RBOOL(GET_THREAD()->vm->thread_report_on_exception);
}

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

Начиная с Ruby 2.5 по умолчанию установлено значение true.

Все потоки, созданные при включённом флаге, выведут сообщение в $stderr, если исключение приведёт к завершению потока.

Thread.new { 1.times { raise } }

выведет в $stderr следующее сообщение:

#<Thread:...> terminated with exception (report_on_exception is true):
Traceback (most recent call last):
        2: from -e:1:in `block in <main>'
        1: from -e:1:in `times'

Это сделано для раннего обнаружения ошибок в потоках. В некоторых случаях такой вывод может быть нежелателен. Есть несколько способов избежать его:

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

  • Если исключение возникло намеренно, лучше перехватить его ближе к месту возникновения, а не позволять ему завершить Thread.

  • Если гарантируется, что для Thread будет вызван Thread#join или Thread#value, то при запуске Thread можно безопасно отключить вывод сообщения с помощью Thread.current.report_on_exception = false. Однако в этом случае исключение может быть обработано значительно позже или не обработано вовсе, если для Thread так и не будет вызван метод присоединения из-за блокировки родительского потока и т. п.

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

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

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

Возвращает новое состояние. Если установлено значение true, все потоки, созданные после этого, унаследуют это условие и выведут сообщение в $stderr, если исключение приведёт к завершению потока:

Thread.report_on_exception = true
t1 = Thread.new do
  puts  "In new thread"
  raise "Exception from thread"
end
sleep(1)
puts "In the main thread"

В результате будет получено:

In new thread
#<Thread:...prog.rb:2> terminated with exception (report_on_exception is true):
Traceback (most recent call last):
prog.rb:4:in `block in <main>': Exception from thread (RuntimeError)
In the main thread

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

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

start([args]*) {|args| block } → thread Показать исходный код
fork([args]*) {|args| block } → thread
static VALUE
thread_start(VALUE klass, VALUE args)
{
    struct thread_create_params params = {
        .type = thread_invoke_type_proc,
        .args = args,
        .proc = rb_block_proc(),
    };
    return thread_create_core(rb_thread_alloc(klass), &params);
}

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

stop → nil Показать исходный код
static VALUE
thread_stop(VALUE _)
{
    return rb_thread_stop();
}

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

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 Показать исходный код
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

Thread#[] и Thread#[]= относятся не к потоку, а к волокну. В 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 Показать исходный код
static VALUE
rb_thread_aset(VALUE self, VALUE id, VALUE val)
{
    return rb_thread_local_aset(self, rb_to_id(id), val);
}

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

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

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

abort_on_exception → true or false Показать исходный код
static VALUE
rb_thread_abort_exc(VALUE thread)
{
    return RBOOL(rb_thread_ptr(thread)->abort_on_exception);
}

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

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

См. также abort_on_exception=.

Также существует метод класса, задающий это значение для всех потоков; см. ::abort_on_exception.

abort_on_exception= boolean → true or false Показать исходный код
static VALUE
rb_thread_abort_exc_set(VALUE thread, VALUE val)
{
    rb_thread_ptr(thread)->abort_on_exception = RTEST(val);
    return val;
}

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

См. также abort_on_exception.

Также существует метод класса, задающий это значение для всех потоков; см. ::abort_on_exception=.

add_trace_func(proc) → proc Показать исходный код
static VALUE
thread_add_trace_func_m(VALUE obj, VALUE trace)
{
    thread_add_trace_func(GET_EC(), rb_thread_ptr(obj), trace);
    return trace;
}

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

См. Thread#set_trace_func и Kernel#set_trace_func.

alive? → true or false Показать исходный код
static VALUE
rb_thread_alive_p(VALUE thread)
{
    return RBOOL(!thread_finished(rb_thread_ptr(thread)));
}

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

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

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

backtrace → array or nil Показать исходный код
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 Показать исходный код
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

Завершает thr и передаёт выполнение другому потоку, возвращая завершённый Thread. Если это главный или последний поток, процесс завершается. Обратите внимание: вызывающий код не ждёт завершения потока, если получатель — не текущий выполняющийся поток. Завершение происходит асинхронно, и перед выходом поток ещё может выполнить небольшой объём кода Ruby.

Псевдоним: kill
fetch(sym) → obj Показать исходный код
fetch(sym) { } → obj
fetch(sym, default) → obj
static VALUE
rb_thread_fetch(int argc, VALUE *argv, VALUE self)
{
    VALUE key, val;
    ID id;
    rb_thread_t *target_th = rb_thread_ptr(self);
    int block_given;

    rb_check_arity(argc, 1, 2);
    key = argv[0];

    block_given = rb_block_given_p();
    if (block_given && argc == 2) {
        rb_warn("block supersedes default value argument");
    }

    id = rb_check_id(&key);

    if (id == recursive_key) {
        return target_th->ec->local_storage_recursive_hash;
    }
    else if (id && target_th->ec->local_storage &&
             rb_id_table_lookup(target_th->ec->local_storage, id, &val)) {
        return val;
    }
    else if (block_given) {
        return rb_yield(key);
    }
    else if (argc == 1) {
        rb_key_err_raise(rb_sprintf("key not found: %+"PRIsVALUE, key), self, key);
    }
    else {
        return argv[1];
    }
}

Возвращает значение локальной переменной волокна для заданного ключа. Если ключ не найден, возможны следующие варианты: если другие аргументы не указаны, будет вызвано исключение KeyError; если задан default, будет возвращено это значение; если указан необязательный блок кода, он будет выполнен, а его результат возвращён. См. Thread#[] и Hash#fetch.

group → thgrp or nil Показать исходный код
VALUE
rb_thread_group(VALUE thread)
{
    return rb_thread_ptr(thread)->thgroup;
}

Возвращает ThreadGroup, содержащую заданный поток.

Thread.main.group   #=> #<ThreadGroup:0x4029d914>
inspect
Псевдоним: to_s
join → thr Показать исходный код
join(limit) → thr
static VALUE
thread_join_m(int argc, VALUE *argv, VALUE self)
{
    VALUE timeout = Qnil;
    rb_hrtime_t rel = 0, *limit = 0;

    if (rb_check_arity(argc, 0, 1)) {
        timeout = argv[0];
    }

    // Convert the timeout eagerly, so it's always converted and deterministic
    /*
     * This supports INFINITY and negative values, so we can't use
     * rb_time_interval right now...
     */
    if (NIL_P(timeout)) {
        /* unlimited */
    }
    else if (FIXNUM_P(timeout)) {
        rel = rb_sec2hrtime(NUM2TIMET(timeout));
        limit = &rel;
    }
    else {
        limit = double2hrtime(&rel, rb_num2dbl(timeout));
    }

    return thread_join(rb_thread_ptr(self), timeout, limit);
}

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

Метод не возвращает управление, пока thr не завершится или пока не пройдёт заданное количество секунд — limit.

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

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

Если 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 Показать исходный код
static VALUE
rb_thread_key_p(VALUE self, VALUE key)
{
    VALUE val;
    ID id = rb_check_id(&key);
    struct rb_id_table *local_storage = rb_thread_ptr(self)->ec->local_storage;

    if (!id || local_storage == NULL) {
        return Qfalse;
    }
    return RBOOL(rb_id_table_lookup(local_storage, id, &val));
}

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

me = Thread.current
me[:oliver] = "a"
me.key?(:oliver)    #=> true
me.key?(:stanley)   #=> false
keys → array Показать исходный код
static VALUE
rb_thread_keys(VALUE self)
{
    struct rb_id_table *local_storage = rb_thread_ptr(self)->ec->local_storage;
    VALUE ary = rb_ary_new();

    if (local_storage) {
        rb_id_table_foreach(local_storage, thread_keys_i, (void *)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]
kill → thr Показать исходный код
VALUE
rb_thread_kill(VALUE thread)
{
    rb_thread_t *target_th = rb_thread_ptr(thread);

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

    RUBY_DEBUG_LOG("target_th:%u", rb_th_serial(target_th));

    if (target_th == GET_THREAD()) {
        /* kill myself immediately */
        rb_threadptr_to_kill(target_th);
    }
    else {
        threadptr_check_pending_interrupt_queue(target_th);
        rb_threadptr_pending_interrupt_enque(target_th, RUBY_FATAL_THREAD_KILLED);
        rb_threadptr_interrupt(target_th);
    }

    return thread;
}

Завершает thr и передаёт выполнение другому потоку, возвращая завершённый Thread. Если это главный или последний поток, процесс завершается. Обратите внимание: вызывающий код не ждёт завершения потока, если получатель — не текущий выполняющийся поток. Завершение происходит асинхронно, и перед выходом поток ещё может выполнить небольшой объём кода Ruby.

Также имеет псевдонимы: terminate, exit
name → string Показать исходный код
static VALUE
rb_thread_getname(VALUE thread)
{
    return rb_thread_ptr(thread)->name;
}

Показывает имя потока.

name=(name) → string Показать исходный код
static VALUE
rb_thread_setname(VALUE thread, VALUE name)
{
    rb_thread_t *target_th = rb_thread_ptr(thread);

    if (!NIL_P(name)) {
        rb_encoding *enc;
        StringValueCStr(name);
        enc = rb_enc_get(name);
        if (!rb_enc_asciicompat(enc)) {
            rb_raise(rb_eArgError, "ASCII incompatible encoding (%s)",
                     rb_enc_name(enc));
        }
        name = rb_str_new_frozen(name);
    }
    target_th->name = name;
    if (threadptr_initialized(target_th) && target_th->has_dedicated_nt) {
        native_set_another_thread_name(target_th->nt->thread_id, name);
    }
    return name;
}

Задаёт указанное имя потоку Ruby. На некоторых платформах имя может быть задано для pthread и/или ядра ОС.

native_thread_id → integer Показать исходный код
static VALUE
rb_thread_native_thread_id(VALUE thread)
{
    rb_thread_t *target_th = rb_thread_ptr(thread);
    if (rb_threadptr_dead(target_th)) return Qnil;
    return native_thread_native_thread_id(target_th);
}

Возвращает идентификатор нативного потока, используемого потоком Ruby.

Идентификатор зависит от ОС. (Это не идентификатор POSIX-потока, возвращаемый pthread_self(3).)

  • В Linux это TID, возвращаемый gettid(2).

  • В macOS это уникальный для всей системы целочисленный идентификатор потока, возвращаемый pthread_threadid_np(3).

  • Во FreeBSD это уникальный целочисленный идентификатор потока, возвращаемый pthread_getthreadid_np(3).

  • В Windows это идентификатор потока, возвращаемый GetThreadId().

  • На других платформах вызывается исключение NotImplementedError.

ПРИМЕЧАНИЕ: если поток ещё не связан с нативным потоком или уже отвязан от него, возвращается nil. Если реализация Ruby использует модель потоков M:N, идентификатор может изменяться в зависимости от момента выполнения.

pending_interrupt?(error = nil) → true/false Показать исходный код
static VALUE
rb_thread_pending_interrupt_p(int argc, VALUE *argv, VALUE target_thread)
{
    rb_thread_t *target_th = rb_thread_ptr(target_thread);

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

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

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

Дополнительные сведения см. в описании ::pending_interrupt?.

priority → integer Показать исходный код
static VALUE
rb_thread_priority(VALUE thread)
{
    return INT2NUM(rb_thread_ptr(thread)->priority);
}

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

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

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

#if USE_NATIVE_THREAD_PRIORITY
    target_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;
    }
    target_th->priority = (int8_t)priority;
#endif
    return INT2NUM(target_th->priority);
}

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

Это лишь подсказка для планировщика потоков 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(exception, message = exception.to_s, backtrace = nil, cause: $!) Показать исходный код
raise(message = nil, cause: $!)
static VALUE
thread_raise_m(int argc, VALUE *argv, VALUE self)
{
    rb_thread_t *target_th = rb_thread_ptr(self);
    const rb_thread_t *current_th = GET_THREAD();

    threadptr_check_pending_interrupt_queue(target_th);
    rb_threadptr_raise(target_th, argc, argv);

    /* To perform Thread.current.raise as Kernel.raise */
    if (current_th == target_th) {
        RUBY_VM_CHECK_INTS(target_th->ec);
    }
    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
report_on_exception → true or false Показать исходный код
static VALUE
rb_thread_report_exc(VALUE thread)
{
    return RBOOL(rb_thread_ptr(thread)->report_on_exception);
}

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

Значением по умолчанию при создании Thread является значение глобального флага Thread.report_on_exception.

См. также report_on_exception=.

Также существует метод класса, задающий это значение для всех новых потоков; см. ::report_on_exception=.

report_on_exception= boolean → true or false Показать исходный код
static VALUE
rb_thread_report_exc_set(VALUE thread, VALUE val)
{
    rb_thread_ptr(thread)->report_on_exception = RTEST(val);
    return val;
}

Если задано значение true, при завершении этого thr из-за исключения в $stderr выводится сообщение. Подробности см. в описании ::report_on_exception.

См. также report_on_exception.

Также существует метод класса, задающий это значение для всех новых потоков; см. ::report_on_exception=.

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.

set_trace_func(proc) → proc Показать исходный код
set_trace_func(nil) → nil
static VALUE
thread_set_trace_func_m(VALUE target_thread, VALUE trace)
{
    rb_execution_context_t *ec = GET_EC();
    rb_thread_t *target_th = rb_thread_ptr(target_thread);

    rb_threadptr_remove_event_hook(ec, target_th, call_trace_func, Qundef);

    if (NIL_P(trace)) {
        return Qnil;
    }
    else {
        thread_add_trace_func(ec, target_th, trace);
        return trace;
    }
}

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

См. Kernel#set_trace_func.

status → string, false or nil Показать исходный код
static VALUE
rb_thread_status(VALUE thread)
{
    rb_thread_t *target_th = rb_thread_ptr(thread);

    if (rb_threadptr_dead(target_th)) {
        if (!NIL_P(target_th->ec->errinfo) &&
            !FIXNUM_P(target_th->ec->errinfo)) {
            return Qnil;
        }
        else {
            return Qfalse;
        }
    }
    else {
        return rb_str_new2(thread_status_name(target_th, FALSE));
    }
}

Возвращает состояние 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 or false Показать исходный код
static VALUE
rb_thread_stop_p(VALUE thread)
{
    rb_thread_t *th = rb_thread_ptr(thread);

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

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

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

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

terminate → thr

Завершает thr и передаёт выполнение другому потоку, возвращая завершённый Thread. Если это главный или последний поток, процесс завершается. Обратите внимание: вызывающий код не ждёт завершения потока, если получатель — не текущий выполняющийся поток. Завершение происходит асинхронно, и перед выходом поток ещё может выполнить небольшой объём кода Ruby.

Псевдоним: kill
thread_variable?(key) → true or false Показать исходный код
static VALUE
rb_thread_variable_p(VALUE thread, VALUE key)
{
    VALUE locals;
    VALUE symbol = rb_to_symbol(key);

    if (LIKELY(!THREAD_LOCAL_STORAGE_INITIALISED_P(thread))) {
        return Qfalse;
    }
    locals = rb_thread_local_storage(thread);

    return RBOOL(rb_hash_lookup(locals, symbol) != Qnil);
}

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

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

Обратите внимание: это не локальные переменные волокна. Дополнительные сведения см. в описаниях Thread#[] и Thread#thread_variable_get.

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

    if (LIKELY(!THREAD_LOCAL_STORAGE_INITIALISED_P(thread))) {
        return Qnil;
    }
    locals = rb_thread_local_storage(thread);
    return rb_hash_aref(locals, symbol);
}

Возвращает значение заданной локальной переменной потока. Обратите внимание: эти значения отличаются от локальных значений волокна. Сведения о локальных значениях волокна см. в описаниях Thread#[] и Thread#[]=.

Локальные значения 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(key, value) Показать исходный код
static VALUE
rb_thread_variable_set(VALUE thread, VALUE key, VALUE val)
{
    VALUE locals;

    if (OBJ_FROZEN(thread)) {
        rb_frozen_error_raise(thread, "can't modify frozen thread locals");
    }

    locals = rb_thread_local_storage(thread);
    return rb_hash_aset(locals, rb_to_symbol(key), val);
}

Задаёт локальную переменную потока с ключом key и значением value. Обратите внимание: эти переменные локальны для потоков, а не для волокон. Дополнительные сведения см. в описаниях Thread#thread_variable_get и Thread#[].

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

    ary = rb_ary_new();
    if (LIKELY(!THREAD_LOCAL_STORAGE_INITIALISED_P(thread))) {
        return ary;
    }
    locals = rb_thread_local_storage(thread);
    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#[] и Thread#thread_variable_get.

to_s → string Показать исходный код
static VALUE
rb_thread_to_s(VALUE thread)
{
    VALUE cname = rb_class_path(rb_obj_class(thread));
    rb_thread_t *target_th = rb_thread_ptr(thread);
    const char *status;
    VALUE str, loc;

    status = thread_status_name(target_th, TRUE);
    str = rb_sprintf("#<%"PRIsVALUE":%p", cname, (void *)thread);
    if (!NIL_P(target_th->name)) {
        rb_str_catf(str, "@%"PRIsVALUE, target_th->name);
    }
    if ((loc = threadptr_invoke_proc_location(target_th)) != Qnil) {
        rb_str_catf(str, " %"PRIsVALUE":%"PRIsVALUE,
                    RARRAY_AREF(loc, 0), RARRAY_AREF(loc, 1));
    }
    rb_str_catf(str, " %s>", status);

    return str;
}

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

Также имеет псевдоним: inspect
value → obj Показать исходный код
static VALUE
thread_value(VALUE self)
{
    rb_thread_t *th = rb_thread_ptr(self);
    thread_join(th, Qnil, 0);
    if (UNDEF_P(th->value)) {
        // If the thread is dead because we forked th->value is still Qundef.
        return Qnil;
    }
    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–2025 Yukihiro Matsumoto
Licensed under the Ruby License.
Ruby Standard Library © contributors
Licensed under their own licenses.

Spec-Zone.ru

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