Исходный код 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() :: 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()), [ 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.17.2/Task.Supervisor.html