class Thread
Потоки являются реализацией Ruby для модели конкурентного программирования.
Программы, требующие нескольких потоков выполнения, идеально подходят для класса Thread в Ruby.
Например, мы можем создать новый поток, отделенный от выполнения основного потока, используя ::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 }
Thread initialization
Для создания новых потоков Ruby предоставляет ::new, ::start и ::fork. Для каждого из этих методов необходимо указать блок, иначе будет выброшено исключение ThreadError.
При создании подкласса класса Thread метод initialize вашего подкласса будет игнорироваться ::start и ::fork. В противном случае обязательно вызовите super в вашем методе initialize.
Thread termination
Для завершения потоков Ruby предоставляет множество способов сделать это.
Метод класса ::kill предназначен для завершения заданного потока:
thr = Thread.new { ... }
Thread.kill(thr) # sends exit() to thr В качестве альтернативы вы можете использовать метод экземпляра exit или любой из его псевдонимов kill или terminate.
thr.exit
Thread status
Ruby предоставляет несколько методов экземпляра для запроса состояния заданного потока. Чтобы получить строку с текущим состоянием потока, используйте status
thr = Thread.new { sleep }
thr.status # => "sleep"
thr.exit
thr.status # => false
Вы также можете использовать alive?, чтобы узнать, работает ли поток или находится в состоянии ожидания, и stop?, если поток завершен или находится в состоянии ожидания.
Thread variables and scope
Поскольку потоки создаются с блоками, к ним применяются те же правила, что и к другим блокам Ruby для области видимости переменных. Любые локальные переменные, созданные внутри этого блока, доступны только этому потоку.
Fiber-local vs. Thread-local
Каждый волокно имеет собственный бак для хранения #[]. Когда вы устанавливаете новое локальное волокно, оно доступно только внутри этого 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 handling
Любой поток может вызвать исключение, используя метод экземпляра raise, который работает аналогично Kernel#raise.
Однако важно отметить, что исключение, возникающее в любом потоке, кроме основного, зависит от abort_on_exception. Этот параметр по умолчанию false, что означает, что любое необработанное исключение приведет к тихому завершению потока при ожидании его join или value. Вы можете изменить это значение по умолчанию, используя abort_on_exception= true или установив $DEBUG в true.
С добавлением метода класса ::handle_interrupt вы теперь можете асинхронно обрабатывать исключения с потоками.
Scheduling
Ruby предоставляет несколько способов поддержки планирования потоков в вашей программе.
Первый способ — использование метода класса ::stop для перевода текущего работающего потока в спящий режим и планирования выполнения другого потока.
После того как поток перешел в спящий режим, вы можете использовать метод экземпляра wakeup для обозначения вашего потока как пригодного для планирования.
Вы также можете попробовать ::pass, который пытается передать выполнение другому потоку, но зависит от ОС, будет ли выполняться переключение работающего потока или нет. То же самое относится к priority, который позволяет вам намекнуть планировщику потоков, какие потоки вы хотите, чтобы они имели приоритет при передаче выполнения. Этот метод также зависит от ОС и может игнорироваться на некоторых платформах.
Публичные методы класса
static VALUE
rb_thread_s_debug(void)
{
return INT2NUM(rb_thread_debug_enabled);
} Возвращает уровень отладки потока. Доступно только при компиляции с THREAD_DEBUG=-1.
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.
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.
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>
# File prelude.rb, line 10
def self.exclusive
warn "Thread.exclusive is deprecated, use Mutex", caller
MUTEX_FOR_THREAD_EXCLUSIVE.synchronize{
yield
}
end Заключает блок в один, глобальный для VM Thread::Mutex#synchronize, возвращая значение блока. Поток, выполняющийся внутри раздела exclusive, будет блокировать только другие потоки, которые также используют механизм ::exclusive.
static VALUE
rb_thread_exit(void)
{
rb_thread_t *th = GET_THREAD();
return rb_thread_kill(th->self);
} Завершает текущий выполняющийся поток и планирует запуск другого потока.
Если этот поток уже помечен для завершения, ::exit возвращает Thread.
Если это главный поток или последний поток, завершить процесс.
static VALUE
rb_thread_s_handle_interrupt(VALUE self, VALUE mask_arg)
{
VALUE mask;
rb_thread_t *th = GET_THREAD();
volatile VALUE r = Qnil;
int state;
if (!rb_block_given_p()) {
rb_raise(rb_eArgError, "block is needed.");
}
mask = 0;
mask_arg = rb_convert_type(mask_arg, T_HASH, "Hash", "to_hash");
rb_hash_foreach(mask_arg, handle_interrupt_arg_check_i, (VALUE)&mask);
if (!mask) {
return rb_yield(Qnil);
}
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);
}
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.
}
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
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>
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_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 (!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);
} Возвращает, пуста ли асинхронная очередь.
Так как ::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 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"
Публичные методы экземпляра
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.
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.
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.
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=.
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 в качестве обработчика трассировки.
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
static VALUE
rb_thread_backtrace_m(int argc, VALUE *argv, VALUE thval)
{
return rb_vm_thread_backtrace(argc, argv, thval);
} Возвращает текущий стек вызовов целевого потока.
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, за исключением того, что он применяется к определенному потоку.
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.
Если это основной поток или последний поток, завершает процесс.
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>
static VALUE
rb_thread_inspect(VALUE thread)
{
VALUE cname = rb_class_path(rb_obj_class(thread));
rb_thread_t *th;
const char *status;
VALUE str;
GetThreadPtr(thread, th);
status = thread_status_name(th);
str = rb_sprintf("#<%"PRIsVALUE":%p", cname, (void *)thread);
if (!NIL_P(th->name)) {
rb_str_catf(str, "@%"PRIsVALUE, th->name);
}
if (!th->first_func && th->first_proc) {
VALUE loc = rb_proc_location(th->first_proc);
if (!NIL_P(loc)) {
const VALUE *ptr = RARRAY_CONST_PTR(loc);
rb_str_catf(str, "@%"PRIsVALUE":%"PRIsVALUE, ptr[0], ptr[1]);
rb_gc_force_recycle(loc);
}
}
rb_str_catf(str, " %s>", status);
OBJ_INFECT(str, thread);
return str;
} Выгружает имя, идентификатор и статус 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...
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
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]
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.
Если это основной поток или последний поток, завершает процесс.
static VALUE
rb_thread_getname(VALUE thread)
{
rb_thread_t *th;
GetThreadPtr(thread, th);
return th->name;
} Показывает имя потока.
static VALUE
rb_thread_setname(VALUE thread, VALUE name)
{
#ifdef SET_ANOTHER_THREAD_NAME
const char *s = "";
#endif
rb_thread_t *th;
GetThreadPtr(thread, th);
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);
#ifdef SET_ANOTHER_THREAD_NAME
s = RSTRING_PTR(name);
#endif
}
th->name = name;
#if defined(SET_ANOTHER_THREAD_NAME)
if (threadptr_initialized(th)) {
SET_ANOTHER_THREAD_NAME(th->thread_id, s);
}
#endif
return name;
} Устанавливает данное имя для потока Ruby. На некоторых платформах это может установить имя для pthread и/или ядра.
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? для получения дополнительной информации.
static VALUE
rb_thread_priority(VALUE thread)
{
rb_thread_t *th;
GetThreadPtr(thread, th);
return INT2NUM(th->priority);
} Возвращает приоритет thr. По умолчанию наследуется от текущего потока, создающего новый поток, или ноль для исходного основного потока; поток с более высоким приоритетом будет выполняться чаще, чем потоки с более низким приоритетом (но потоки с более низким приоритетом также могут выполняться).
Это всего лишь подсказка для планировщика потоков Ruby. Она может игнорироваться на некоторых платформах.
Thread.current.priority #=> 0
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 в 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
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
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
rb_thread_safe_level(VALUE thread)
{
rb_thread_t *th;
GetThreadPtr(thread, th);
return INT2NUM(th->safe_level);
} Возвращает уровень безопасности, действующий для thr. Установка уровней безопасности, локальных для потока, может помочь при реализации песочниц, которые запускают небезопасный код.
thr = Thread.new { $SAFE = 1; sleep }
Thread.current.safe_level #=> 0
thr.safe_level #=> 1
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.
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"
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
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.
Если это основной поток или последний поток, завершает процесс.
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 для получения более подробной информации.
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 возвращается для локальной переменной волокна. Волокно выполняется в том же потоке, поэтому локальные значения потока доступны.
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 и #[] для получения дополнительной информации.
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 для получения более подробной информации.
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
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.