класс Ractor
Ractor — это абстракция модели Актора для Ruby, которая обеспечивает безопасную параллельную обработку в потоках.
Ractor.new позволяет создать новый Ractor, который будет выполняться параллельно.
# The simplest ractor
r = Ractor.new {puts "I am in Ractor!"}
r.take # wait for it to finish
# here "I am in Ractor!" would be printed
Реакторы не используют общие объекты, поэтому проблемы безопасности потоков, такие как гонки данных и условия гонки, не возникают при программировании с несколькими реакторами.
Для достижения этого реакторы сильно ограничивают совместное использование объектов между разными реакторами. Например, в отличие от потоков, реакторы не могут получить доступ к объектам друг друга, а также к объектам через переменные внешней области видимости.
a = 1
r = Ractor.new {puts "I am in Ractor! a=#{a}"}
# fails immediately with
# ArgumentError (can not isolate a Proc because it accesses outer variables (a).)
В CRuby (по умолчанию) блокировка глобальной виртуальной машины (GVL) используется для каждого реактора, поэтому реакторы выполняются параллельно без блокировки друг друга.
Вместо доступа к общему состоянию объекты должны передаваться между реакторами посредством отправки и получения объектов в качестве сообщений.
a = 1
r = Ractor.new do
a_in_ractor = receive # receive blocks till somebody will pass message
puts "I am in Ractor! a=#{a_in_ractor}"
end
r.send(a) # pass it
r.take
# here "I am in Ractor! a=1" would be printed
Существует две пары методов для отправки/приёма сообщений:
-
Ractor#sendиRactor.receiveдля случаев, когда отправитель знает получателя (push); -
Ractor.yieldиRactor#takeдля случаев, когда получатель знает отправителя (pull);
Кроме того, аргумент к Ractor.new будет передан в блок и доступен в нём, как если бы он был получен Ractor.receive, а последнее значение блока будет отправлено за пределы реактора, как если бы оно было отправлено Ractor.yield.
Небольшой пример с классической игрой «пинг-понг»:
server = Ractor.new do
puts "Server starts: #{self.inspect}"
puts "Server sends: ping"
Ractor.yield 'ping' # The server doesn't know the receiver and sends to whoever interested
received = Ractor.receive # The server doesn't know the sender and receives from whoever sent
puts "Server received: #{received}"
end
client = Ractor.new(server) do |srv| # The server is sent inside client, and available as srv
puts "Client starts: #{self.inspect}"
received = srv.take # The Client takes a message specifically from the server
puts "Client received from " \
"#{srv.inspect}: #{received}"
puts "Client sends to " \
"#{srv.inspect}: pong"
srv.send 'pong' # The client sends a message specifically to the server
end
[client, server].each(&:take) # Wait till they both finish
Это выведет:
Server starts: #<Ractor:#2 test.rb:1 running> Server sends: ping Client starts: #<Ractor:#3 test.rb:8 running> Client received from #<Ractor:#2 rac.rb:1 blocking>: ping Client sends to #<Ractor:#2 rac.rb:1 blocking>: pong Server received: pong
Считается, что Ractor получает сообщения через входящий порт и отправляет их в исходящий порт. Любой из них можно отключить с помощью Ractor#close_incoming и Ractor#close_outgoing соответственно. Если реактор завершил работу, его порты будут автоматически закрыты.
Делимые и не делимые объекты
При отправке и получении объекта в реактор важно понимать, является ли объект делимым или нет. Большинство объектов — это не делимые объекты.
Делимые объекты — это, по сути, те, которые могут использоваться несколькими потоками без нарушения безопасности потоков; например, неизменяемые. Ractor.shareable? позволяет проверить это, а Ractor.make_shareable пытается сделать объект делимым, если он таковым не является.
Ractor.shareable?(1) #=> true -- numbers and other immutable basic values are
Ractor.shareable?('foo') #=> false, unless the string is frozen due to # freeze_string_literals: true
Ractor.shareable?('foo'.freeze) #=> true
ary = ['hello', 'world']
ary.frozen? #=> false
ary[0].frozen? #=> false
Ractor.make_shareable(ary)
ary.frozen? #=> true
ary[0].frozen? #=> true
ary[1].frozen? #=> true
При отправке делимого объекта (через send или Ractor.yield) дополнительной обработки не происходит, и он становится доступным для обоих реакторов. При отправке не делимого объекта он может быть скопирован или перемещён. Первое является значением по умолчанию, и оно создаёт полную копию объекта, глубоко клонируя не делимые части его структуры.
data = ['foo', 'bar'.freeze]
r = Ractor.new do
data2 = Ractor.receive
puts "In ractor: #{data2.object_id}, #{data2[0].object_id}, #{data2[1].object_id}"
end
r.send(data)
r.take
puts "Outside : #{data.object_id}, #{data[0].object_id}, #{data[1].object_id}"
Это выведет:
In ractor: 340, 360, 320 Outside : 380, 400, 320
(Обратите внимание, что идентификаторы объектов обоих массивов и не замороженных строк внутри массивов изменились внутри реактора, показывая, что это разные объекты. Но элемент второго массива, который является делимой замороженной строкой, имеет тот же идентификатор объекта.)
Глубокое клонирование объектов может быть медленным и иногда невозможным. В качестве альтернативы можно использовать move: true при отправке. Это переместит объект в получающий реактор, сделав его недоступным для отправляющего реактора.
data = ['foo', 'bar']
r = Ractor.new do
data_in_ractor = Ractor.receive
puts "In ractor: #{data_in_ractor.object_id}, #{data_in_ractor[0].object_id}"
end
r.send(data, move: true)
r.take
puts "Outside: moved? #{Ractor::MovedObject === data}"
puts "Outside: #{data.inspect}"
Это выведет:
In ractor: 100, 120 Outside: moved? true test.rb:9:in `method_missing': can not send any methods to a moved object (Ractor::MovedError)
Обратите внимание, что даже inspect (и более базовые методы, такие как __id__) недоступны для перемещённого объекта.
Помимо замороженных объектов, существуют делимые объекты. Объекты Class и Module являются делимыми, поэтому определения Класса/Модуля разделяются между реакторами. Объекты Ractor также являются делимыми объектами. Все операции с делимыми изменяемыми объектами безопасны для потоков, поэтому свойство безопасности потоков сохраняется. В Ruby невозможно определить изменяемые делимые объекты, но расширения C могут их реализовать.
Запрещено обращаться к переменным экземпляра изменяемых делимых объектов (особенно Модулей и классов) из реакторов, отличных от главного:
class C
class << self
attr_accessor :tricky
end
end
C.tricky = 'test'
r = Ractor.new(C) do |cls|
puts "I see #{cls}"
puts "I can't see #{cls.tricky}"
end
r.take
# I see C
# can not access instance variables of classes/modules from non-main Ractors (RuntimeError)
Реакторы могут получать доступ к константам, если они делимые. Главный Ractor — единственный, кто может получить доступ к не делимым константам.
GOOD = 'good'.freeze
BAD = 'bad'
r = Ractor.new do
puts "GOOD=#{GOOD}"
puts "BAD=#{BAD}"
end
r.take
# GOOD=good
# can not access non-shareable objects in constant Object::BAD by non-main Ractor. (NameError)
# Consider the same C class from above
r = Ractor.new do
puts "I see #{C}"
puts "I can't see #{C.tricky}"
end
r.take
# I see C
# can not access instance variables of classes/modules from non-main Ractors (RuntimeError)
См. также описание псевдонима # shareable_constant_value в объяснении синтаксиса комментариев в синтаксисе комментариев.
Реакторы против потоков
Каждый реактор создаёт свой собственный поток. Новые потоки могут быть созданы внутри реактора (и, в CRuby, делят GVL с другими потоками этого реактора).
r = Ractor.new do
a = 1
Thread.new {puts "Thread in ractor: a=#{a}"}.join
end
r.take
# Here "Thread in ractor: a=1" will be printed
Примечание по примерам кода
В примерах ниже иногда используется следующий метод, чтобы дождаться завершения реакторов, которые в данный момент не заблокированы (или обработки до следующей блокировки) метод.
def wait sleep(0.1) end
Это **только для демонстрационных целей** и не должно использоваться в реальном коде. В большинстве случаев используется просто take, чтобы дождаться завершения реактора.
Ссылка
См. документ проектирования реакторов для получения более подробной информации.
Методы публичного класса
# File ractor.rb, line 287
def self.count
__builtin_cexpr! %q{
ULONG2NUM(GET_VM()->ractor.cnt);
}
end Возвращает общее количество Ractor, которые в данный момент выполняются.
Ractor.count #=> 1
r = Ractor.new(name: 'example') { Ractor.yield(1) }
Ractor.count #=> 2 (main + example ractor)
r.take # wait for Ractor.yield(1)
r.take # wait till r will finish
Ractor.count #=> 1
# File ractor.rb, line 273
def self.current
__builtin_cexpr! %q{
rb_ractor_self(rb_ec_ractor_ptr(ec));
}
end Возвращает текущий выполняемый Ractor.
Ractor.current #=> #<Ractor:#1 running>
# File ractor.rb, line 833
def self.main
__builtin_cexpr! %q{
rb_ractor_self(GET_VM()->ractor.main_ractor);
}
end возвращает главный рактор
Сделать obj доступным для использования между ракторами.
obj и все объекты, на которые он ссылается, будут заморожены, если они еще не являются разделяемыми.
Если copy ключевое слово равно true, метод скопирует объекты перед заморозкой. Это более безопасный вариант, но он может быть медленнее.
Обратите внимание, что спецификация и реализация этого метода не являются окончательными и могут быть изменены в будущем.
obj = ['test'] Ractor.shareable?(obj) #=> false Ractor.make_shareable(obj) #=> ["test"] Ractor.shareable?(obj) #=> true obj.frozen? #=> true obj[0].frozen? #=> true # Copy vs non-copy versions: obj1 = ['test'] obj1s = Ractor.make_shareable(obj1) obj1.frozen? #=> true obj1s.object_id == obj1.object_id #=> true obj2 = ['test'] obj2s = Ractor.make_shareable(obj2, copy: true) obj2.frozen? #=> false obj2s.frozen? #=> true obj2s.object_id == obj2.object_id #=> false obj2s[0].object_id == obj2[0].object_id #=> false
См. также раздел «Разделяемые и неразделяемые объекты» в документации класса Ractor.
# File ractor.rb, line 262
def self.new(*args, name: nil, &block)
b = block # TODO: builtin bug
raise ArgumentError, "must be called with a block" unless block
loc = caller_locations(1, 1).first
loc = "#{loc.path}:#{loc.lineno}"
__builtin_ractor_create(loc, name, args, b)
end Создаёт новый Ractor с аргументами и блоком.
Блок (Proc) будет изолирован (не может получить доступ к внешним переменным). self внутри блока будет ссылаться на текущий Ractor.
r = Ractor.new { puts "Hi, I am #{self.inspect}" }
r.take
# Prints "Hi, I am #<Ractor:#2 test.rb:1 running>"
args переданные в метод, будут переданы аргументам блока по тем же правилам, что и объекты, переданные через send/Ractor.receive: если args не разделяемые, они будут скопированы (с помощью глубокого клонирования, что может быть неэффективно).
arg = [1, 2, 3]
puts "Passing: #{arg} (##{arg.object_id})"
r = Ractor.new(arg) {|received_arg|
puts "Received: #{received_arg} (##{received_arg.object_id})"
}
r.take
# Prints:
# Passing: [1, 2, 3] (#280)
# Received: [1, 2, 3] (#300)
name рактора можно установить для отладки:
r = Ractor.new(name: 'my ractor') {}
p r
#=> #<Ractor:#3 my ractor test.rb:1 terminated>
# File ractor.rb, line 415
def self.receive
__builtin_cexpr! %q{
ractor_receive(ec, rb_ec_ractor_ptr(ec))
}
end Получает входящее сообщение из очереди текущего порта Ractor, которое было отправлено туда методом send.
r = Ractor.new do
v1 = Ractor.receive
puts "Received: #{v1}"
end
r.send('message1')
r.take
# Here will be printed: "Received: message1"
В качестве альтернативы можно использовать частный метод экземпляра receive:
r = Ractor.new do
v1 = receive
puts "Received: #{v1}"
end
r.send('message1')
r.take
# Here will be printed: "Received: message1"
Метод блокируется, если очередь пуста.
r = Ractor.new do
puts "Before first receive"
v1 = Ractor.receive
puts "Received: #{v1}"
v2 = Ractor.receive
puts "Received: #{v2}"
end
wait
puts "Still not received"
r.send('message1')
wait
puts "Still received only one"
r.send('message2')
r.take
Вывод:
Before first receive Still not received Received: message1 Still received only one Received: message2
Если для рактора был вызван метод close_incoming, метод вызывает Ractor::ClosedError, если входящей очереди больше нет сообщений:
Ractor.new do close_incoming receive end wait # in `receive': The incoming port is already closed => #<Ractor:#2 test.rb:1 running> (Ractor::ClosedError)
# File ractor.rb, line 493 def self.receive_if &b Primitive.ractor_receive_if b end
Получить только определенное сообщение.
Вместо Ractor.receive, Ractor.receive_if может предоставить шаблон с помощью блока, и вы можете выбрать сообщение для получения.
r = Ractor.new do
p Ractor.receive_if{|msg| msg.match?(/foo/)} #=> "foo3"
p Ractor.receive_if{|msg| msg.match?(/bar/)} #=> "bar1"
p Ractor.receive_if{|msg| msg.match?(/baz/)} #=> "baz2"
end
r << "bar1"
r << "baz2"
r << "foo3"
r.take
Это выведет:
foo3 bar1 baz2
Если блок возвращает истинное значение, сообщение удаляется из очереди входящих сообщений и возвращается. В противном случае сообщение остается в очереди входящих сообщений, и проверяются следующие полученные сообщения с помощью данного блока.
Если сообщений в очереди входящих сообщений больше нет, метод будет блокироваться до тех пор, пока не прибудут новые сообщения.
Если блок выходит с помощью break/return/исключения/throw, сообщение удаляется из очереди входящих сообщений так, как будто было возвращено истинное значение.
r = Ractor.new do
val = Ractor.receive_if{|msg| msg.is_a?(Array)}
puts "Received successfully: #{val}"
end
r.send(1)
r.send('test')
wait
puts "2 non-matching sent, nothing received"
r.send([1, 2, 3])
wait
Выводит:
2 non-matching sent, nothing received Received successfully: [1, 2, 3]
Обратите внимание, что вы не можете вызвать receive/receive_if рекурсивно в данном блоке. Это означает, что вы не должны выполнять какие-либо задачи в блоке.
Ractor.current << true
Ractor.receive_if{|msg| Ractor.receive}
#=> `receive': can not call receive/receive_if recursively (Ractor::Error)
# File ractor.rb, line 342
def self.select(*ractors, yield_value: yield_unspecified = true, move: false)
raise ArgumentError, 'specify at least one ractor or `yield_value`' if yield_unspecified && ractors.empty?
__builtin_cstmt! %q{
const VALUE *rs = RARRAY_CONST_PTR_TRANSIENT(ractors);
VALUE rv;
VALUE v = ractor_select(ec, rs, RARRAY_LENINT(ractors),
yield_unspecified == Qtrue ? Qundef : yield_value,
(bool)RTEST(move) ? true : false, &rv);
return rb_ary_new_from_args(2, rv, v);
}
end Ожидает, пока у первого рактора появится что-то в его исходящем порте, считывает данные от этого рактора и возвращает этот рактор и полученный объект.
r1 = Ractor.new {Ractor.yield 'from 1'}
r2 = Ractor.new {Ractor.yield 'from 2'}
r, obj = Ractor.select(r1, r2)
puts "received #{obj.inspect} from #{r.inspect}"
# Prints: received "from 1" from #<Ractor:#2 test.rb:1 running>
Если один из заданных ракторов — текущий рактор, и он будет выбран, r будет содержать :receive символ вместо объекта рактора.
r1 = Ractor.new(Ractor.current) do |main|
main.send 'to main'
Ractor.yield 'from 1'
end
r2 = Ractor.new do
Ractor.yield 'from 2'
end
r, obj = Ractor.select(r1, r2, Ractor.current)
puts "received #{obj.inspect} from #{r.inspect}"
# Prints: received "to main" from :receive
Если yield_value предоставлен, это значение может быть передано, если другой Ractor вызывает take. В этом случае пара [:yield, nil] будет возвращена:
r1 = Ractor.new(Ractor.current) do |main|
puts "Received from main: #{main.take}"
end
puts "Trying to select"
r, obj = Ractor.select(r1, Ractor.current, yield_value: 123)
wait
puts "Received #{obj.inspect} from #{r.inspect}"
Это выведет:
Trying to select Received from main: 123 Received nil from :yield
move логический флаг определяет, должно ли переданное значение копироваться (по умолчанию) или перемещаться.
Проверяет, является ли объект разделяемым ракторами.
Ractor.shareable?(1) #=> true -- numbers and other immutable basic values are frozen
Ractor.shareable?('foo') #=> false, unless the string is frozen due to # freeze_string_literals: true
Ractor.shareable?('foo'.freeze) #=> true
См. также раздел «Разделяемые и неразделяемые объекты» в документации класса Ractor.
# File ractor.rb, line 626
def self.yield(obj, move: false)
__builtin_cexpr! %q{
ractor_yield(ec, rb_ec_ractor_ptr(ec), obj, move)
}
end Отправляет сообщение текущему рактору в исходящий порт для обработки методом take.
r = Ractor.new {Ractor.yield 'Hello from ractor'}
puts r.take
# Prints: "Hello from ractor"
Метод блокируется и возвращается только тогда, когда кто-то получит отправленное сообщение.
r = Ractor.new do Ractor.yield 'Hello from ractor' puts "Ractor: after yield" end wait puts "Still not taken" puts r.take
Это выведет:
Still not taken Hello from ractor Ractor: after yield
Если исходящий порт был закрыт методом close_outgoing, метод вызывает исключение:
r = Ractor.new do close_outgoing Ractor.yield 'Hello from ractor' end wait # `yield': The outgoing-port is already closed (Ractor::ClosedError)
Значение аргумента move такое же, как и для send.
Публичные методы экземпляра
# File ractor.rb, line 823 def [](sym) Primitive.ractor_local_value(sym) end
Получить значение из локального хранилища рактора.
# File ractor.rb, line 828 def []=(sym, val) Primitive.ractor_local_value_set(sym, val) end
Установить значение в локальном хранилище рактора.
# File ractor.rb, line 733
def close_incoming
__builtin_cexpr! %q{
ractor_close_incoming(ec, RACTOR_PTR(self));
}
end Закрывает входящий порт и возвращает его предыдущее состояние. Все последующие попытки Ractor.receive в ракторе и send в рактор завершатся ошибкой Ractor::ClosedError.
r = Ractor.new {sleep(500)}
r.close_incoming #=> false
r.close_incoming #=> true
r.send('test')
# Ractor::ClosedError (The incoming-port is already closed)
# File ractor.rb, line 752
def close_outgoing
__builtin_cexpr! %q{
ractor_close_outgoing(ec, RACTOR_PTR(self));
}
end Закрывает исходящий порт и возвращает его предыдущее состояние. Все последующие попытки Ractor.yield в ракторе и take из рактора завершатся ошибкой Ractor::ClosedError.
r = Ractor.new {sleep(500)}
r.close_outgoing #=> false
r.close_outgoing #=> true
r.take
# Ractor::ClosedError (The outgoing-port is already closed)
# File ractor.rb, line 699
def inspect
loc = __builtin_cexpr! %q{ RACTOR_PTR(self)->loc }
name = __builtin_cexpr! %q{ RACTOR_PTR(self)->name }
id = __builtin_cexpr! %q{ INT2FIX(rb_ractor_id(RACTOR_PTR(self))) }
status = __builtin_cexpr! %q{
rb_str_new2(ractor_status_str(RACTOR_PTR(self)->status_))
}
"#<Ractor:##{id}#{name ? ' '+name : ''}#{loc ? " " + loc : ''} #{status}>"
end # File ractor.rb, line 712
def name
__builtin_cexpr! %q{RACTOR_PTR(self)->name}
end Имя, заданное в Ractor.new, или nil.
# File ractor.rb, line 582
def send(obj, move: false)
__builtin_cexpr! %q{
ractor_send(ec, RACTOR_PTR(self), obj, move)
}
end Отправка сообщения в очередь входящих сообщений рактора для обработки Ractor.receive.
r = Ractor.new do
value = Ractor.receive
puts "Received #{value}"
end
r.send 'message'
# Prints: "Received: message"
Метод неблокирующий (возвращает немедленно, даже если рактор не готов принять что-либо):
r = Ractor.new {sleep(5)}
r.send('test')
puts "Sent successfully"
# Prints: "Sent successfully" immediately
Попытка отправки в рактор, который уже завершил свою работу, вызовет Ractor::ClosedError.
r = Ractor.new {}
r.take
p r
# "#<Ractor:#6 (irb):23 terminated>"
r.send('test')
# Ractor::ClosedError (The incoming-port is already closed)
Если для рактора был вызван close_incoming, метод также вызовет Ractor::ClosedError.
r = Ractor.new do
sleep(500)
receive
end
r.close_incoming
r.send('test')
# Ractor::ClosedError (The incoming-port is already closed)
# The error would be raised immediately, not when ractor will try to receive
Если объект obj является несовместимым, по умолчанию он будет скопирован в рактор с помощью глубокого клонирования. Если передается move: true, объект переносится в рактор и становится недоступным для отправителя.
r = Ractor.new {puts "Received: #{receive}"}
msg = 'message'
r.send(msg, move: true)
r.take
p msg
Это выведет:
Received: message in `p': undefined method `inspect' for #<Ractor::MovedObject:0x000055c99b9b69b8>
Все ссылки на объект и его части станут недействительными для отправителя.
r = Ractor.new {puts "Received: #{receive}"}
s = 'message'
ary = [s]
copy = ary.dup
r.send(ary, move: true)
s.inspect
# Ractor::MovedError (can not send any methods to a moved object)
ary.class
# Ractor::MovedError (can not send any methods to a moved object)
copy.class
# => Array, it is different object
copy[0].inspect
# Ractor::MovedError (can not send any methods to a moved object)
# ...but its item was still a reference to `s`, which was moved
Если объект был совместимым, move: true не оказывает на него влияния:
r = Ractor.new {puts "Received: #{receive}"}
s = 'message'.freeze
r.send(s, move: true)
s.inspect #=> "message", still available
# File ractor.rb, line 693
def take
__builtin_cexpr! %q{
ractor_take(ec, RACTOR_PTR(self))
}
end Получить сообщение из исходящего порта рактора, которое было помещено туда с помощью Ractor.yield или во время завершения рактора.
r = Ractor.new do Ractor.yield 'explicit yield' 'last value' end puts r.take #=> 'explicit yield' puts r.take #=> 'last value' puts r.take # Ractor::ClosedError (The outgoing-port is already closed)
Тот факт, что последнее значение также помещается в исходящий порт, означает, что take может использоваться как аналог Thread#join («просто подождите, пока рактор завершит свою работу»), но не забудьте, что он вызовет ошибку, если кто-то уже израсходовал все, что произвел рактор.
Если исходящий порт был закрыт с помощью close_outgoing, метод вызовет Ractor::ClosedError.
r = Ractor.new do sleep(500) Ractor.yield 'Hello from ractor' end r.close_outgoing r.take # Ractor::ClosedError (The outgoing-port is already closed) # The error would be raised immediately, not when ractor will try to receive
Если в ракторе возникает необработанное исключение, оно передаётся при вызове take как Ractor::RemoteError.
r = Ractor.new {raise "Something weird happened"}
begin
r.take
rescue => e
p e # => #<Ractor::RemoteError: thrown by remote Ractor.>
p e.ractor == r # => true
p e.cause # => #<RuntimeError: Something weird happened>
end
Ractor::ClosedError является потомком StopIteration, поэтому закрытие рактора прервёт циклы без распространения ошибки:
r = Ractor.new do
3.times {|i| Ractor.yield "message #{i}"}
"finishing"
end
loop {puts "Received: " + r.take}
puts "Continue successfully"
Это выведет:
Received: message 0 Received: message 1 Received: message 2 Received: finishing Continue successfully
Приватные методы экземпляра
# File ractor.rb, line 426
def receive
__builtin_cexpr! %q{
ractor_receive(ec, rb_ec_ractor_ptr(ec))
}
end то же самое, что и Ractor.receive
# File ractor.rb, line 497
def receive_if &b
Primitive.ractor_receive_if b
end
Ruby Core © 1993–2022 Yukihiro Matsumoto
Licensed under the Ruby License.
Ruby Standard Library © contributors
Licensed under their own licenses.