Процессы и обмен сообщениями
Процессы и обмен сообщениями
Заголовок раздела «Процессы и обмен сообщениями»Создание изолированных процессов и взаимодействие через передачу сообщений.
Процессы предоставляют изолированные единицы выполнения, которые взаимодействуют через передачу сообщений. Каждый процесс имеет собственный почтовый ящик и может подписываться на конкретные темы сообщений.
Эта страница — введение: каждый сниппет демонстрирует один API изолированно. Полное рабочее приложение, связывающее порождение, мониторинг и обмен сообщениями воедино, см. в руководстве Echo Service.
Основные концепции:
- Создание процессов с помощью
process.spawn()и его вариантов - Отправка сообщений по PID или зарегистрированным именам через темы
- Получение сообщений через
process.listen()илиprocess.inbox() - Мониторинг жизненного цикла процессов через события
- Связывание процессов для координированной обработки отказов
Создание процессов
Заголовок раздела «Создание процессов»Создание нового процесса по ссылке на запись.
local pid, err = process.spawn("app.test.process:echo_worker", "app:processes", "hello")if err then return false, "spawn failed: " .. errend
-- pid is a string identifier for the spawned processprint("Started worker:", pid)Параметры:
- Ссылка на запись (например,
"app.test.process:echo_worker") - Ссылка на хост (например,
"app:processes") - Необязательные аргументы, передаваемые в функцию main воркера
Получение собственного PID
Заголовок раздела «Получение собственного PID»local my_pid = process.pid()-- Returns string PID of current processПередача сообщений
Заголовок раздела «Передача сообщений»Сообщения используют маршрутизацию на основе тем. Отправляйте сообщения на PID с указанием темы, затем получайте через подписку на тему или через почтовый ящик.
Отправка сообщений
Заголовок раздела «Отправка сообщений»-- Send to process by PIDlocal sent, err = process.send(worker_pid, "messages", "hello from parent")if err then return false, "send failed: " .. errend
-- send returns (bool, error)Получение через подписку на тему
Заголовок раздела «Получение через подписку на тему»Подписывайтесь на конкретные темы с помощью process.listen():
-- Worker that listens for messages on "messages" topiclocal function main() local ch = process.listen("messages")
local msg = ch:receive() if msg then -- msg is the payload directly print("Received:", msg) return true end
return falseend
return { main = main }Получение через почтовый ящик
Заголовок раздела «Получение через почтовый ящик»Почтовый ящик получает сообщения, не соответствующие ни одному слушателю тем:
local function main() local inbox_ch = process.inbox() local specific_ch = process.listen("specific_topic")
while true do local result = channel.select({ specific_ch:case_receive(), inbox_ch:case_receive() })
if result.channel == specific_ch then -- Messages to "specific_topic" arrive here local payload = result.value elseif result.channel == inbox_ch then -- Messages to any OTHER topic arrive here local msg = result.value print("Inbox got:", msg.topic, msg.payload) end endendРежим сообщений для информации об отправителе
Заголовок раздела «Режим сообщений для информации об отправителе»Используйте { message = true } для доступа к PID отправителя и теме:
-- Worker that echoes messages back to senderlocal function main() local ch = process.listen("echo", { message = true })
local msg = ch:receive() if msg then local sender = msg:from() local payload = msg:payload()
if sender then process.send(sender, "reply", payload) end return true end
return falseend
return { main = main }Мониторинг процессов
Заголовок раздела «Мониторинг процессов»Мониторинг процессов для получения событий EXIT при их завершении.
Создание с мониторингом
Заголовок раздела «Создание с мониторингом»local events_ch = process.events()
local worker_pid, err = process.spawn_monitored( "app.test.process:events_exit_worker", "app:processes")if err then return false, "spawn failed: " .. errend
-- Wait for EXIT eventlocal timeout = time.after("3s")local result = channel.select { events_ch:case_receive(), timeout:case_receive(),}
if result.channel == timeout then return false, "timeout waiting for EXIT event"end
local event = result.valueif event.kind == process.event.EXIT then print("Worker exited:", event.from) if event.result and event.result.error then print("Exit error:", event.result.error) elseif event.result then print("Return value:", event.result.value) endendЯвный мониторинг
Заголовок раздела «Явный мониторинг»Мониторинг уже запущенного процесса:
local events_ch = process.events()
-- Spawn without monitoringlocal worker_pid, err = process.spawn("app.test.process:long_worker", "app:processes")if err then return false, "spawn failed: " .. errend
-- Add monitoring explicitlylocal ok, monitor_err = process.monitor(worker_pid)if monitor_err then return false, "monitor failed: " .. monitor_errend
-- Now will receive EXIT events for this workerПрекращение мониторинга:
local ok, err = process.unmonitor(worker_pid)Связывание процессов
Заголовок раздела «Связывание процессов»Связывание процессов для координированного управления жизненным циклом. Связанные процессы получают события LINK_DOWN, когда связанный процесс завершается с ошибкой.
Создание связанного процесса
Заголовок раздела «Создание связанного процесса»-- Child terminates if parent crashes (unless trap_links is set)local pid, err = process.spawn_linked("app.test.process:child_worker", "app:processes")if err then return false, "spawn_linked failed: " .. errendЯвное связывание
Заголовок раздела «Явное связывание»-- Link to existing processlocal ok, err = process.link(target_pid)if err then return false, "link failed: " .. errend
-- Unlinklocal ok, err = process.unlink(target_pid)Обработка событий LINK_DOWN
Заголовок раздела «Обработка событий LINK_DOWN»По умолчанию LINK_DOWN приводит к завершению процесса с ошибкой. Включите trap_links для получения события вместо падения:
local function main() -- Enable trap_links to receive LINK_DOWN events instead of crashing local ok, err = process.set_options({ trap_links = true }) if not ok then return false, "set_options failed: " .. err end
-- Verify trap_links is enabled local opts = process.get_options() if not opts.trap_links then return false, "trap_links should be true" end
local events_ch = process.events()
-- Spawn a linked process that will fail local error_pid, err2 = process.spawn_linked( "app.test.process:error_exit_worker", "app:processes" ) if err2 then return false, "spawn error worker failed: " .. err2 end
-- Wait for LINK_DOWN event local timeout = time.after("2s") local result = channel.select { events_ch:case_receive(), timeout:case_receive(), }
if result.channel == timeout then return false, "timeout waiting for LINK_DOWN" end
local event = result.value if event.kind == process.event.LINK_DOWN then print("Linked process died:", event.from) -- Handle gracefully instead of crashing return true end
return false, "expected LINK_DOWN, got: " .. tostring(event.kind)end
return { main = main }Реестр процессов
Заголовок раздела «Реестр процессов»Регистрация имён для процессов для поиска и обмена сообщениями по имени.
Регистрация имён
Заголовок раздела «Регистрация имён»local function main() local test_name = "my_service_" .. tostring(os.time())
-- Register current process with a name local ok, err = process.registry.register(test_name) if err then return false, "register failed: " .. err end
-- Lookup the registered name local pid, lookup_err = process.registry.lookup(test_name) if lookup_err then return false, "lookup failed: " .. lookup_err end
-- Verify it resolves to our PID if pid ~= process.pid() then return false, "lookup returned wrong pid" end
return trueend
return { main = main }Снятие регистрации
Заголовок раздела «Снятие регистрации»-- Unregister explicitlylocal unregistered = process.registry.unregister(test_name)if not unregistered then print("Name was not registered")end
-- Lookup after unregister returns nil + errorlocal pid, err = process.registry.lookup(test_name)-- pid will be nil, err will be non-nilИмена автоматически освобождаются при завершении процесса.
Полный пример: пул воркеров с мониторингом
Заголовок раздела «Полный пример: пул воркеров с мониторингом»Этот пример показывает родительский процесс, создающий несколько воркеров с мониторингом и отслеживающий их завершение.
-- Parent processlocal time = require("time")
local function main() local events_ch = process.events()
-- Track spawned workers local workers = {} local worker_count = 5
-- Spawn multiple monitored workers for i = 1, worker_count do local worker_pid, err = process.spawn_monitored( "app.test.process:task_worker", "app:processes", { task_id = i, value = i * 10 } )
if err then return false, "spawn worker " .. i .. " failed: " .. err end
workers[worker_pid] = { task_id = i, started = os.time() } end
-- Wait for all workers to complete local completed = 0 local timeout = time.after("10s")
while completed < worker_count do local result = channel.select { events_ch:case_receive(), timeout:case_receive(), }
if result.channel == timeout then return false, "timeout waiting for workers" end
local event = result.value if event.kind == process.event.EXIT then local worker = workers[event.from] if worker then if event.result and event.result.error then print("Worker " .. worker.task_id .. " failed:", event.result.error) else print("Worker " .. worker.task_id .. " completed:", event.result and event.result.value) end completed = completed + 1 end end end
return trueend
return { main = main }Процесс-воркер:
-- task_worker.lualocal time = require("time")
local function main(task) -- Simulate work time.sleep("100ms")
-- Process task local result = task.value * 2
return resultend
return { main = main }Следующие шаги
Заголовок раздела «Следующие шаги»- Справочник модуля процессов - Полная документация API
- Каналы - Операции с каналами для обработки сообщений