Spec-Zone.ru › Elixir 1.14

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: Task.Supervisor}
]

и:

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.

Краткое описание

Типы

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()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/3 для получения дополнительной информации и 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/3 для получения дополнительной информации и 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/3 для получения дополнительной информации.

Вызывает ошибку, если supervisor достигла максимального количества дочерних задач.

Параметры

  • :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/3 для получения дополнительной информации.

Вызывает ошибку, если 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}.
  • :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()]

Возвращает все идентификаторы процессов дочерних задач, кроме тех, которые перезапускаются.

Обратите внимание, что вызов этой функции при наблюдении за большим количеством дочерних задач в условиях низкой доступности памяти может вызвать исключение «недостаточно памяти».

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 Plataformatec
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.14.1/Task.Supervisor.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API