Исходный код 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 для получения дополнительных примеров.
Масштабируемость и разделение
Наблюдатель Task.Supervisor — это один процесс, ответственный за запуск других процессов. В некоторых приложениях наблюдатель Task.Supervisor может стать узким местом. Для решения этой проблемы можно запустить несколько экземпляров наблюдателя Task.Supervisor и затем случайным образом выбрать экземпляр для запуска задачи.
Вместо:
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
)
В приведенном коде мы запускаем наблюдателя разбиений, который по умолчанию запускает динамический наблюдатель для каждого ядра вашего компьютера. Затем, вместо вызова Task.Supervisor по имени, вы вызываете его через наблюдателя разбиений, используя формат {:via, PartitionSupervisor, {name, key}}, где name — имя наблюдателя разбиений, а key — ключ маршрутизации. Мы выбрали self() в качестве ключа маршрутизации, что означает, что каждый процесс будет назначен одному из существующих наблюдателей задач. Подробнее об этом читайте в документации PartitionSupervisor.
Регистрация имен
Наблюдатель Task.Supervisor подчиняется тем же правилам регистрации имен, что и GenServer. Подробнее об этом читайте в документации GenServer.
Резюме
Типы
- async_stream_option()
Опции, заданные для функций
async_streamиasync_stream_nolink.- 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.
Типы
async_stream_option()Source
@type async_stream_option() ::
Task.async_stream_option() | {:shutdown, Supervisor.shutdown()} Опции, заданные для функций async_stream и async_stream_nolink.
option()Source
@type option() :: GenServer.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()), [ async_stream_option() ]) :: 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()],
[
async_stream_option()
]
) :: 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()),
[
async_stream_option()
]
) :: Enumerable.t() Возвращает поток, который выполняет заданную function асинхронно на каждом элементе в enumerable.
Каждый элемент в enumerable передаётся в качестве аргумента заданной функции fun и обрабатывается своей задачей. Задачи будут созданы под заданным supervisor и не будут связаны с вызывающим процессом, подобно async_nolink/3.
См. async_stream/6 для обсуждения и примеров.
Обработка ошибок и очистка
Даже если задачи не связаны с вызывающим процессом, нет риска оставить зависшие задачи, работающие после остановки потока.
Рассмотрим следующий пример:
Task.Supervisor.async_stream_nolink(MySupervisor, collection, fun, on_timeout: :kill_task, ordered: false)
|> Enum.each(fn
{:ok, _} -> :ok
{:exit, reason} -> raise "Task exited: #{Exception.format_exit(reason)}"
end)
Если одна задача вызывает исключение или истекает время ожидания:
- выполняется второй блок кода
- вызывается исключение
- поток останавливается
- все активные задачи будут остановлены
Вот ещё один пример:
Task.Supervisor.async_stream_nolink(MySupervisor, collection, fun, on_timeout: :kill_task, ordered: false)
|> Stream.filter(&match?({:ok, _}, &1))
|> Enum.take(3)
Это вернёт три первые успешно выполненные задачи, игнорируя истечение времени ожидания и ошибки, и остановит все активные задачи.
Просто запустив поток с Stream.run/1, с другой стороны, ошибки будут проигнорированы, и весь поток будет обработан.
async_stream_nolink(supervisor, enumerable, module, function, args, options \\ [])Source
@spec async_stream_nolink(
Supervisor.supervisor(),
Enumerable.t(),
module(),
atom(),
[term()],
[
async_stream_option()
]
) :: Enumerable.t() Возвращает поток, где заданная функция (module и function) применяется асинхронно к каждому элементу в enumerable.
Каждый элемент в enumerable будет добавлен в качестве префикса к заданной args и обработан своей задачей. Задачи будут созданы под заданным supervisor и не будут связаны с вызывающим процессом, подобно async_nolink/5.
См. async_stream/6 для обсуждения, параметров и примеров.
children(supervisor)Source
@spec children(Supervisor.supervisor()) :: [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.18.1/Task.Supervisor.html