Источник Task.Supervisor
Наблюдатель задач.
Этот модуль определяет наблюдателя, который может использоваться для динамического наблюдения за задачами.
Наблюдатель задач запускается без потомков, часто под управлением другого наблюдателя и с именем:
children = [
{Task.Supervisor, name: MyApp.TaskSupervisor}
]
Supervisor.start_link(children, strategy: :one_for_one)
Параметры, указанные в спецификации потомка, документированы в start_link/1.
После запуска вы можете запускать задачи непосредственно под наблюдателем, например:
task = Task.Supervisor.async(MyApp.TaskSupervisor, fn -> :do_some_work end)
См. модуль Task для получения дополнительных примеров.
Масштабируемость и разделение
Наблюдатель задач — это один процесс, ответственный за запуск других процессов. В некоторых приложениях наблюдатель задач может стать узким местом. Для решения этой проблемы вы можете запустить несколько экземпляров наблюдателя задач, а затем случайным образом выбрать один для запуска задачи.
Вместо:
children = [
{Task.Supervisor, name: MyApp.TaskSupervisor}
]
и:
Task.Supervisor.async(MyApp.TaskSupervisor, fn -> :do_some_work end)
Вы можете сделать так:
children = [
{PartitionSupervisor,
child_spec: Task.Supervisor,
name: MyApp.TaskSupervisors}
]
а затем:
Task.Supervisor.async(
{:via, PartitionSupervisor, {MyApp.TaskSupervisors, self()}},
fn -> :do_some_work end
)
В приведенном коде мы запускаем наблюдателя разделения, который по умолчанию запускает динамического наблюдателя для каждого ядра вашего компьютера. Затем, вместо вызова наблюдателя задач по имени, вы вызываете его через наблюдателя разделения, используя формат {:via, PartitionSupervisor, {name, key}}, где name — имя наблюдателя разделения, а key — ключ маршрутизации. Мы выбрали self() в качестве ключа маршрутизации, что означает, что каждый процесс будет назначен одному из существующих наблюдателей задач. Подробнее см. в документации PartitionSupervisor.
Регистрация имени
Наблюдатель задач подчиняется тем же правилам регистрации имен, что и GenServer. Подробнее см. в документации GenServer.
Краткое описание
Типы
- option()
Значения параметров, используемые
start_link
Функции
- async(supervisor, fun, options \\ [])
Запускает задачу, за которой можно наблюдать.
- async(supervisor, module, fun, args, options \\ [])
Запускает задачу, за которой можно наблюдать.
- async_nolink(supervisor, fun, options \\ [])
Запускает задачу, за которой можно наблюдать.
- async_nolink(supervisor, module, fun, args, options \\ [])
Запускает задачу, за которой можно наблюдать.
- async_stream(supervisor, enumerable, fun, options \\ [])
Возвращает поток, выполняющий заданную функцию
funпараллельно на каждом элементе вenumerable.- async_stream(supervisor, enumerable, module, function, args, options \\ [])
Возвращает поток, где заданная функция (
moduleиfunction) отображается параллельно на каждом элементе вenumerable.- async_stream_nolink(supervisor, enumerable, fun, options \\ [])
Возвращает поток, выполняющий заданную
functionпараллельно на каждом элементе вenumerable.- async_stream_nolink(supervisor, enumerable, module, function, args, options \\ [])
Возвращает поток, где заданная функция (
moduleиfunction) отображается параллельно на каждом элементе вenumerable.- children(supervisor)
Возвращает все идентификаторы потомков, кроме тех, которые перезапускаются.
- start_child(supervisor, fun, options \\ [])
Запускает задачу как потомка данного
supervisor.- start_child(supervisor, module, fun, args, options \\ [])
Запускает задачу как потомка данного
supervisor.- start_link(options \\ [])
Запускает новый наблюдатель.
- terminate_child(supervisor, pid)
Завершает потомка с заданным
pid.
Типы
option()Источник
@type option() :: DynamicSupervisor.option() | DynamicSupervisor.init_option()
Значения параметров, используемые start_link
Функции
async(supervisor, fun, options \\ [])Source
@spec async(Supervisor.supervisor(), (-> any()), Keyword.t()) :: Task.t()
Запускает задачу, на которой можно дождаться результата.
Аргумент supervisor должен быть ссылкой, как определено в Supervisor. Задача по-прежнему будет связана с вызывающим её процессом, см. Task.async/1 для получения дополнительной информации и async_nolink/3 для варианта без связи.
Вызывает ошибку, если supervisor достигла максимального числа дочерних задач.
Параметры
-
:shutdown-:brutal_killесли задачи необходимо убить при завершении работы или целое число, обозначающее значение таймаута, по умолчанию 5000 миллисекунд. Задачи должны перехватывать выхода для того, чтобы таймаут имел эффект.
async(supervisor, module, fun, args, options \\ [])Source
@spec async(Supervisor.supervisor(), module(), atom(), [term()], Keyword.t()) :: Task.t()
Запускает задачу, на которой можно дождаться результата.
Аргумент supervisor должен быть ссылкой, как определено в Supervisor. Задача по-прежнему будет связана с вызывающим её процессом, см. Task.async/1 для получения дополнительной информации и async_nolink/3 для варианта без связи.
Вызывает ошибку, если supervisor достигла максимального числа дочерних задач.
Параметры
-
:shutdown-:brutal_killесли задачи необходимо убить при завершении работы или целое число, обозначающее значение таймаута, по умолчанию 5000 миллисекунд. Задачи должны перехватывать выхода для того, чтобы таймаут имел эффект.
async_nolink(supervisor, fun, options \\ [])Source
@spec async_nolink(Supervisor.supervisor(), (-> any()), Keyword.t()) :: Task.t()
Запускает задачу, на которой можно дождаться результата.
Аргумент supervisor должен быть ссылкой, как определено в Supervisor. Задача не будет связана с вызывающим её процессом, см. Task.async/1 для получения дополнительной информации.
Вызывает ошибку, если supervisor достигла максимального числа дочерних задач.
Обратите внимание, что эта функция требует от надзорщика задач наличия :temporary в качестве параметра :restart (по умолчанию), так как async_nolink/3 сохраняет прямую ссылку на задачу, которая теряется, если задача перезапускается.
Параметры
-
:shutdown-:brutal_killесли задачи необходимо убить при завершении работы или целое число, обозначающее значение таймаута, по умолчанию 5000 миллисекунд. Задачи должны перехватывать выхода для того, чтобы таймаут имел эффект.
Совместимость с поведением OTP
Если вы создаете задачу с помощью async_nolink внутри поведения OTP, например, GenServer, вы должны обработать сообщение, пришедшее от задачи внутри вашего вызова GenServer.handle_info/2.
Ответ, отправленный задачей, будет в формате {ref, result}, где ref - ссылка на монитор, хранящийся в структуре задачи, а result - возвращаемое значение функции задачи.
Помните, что независимо от того, как завершается задача, созданная с помощью async_nolink, вызывающий процесс всегда получит сообщение :DOWN со значением ref, хранящимся в структуре задачи. Если задача завершается нормально, причина в сообщении :DOWN будет :normal.
Примеры
Как правило, вы используете async_nolink/3 в случае, когда есть разумное предположение, что задача может завершиться ошибкой, и вы не хотите, чтобы она приводила к завершению вызывающего процесса. Давайте рассмотрим пример, где GenServer должен запустить единственную задачу и отслеживать её статус:
defmodule MyApp.Server do
use GenServer
# ...
def start_task do
GenServer.call(__MODULE__, :start_task)
end
# In this case the task is already running, so we just return :ok.
def handle_call(:start_task, _from, %{ref: ref} = state) when is_reference(ref) do
{:reply, :ok, state}
end
# The task is not running yet, so let's start it.
def handle_call(:start_task, _from, %{ref: nil} = state) do
task =
Task.Supervisor.async_nolink(MyApp.TaskSupervisor, fn ->
...
end)
# We return :ok and the server will continue running
{:reply, :ok, %{state | ref: task.ref}}
end
# The task completed successfully
def handle_info({ref, answer}, %{ref: ref} = state) do
# We don't care about the DOWN message now, so let's demonitor and flush it
Process.demonitor(ref, [:flush])
# Do something with the result and then return
{:noreply, %{state | ref: nil}}
end
# The task failed
def handle_info({:DOWN, ref, :process, _pid, _reason}, %{ref: ref} = state) do
# Log and possibly restart the task...
{:noreply, %{state | ref: nil}}
end
end async_nolink(supervisor, module, fun, args, options \\ [])Source
@spec async_nolink(Supervisor.supervisor(), module(), atom(), [term()], Keyword.t()) :: Task.t()
Запускает задачу, на которой можно дождаться результата.
Аргумент supervisor должен быть ссылкой, как определено в Supervisor. Задача не будет связана с вызывающим её процессом, см. Task.async/1 для получения дополнительной информации.
Вызывает ошибку, если supervisor достигла максимального числа дочерних задач.
Обратите внимание, что эта функция требует от надзорщика задач наличия :temporary в качестве параметра :restart (по умолчанию), так как async_nolink/5 сохраняет прямую ссылку на задачу, которая теряется, если задача перезапускается.
async_stream(supervisor, enumerable, fun, options \\ [])Source
@spec async_stream( Supervisor.supervisor(), Enumerable.t(), (term() -> term()), keyword() ) :: Enumerable.t()
Возвращает поток, который запускает данную функцию fun асинхронно для каждого элемента в enumerable.
Каждый элемент в enumerable передаётся в качестве аргумента заданной функции fun и обрабатывается своей собственной задачей. Задачи будут созданы под управлением заданного supervisor и связаны с вызывающим процессом, аналогично async/3.
См. async_stream/6 для обсуждения, параметров и примеров.
async_stream(supervisor, enumerable, module, function, args, options \\ [])Source
@spec async_stream( Supervisor.supervisor(), Enumerable.t(), module(), atom(), [term()], keyword() ) :: Enumerable.t()
Возвращает поток, где заданная функция (module и function) применяется асинхронно к каждому элементу в enumerable.
Каждый элемент будет добавлен в начало указанного args и обработан своей собственной задачей. Задачи будут созданы под управлением заданного supervisor и связаны с вызывающим процессом, аналогично async/5.
При потоковой передаче каждая задача будет выводить {:ok, value} при успешном завершении или {:exit, reason} если вызывающий процесс перехватывает исключения. Порядок результатов зависит от значения параметра :ordered.
Уровень одновременности и время выполнения задач можно контролировать с помощью параметров (см. раздел "Параметры" ниже).
Если вам нужно перехватывать исключения для обработки исключений внутри асинхронного потока, рассмотрите использование async_stream_nolink/6 для запуска задач, не связанных с вызывающим процессом.
Параметры
:max_concurrency- устанавливает максимальное количество задач, которые могут выполняться одновременно. По умолчаниюSystem.schedulers_online/0.:ordered- определяет, должны ли результаты возвращаться в том же порядке, что и входной поток. Этот параметр полезен для больших потоков и не требует буферизации результатов перед их передачей. Также это полезно, когда вы используете задачи для побочных эффектов. По умолчаниюtrue.:timeout- максимальное время ожидания (в миллисекундах) без получения ответа от задачи (по всем выполняющимся задачам). По умолчанию5000.-
:on_timeout- действие при истечении времени ожидания задачи. Возможные значения:-
:exit(по умолчанию) - процесс, запустивший задачи, завершается. -
:kill_task- истекшая по времени задача уничтожается. Выводимое значение для этой задачи{:exit, :timeout}.
-
:zip_input_on_exit- (с версии v1.14.0) добавляет исходный вход в кортежи:exit. Выводимое значение для этой задачи{:exit, {input, reason}}, гдеinput- элемент коллекции, вызвавший выход во время обработки. По умолчаниюfalse.:shutdown-:brutal_killесли задачи необходимо убить при завершении работы или целое число, обозначающее значение таймаута. По умолчанию5000миллисекунд. Задачи должны перехватывать выхода для того, чтобы таймаут имел эффект.
Примеры
Давайте создадим поток и перечислим его:
stream = Task.Supervisor.async_stream(MySupervisor, collection, Mod, :expensive_fun, []) Enum.to_list(stream)
async_stream_nolink(supervisor, enumerable, fun, options \\ [])Source
@spec async_stream_nolink( Supervisor.supervisor(), Enumerable.t(), (term() -> term()), keyword() ) :: Enumerable.t()
Возвращает поток, который запускает данную function асинхронно для каждого элемента в enumerable.
Каждый элемент в enumerable передаётся в качестве аргумента заданной функции fun и обрабатывается своей собственной задачей. Задачи будут созданы под управлением заданного supervisor и не будут связаны с вызывающим процессом, аналогично async_nolink/3.
См. async_stream/6 для обсуждения и примеров.
async_stream_nolink(supervisor, enumerable, module, function, args, options \\ [])Source
@spec async_stream_nolink( Supervisor.supervisor(), Enumerable.t(), module(), atom(), [term()], keyword() ) :: Enumerable.t()
Возвращает поток, где заданная функция (module и function) применяется асинхронно к каждому элементу в enumerable.
Каждый элемент в enumerable будет добавлен в начало указанного args и обработан своей собственной задачей. Задачи будут созданы под управлением заданного supervisor и не будут связаны с вызывающим процессом, аналогично async_nolink/5.
См. async_stream/6 для обсуждения, параметров и примеров.
children(supervisor)Source
@spec children(Supervisor.supervisor()) :: [pid()]
Возвращает все идентификаторы процессов (PID) дочерних процессов, кроме тех, которые перезапускаются.
Обратите внимание, что вызов этой функции при наблюдении за большим количеством дочерних процессов в условиях низкой доступности памяти может привести к исключению из-за нехватки памяти.
start_child(supervisor, fun, options \\ [])Source
@spec start_child(Supervisor.supervisor(), (-> any()), keyword()) :: DynamicSupervisor.on_start_child()
Запускает задачу как дочерний процесс данного supervisor.
Task.Supervisor.start_child(MyTaskSupervisor, fn -> IO.puts "I am running in a task" end)
Обратите внимание, что запущенный процесс не связан с вызывающим процессом, а только с управляющей задачей. Эта команда полезна в случае, если задаче требуется выполнить побочные эффекты (например, ввод-вывод), и вам не интересны ее результаты или успешное завершение.
Параметры
:restart- стратегия перезапуска, может быть:temporary(по умолчанию),:transientили:permanent.:temporaryозначает, что задача никогда не перезапускается,:transientозначает, что она перезапускается, если выход не:normal,:shutdownили{:shutdown, reason}. Стратегия перезапуска:permanentозначает, что она всегда перезапускается.:shutdown-:brutal_killесли задачу необходимо убить непосредственно при завершении, или целое число, указывающее значение таймаута, по умолчанию 5000 миллисекунд. Задача должна перехватывать выходы, чтобы таймаут имел эффект.
start_child(supervisor, module, fun, args, options \\ [])Source
@spec start_child(Supervisor.supervisor(), module(), atom(), [term()], keyword()) :: DynamicSupervisor.on_start_child()
Запускает задачу как дочерний процесс данного supervisor.
Аналогично start_child/3, за исключением того, что задача задается указанным module, fun и args.
start_link(options \\ [])Source
@spec start_link([option()]) :: Supervisor.on_start()
Запускает новую управляющую задачу.
Примеры
Управляющая задача для задач обычно запускается в дереве управления с использованием кортежа:
{Task.Supervisor, name: MyApp.TaskSupervisor}
Вы также можете запустить ее, вызвав start_link/1 напрямую:
Task.Supervisor.start_link(name: MyApp.TaskSupervisor)
Но это рекомендуется только для сценариев и следует избегать в рабочем коде. Как правило, процессы всегда должны запускаться внутри деревьев управления.
Параметры
:name- используется для регистрации имени управляющей задачи, поддерживаемые значения описаны в разделеName Registrationв документации модуляGenServer;:max_restarts,:max_seconds, и:max_children- как указано вDynamicSupervisor;
Эта функция также может принимать :restart и :shutdown в качестве параметров, но эти два параметра устарели, и теперь предпочтительнее передавать их непосредственно в start_child.
terminate_child(supervisor, pid)Source
@spec terminate_child(Supervisor.supervisor(), pid()) :: :ok | {:error, :not_found} Прекращает выполнение дочернего процесса с заданным pid.
© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.16.3/Task.Supervisor.html