класс Thread
Потоки — это реализация 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-локальные vs. Thread-локальные
Каждый волокно имеет свой собственный контейнер для хранения Thread#[]. Когда вы устанавливаете новую волокно-локальную переменную, она доступна только внутри этого Fiber.
Thread.new {
Thread.current[:foo] = "bar"
Fiber.new {
p Thread.current[:foo] # => nil
}.resume
}.join
В этом примере используется [] для получения и []= для установки волокно-локальных переменных. Вы также можете использовать keys для перечисления волокно-локальных переменных для заданного потока и key? для проверки наличия волокно-локальной переменной.
Что касается 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-локальной переменной.
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, который позволяет указать планировщику потоков, какие потоки вы хотите сделать приоритетными при передаче выполнения. Этот метод также зависит от операционной системы и может быть проигнорирован на некоторых платформах.
Публичные методы класса
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.
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=.
static VALUE
thread_s_current(VALUE klass)
{
return rb_thread_current();
} Возвращает текущий выполняющийся поток.
Thread.current #=> #<Thread:0x401bdf4c run>
static VALUE
each_caller_location(VALUE unused)
{
rb_ec_partial_backtrace_object(GET_EC(), 2, ALL_BACKTRACE_LINES, NULL, FALSE, TRUE);
return Qnil;
} Передает каждый кадр текущего стека выполнения как объект местоположения трассировки стека.
static VALUE
rb_thread_exit(VALUE _)
{
rb_thread_t *th = GET_THREAD();
return rb_thread_kill(th->self);
} Завершает текущий выполняющийся поток и планирует выполнение другого потока.
Если этот поток уже помечен для завершения, ::exit возвращает 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), ¶ms);
} В основном то же самое, что и ::new. Однако, если класс Thread является подклассом, то вызов start в этом подклассе не вызовет метод initialize подкласса.
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_RAW(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 -
Вызывать прерывания во время BlockingOperation.
-
:never -
Никогда не вызывать все прерывания.
BlockingOperation означает, что операция будет блокировать вызывающий поток, например, чтение и запись. В реализации CRuby, BlockingOperation — это любая операция, выполняемая без GVL.
Замаскированные асинхронные прерывания откладываются до тех пор, пока они не будут включены. Этот метод аналогичен sigprocmask(3).
ЗАМЕТКА
Асинхронные прерывания сложно использовать.
Если вам нужно взаимодействовать между потоками, рассмотрите другой способ, например, Queue.
Или используйте их с глубоким пониманием этого метода.
Использование
В этом примере мы можем защититься от исключений Thread#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.
}
Для обработки всех прерываний используйте Object а не Exception в качестве ExceptionClass, поскольку прерывания kill/terminate не обрабатываются Exception.
static VALUE
rb_thread_s_ignore_deadlock(VALUE _)
{
return RBOOL(GET_THREAD()->vm->thread_ignore_deadlock);
} Возвращает состояние глобального условия «игнорировать взаимоблокировку». По умолчанию false, поэтому условия взаимоблокировки не игнорируются.
См. также ::ignore_deadlock=.
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.
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
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>
static VALUE
rb_thread_s_main(VALUE klass)
{
return rb_thread_main();
} Возвращает главный поток.
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]
Исключение ThreadError возникает, если ::new вызывается без блока кода.
Если вы собираетесь создавать подкласс Thread, обязательно вызовите super в методе initialize, иначе будет выброшено исключение ThreadError.
static VALUE
thread_s_pass(VALUE klass)
{
rb_thread_schedule();
return Qnil;
} Дает указание планировщику потоков передать выполнение другому потоку. Действующий поток может или не может переключится, это зависит от ОС и процессора.
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 static VALUE
rb_thread_s_report_exc(VALUE _)
{
return RBOOL(GET_THREAD()->vm->thread_report_on_exception);
} Возвращает состояние глобальной опции «сообщать об исключениях».
По умолчанию — true с 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.current.report_on_exception = falseпри запускеThread. Однако это может обработать исключение гораздо позже или вообще не обработать, еслиThreadне будет объединён, так как родительский поток заблокирован и т. д.
См. также ::report_on_exception=.
Также есть метод уровня экземпляра для установки этого значения для конкретного потока, см. report_on_exception=.
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=.
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), ¶ms);
} По сути, то же самое, что и ::new. Однако, если класс Thread унаследован, вызов start в этом подклассе не вызовет метод initialize подкласса.
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"
Публичные методы экземпляра
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.
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.
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.
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=.
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 в качестве обработчика для трассировки.
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
static VALUE
rb_thread_backtrace_m(int argc, VALUE *argv, VALUE thval)
{
return rb_vm_thread_backtrace(argc, argv, thval);
} Возвращает текущий backtrace целевого потока.
static VALUE
rb_thread_backtrace_locations_m(int argc, VALUE *argv, VALUE thval)
{
return rb_vm_thread_backtrace_locations(argc, argv, thval);
} Возвращает стек выполнения для целевого потока—массив, содержащий объекты местоположения backtrace.
См. Thread::Backtrace::Location для получения дополнительной информации.
Этот метод ведет себя аналогично Kernel#caller_locations, за исключением того, что он применяется к определенному потоку.
Завершает thr и планирует выполнение другого потока, возвращая завершенный Thread. Если это основной поток или последний поток, завершает процесс.
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.
VALUE
rb_thread_group(VALUE thread)
{
return rb_thread_ptr(thread)->thgroup;
} Возвращает ThreadGroup, содержащий данный поток.
Thread.main.group #=> #<ThreadGroup:0x4029d914>
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.
Все потоки, не объединенные, будут завершены при выходе из основной программы.
Если 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...
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
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]
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. Если это основной поток или последний поток, завершает процесс.
static VALUE
rb_thread_getname(VALUE thread)
{
return rb_thread_ptr(thread)->name;
} отобразить имя потока.
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 и/или ядра.
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, идентификатор может измениться в зависимости от времени.
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? для получения дополнительной информации.
static VALUE
rb_thread_priority(VALUE thread)
{
return INT2NUM(rb_thread_ptr(thread)->priority);
} Возвращает приоритет thr. По умолчанию наследуется от текущего потока, создавшего новый поток, или ноль для начального главного потока; поток с более высоким приоритетом будет выполняться чаще, чем потоки с более низким приоритетом (но потоки с более низким приоритетом также могут выполняться).
Это просто подсказка для планировщика потоков Ruby. Может быть проигнорировано в некоторых платформах.
Thread.current.priority #=> 0
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 на целое число. Потоки с более высоким приоритетом будут выполняться чаще, чем потоки с более низким приоритетом (но потоки с более низким приоритетом также могут выполняться).
Это просто подсказка для планировщика потоков 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
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
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=.
static VALUE
rb_thread_report_exc_set(VALUE thread, VALUE val)
{
rb_thread_ptr(thread)->report_on_exception = RTEST(val);
return val;
} При установке в true, сообщение выводится в $stderr, если исключение останавливает этот thr. См. ::report_on_exception для получения подробной информации.
См. также report_on_exception.
Также есть метод уровня класса для установки этого значения для всех новых потоков, см. ::report_on_exception=.
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.
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.
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"
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
Прерывает thr и планирует запуск другого потока, возвращая завершенный Thread. Если это главный поток или последний поток, завершает процесс.
static VALUE
rb_thread_variable_p(VALUE thread, VALUE key)
{
VALUE locals;
if (LIKELY(!THREAD_LOCAL_STORAGE_INITIALISED_P(thread))) {
return Qfalse;
}
locals = rb_thread_local_storage(thread);
return RBOOL(rb_hash_lookup(locals, rb_to_symbol(key)) != 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 для получения дополнительной информации.
static VALUE
rb_thread_variable_get(VALUE thread, VALUE key)
{
VALUE locals;
if (LIKELY(!THREAD_LOCAL_STORAGE_INITIALISED_P(thread))) {
return Qnil;
}
locals = rb_thread_local_storage(thread);
return rb_hash_aref(locals, rb_to_symbol(key));
} Возвращает значение локальной переменной потока, если оно было установлено. Обратите внимание, что эти переменные отличаются от локальных переменных волокна. Для локальных переменных волокна см. 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. Волокно выполняется в том же потоке, поэтому доступны значения локальных переменных потока.
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#[].
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.
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.
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
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–2022 Yukihiro Matsumoto
Licensed under the Ruby License.
Ruby Standard Library © contributors
Licensed under their own licenses.