Spec-Zone.ru › Elixir 1.15

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.

Типы

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)

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

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 достигло максимального числа дочерних элементов.

Обратите внимание, что для этой функции требуется, чтобы у надзирателя задач был параметр :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/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.15.4/Task.Supervisor.html

Spec-Zone.ru

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