Spec-Zone.ru › Elixir 1.10

Задача

Удобства для запуска и ожидания задач.

Задачи — это процессы, предназначенные для выполнения одного определенного действия на протяжении всего их жизненного цикла, часто с минимальным или нулевым взаимодействием с другими процессами. Наиболее распространённый случай использования задач — преобразование последовательного кода в конкурентный код путём асинхронного вычисления значения:

task = Task.async(fn -> do_some_work() end)
res = do_some_other_work()
res + Task.await(task)

Запущенные задачи с async могут быть ожидаемы вызывающим процессом (и только им), как показано в примере выше. Они реализуются путём запуска процесса, который отправляет сообщение вызывающему процессу после выполнения заданного вычисления.

Помимо async/1 и await/2, задачи также могут быть запущены в рамках дерева управления и динамически запущены на удалённых узлах. Мы рассмотрим все три сценария далее.

async и await

Одно из распространённых применений задач — преобразование последовательного кода в конкурентный код с помощью Task.async/1, сохраняя при этом его семантику. При вызове будет создан новый процесс, связанный и контролируемый вызывающим процессом. После завершения действия задачи будет отправлено сообщение вызывающему процессу с результатом.

Task.await/2 используется для чтения сообщения, отправленного задачей.

При использовании async следует учитывать две важные вещи:

  1. Если вы используете асинхронные задачи, вы должны ожидать ответа, поскольку они всегда отправляются. Если вы не ожидаете ответа, рассмотрите использование Task.start_link/1, описанного ниже.

  2. Асинхронные задачи связывают вызывающий процесс и запущенный процесс. Это означает, что если вызывающий процесс завершается аварийно, задача также завершится аварийно, и наоборот. Это сделано намеренно: если процесс, который должен получить результат, больше не существует, нет смысла завершать вычисление.

    Если это нежелательно, используйте Task.start/1 или рассмотрите возможность запуска задачи под Task.Supervisor с помощью async_nolink или start_child.

Task.yield/2 — это альтернатива await/2, где вызывающий процесс временно заблокируется, ожидая ответа или аварийного завершения задачи. Если результат не придёт в течение установленного таймаута, его можно снова вызвать позже. Это позволяет проверять результат задачи несколько раз. Если ответ не приходит в течение заданного времени, Task.shutdown/2 может быть использована для остановки задачи.

Надзорные задачи

Также можно запустить задачу под управлением надзорщика. Модуль Task реализует функцию child_spec/1, которая позволяет запустить её непосредственно под надзорщиком, передав кортеж с функцией для выполнения:

Supervisor.start_link([
  {Task, fn -> :some_work end}
], strategy: :one_for_one)

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

defmodule MyTask do
  use Task

  def start_link(arg) do
    Task.start_link(__MODULE__, :run, [arg])
  end

  def run(arg) do
    # ...
  end
end

И затем передать его надзорщику:

Supervisor.start_link([
  {MyTask, arg}
], strategy: :one_for_one)

Поскольку эти задачи находятся под управлением надзорщика и не связаны напрямую с вызывающим процессом, к ним нельзя обратиться с помощью ожидания. start_link/1, в отличие от async/1, возвращает {:ok, pid} (который ожидается надзорщиками).

use Task определяет функцию child_spec/1, позволяющую поместить определённый модуль в дерево управления. Сгенерированную функцию child_spec/1 можно настроить с помощью следующих опций:

  • :id — идентификатор спецификации дочернего элемента, по умолчанию — текущий модуль
  • :restart — когда дочерний элемент должен быть перезапущен, по умолчанию — :temporary
  • :shutdown — как остановить дочерний элемент, либо сразу, либо дав ему время на завершение

В отличие от GenServer, Agent и Supervisor, задача имеет значение по умолчанию :restart :temporary. Это означает, что задача не будет перезапущена даже в случае аварийного завершения. Если вы хотите, чтобы задача перезапускалась при неудачном завершении, сделайте следующее:

use Task, restart: :transient

Если вы хотите, чтобы задача всегда перезапускалась:

use Task, restart: :permanent

См. раздел «Спецификация дочернего элемента» в модуле Supervisor для получения более подробной информации. Аннотация @doc непосредственно перед use Task будет прикреплена к сгенерированной функции child_spec/1.

Динамически контролируемые задачи

Модуль Task.Supervisor позволяет разработчикам динамически создавать несколько контролируемых задач.

Краткий пример:

{:ok, pid} = Task.Supervisor.start_link()

task =
  Task.Supervisor.async(pid, fn ->
    # Do something
  end)

Task.await(task)

Однако в большинстве случаев вы хотите добавить надзорщик задач в своё дерево управления:

Supervisor.start_link([
  {Task.Supervisor, name: MyApp.TaskSupervisor}
], strategy: :one_for_one)

Теперь вы можете динамически запускать контролируемые задачи:

Task.Supervisor.start_child(MyApp.TaskSupervisor, fn ->
  # Do something
end)

Или даже использовать шаблон async/await:

Task.Supervisor.async(MyApp.TaskSupervisor, fn ->
  # Do something
end)
|> Task.await()

Наконец, см. Task.Supervisor для других поддерживаемых операций.

Распределённые задачи

Так как Elixir предоставляет Task.Supervisor, легко использовать его для динамического запуска задач на разных узлах:

# On the remote node
Task.Supervisor.start_link(name: MyApp.DistSupervisor)

# On the client
supervisor = {MyApp.DistSupervisor, :remote@local}
Task.Supervisor.async(supervisor, MyMod, :my_fun, [arg1, arg2, arg3])

Обратите внимание, что при работе с распределёнными задачами следует использовать функцию Task.Supervisor.async/4, которая ожидает явных модуль, функцию и аргументы, а не Task.Supervisor.async/2, работающую с анонимными функциями. Это связано с тем, что анонимные функции ожидают, что версия того же модуля будет существовать на всех вовлечённых узлах. См. документацию модуля Agent для получения более подробной информации о распределённых процессах, поскольку описанные там ограничения применяются ко всему экосистеме.

Отслеживание предка и вызывающего процесса

Всякий раз, когда вы запускаете новый процесс, Elixir аннотирует родительский процесс с помощью ключа $ancestors в словаре процесса. Это часто используется для отслеживания иерархии в дереве управления.

Например, мы рекомендуем разработчикам всегда запускать задачи под управлением надзорщика. Это обеспечивает лучшую видимость и позволяет контролировать завершение этих задач при завершении работы узла. Это может выглядеть примерно так: Task.Supervisor.start_child(MySupervisor, task_specification). Это означает, что, хотя ваш код и вызывает задачу, фактическим предком задачи является надзорщик, поскольку надзорщик её фактически запускает.

Для отслеживания связи между вашим кодом и задачей используется ключ $callers в словаре процесса. Следовательно, предполагая вызов Task.Supervisor выше, мы имеем:

[your code] -- calls --> [supervisor] ---- spawns --> [task]

Что означает, что мы храним следующие отношения:

[your code]              [supervisor] <-- ancestor -- [task]
    ^                                                  |
    |--------------------- caller ---------------------|

Список вызывающих процессов текущего процесса можно получить из словаря процесса с помощью Process.get(:"$callers"). Это вернёт либо nil , либо список [pid_n, ..., pid2, pid1] с по крайней мере одной записью, где pid_n — PID, вызвавший текущий процесс, pid2 вызвал pid_n, и pid2 был вызван pid1.

Сводка

Типы

t()

Тип задачи.

Функции

%Task{}

Структура задачи.

async(fun)

Запускает задачу, которая должна быть ожидаемой.

async(module, function_name, args)

Запускает задачу, которая должна быть ожидаемой.

async_stream(enumerable, fun, options \\ [])

Возвращает поток, который выполняет данную функцию fun параллельно на каждом элементе в enumerable.

async_stream(enumerable, module, function_name, args, options \\ [])

Возвращает поток, где данная функция (module и function_name) отображается параллельно на каждом элементе в enumerable.

await(task, timeout \\ 5000)

Ожидает ответа задачи и возвращает его.

child_spec(arg)

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

shutdown(task, shutdown \\ 5000)

Отсоединяет и завершает задачу, а затем проверяет ответ.

start(fun)

Запускает задачу.

start(module, function_name, args)

Запускает задачу.

start_link(fun)

Запускает процесс, связанный с текущим процессом.

start_link(module, function_name, args)

Запускает задачу в рамках дерева управления.

yield(task, timeout \\ 5000)

Временно блокирует текущий процесс, ожидая ответа задачи.

yield_many(tasks, timeout \\ 5000)

Ожидает ответы от нескольких задач в заданном интервале времени.

END_OF_DOCUMENT_MARKER

Типы

t()

Характеристики

t() :: %Task{owner: pid() | nil, pid: pid() | nil, ref: reference() | nil}

Тип задачи.

См. %Task{} для информации о каждом поле структуры.

Функции

%Task{}

Структура Task.

Она содержит следующие поля:

  • :pid - идентификатор процесса задачи; nil если задача не использует процесс задачи

  • :ref - ссылка на монитор задачи

  • :owner - идентификатор процесса, который запустил задачу

async(fun)

Спецификации

async((() -> any())) :: t()

Запускает задачу, за которой необходимо ожидать.

fun должна быть анонимной функцией нулевой арности. Эта функция запускает процесс, связанный и отслеживаемый вызывающим процессом. Возвращается структура Task, содержащая соответствующую информацию.

Обратитесь к документации модуля Task для получения дополнительной информации об общем использовании async/1 и async/3.

См. также async/3.

async(module, function_name, args)

Спецификации

async(module(), atom(), [term()]) :: t()

Запускает задачу, за которой необходимо ожидать.

Возвращается структура Task, содержащая соответствующую информацию. Разработчики должны в конечном итоге вызвать Task.await/2 или Task.yield/2, а затем Task.shutdown/2 для возвращенной задачи.

Обратитесь к документации модуля Task для получения дополнительной информации об общем использовании async/1 и async/3.

Связывание

Эта функция запускает процесс, связанный и отслеживаемый вызывающим процессом. Часть связывания важна, потому что она прерывает задачу, если родительский процесс завершается. Она также гарантирует, что код перед async/await будет иметь те же свойства после добавления вызова async. Например, представьте, что у вас есть следующее:

x = heavy_fun()
y = some_fun()
x + y

Теперь вы хотите сделать heavy_fun() асинхронным:

x = Task.async(&heavy_fun/0)
y = some_fun()
Task.await(x) + y

Как и прежде, если heavy_fun/0 завершится неудачно, вся вычисление завершится неудачно, включая родительский процесс. Если вы не хотите, чтобы задача завершалась неудачно, вы должны изменить код heavy_fun/0 таким же образом, как вы бы это сделали, если бы у вас не было асинхронного вызова. Например, вернуть {:ok, val} | :error результаты или, в более крайних случаях, используя try/rescue. Другими словами, асинхронная задача должна рассматриваться как расширение процесса, а не как механизм для изоляции его от всех ошибок.

Если вы не хотите связывать вызывающий процесс с задачей, используйте управляемую задачу с Task.Supervisor и вызовите Task.Supervisor.async_nolink/2.

В любом случае избегайте следующего:

  • Установки :trap_exit в true — обработка завершений должна использоваться только в особых случаях, так как она сделает ваш процесс неуязвимым не только для завершений задачи, но и для любых других процессов.

    Кроме того, даже при обработке завершений вызов await все равно завершит процесс, если задача завершилась без возврата результата.

  • Отключение процесса задачи, запущенного с помощью async/await. Если вы отключаете процессы, а задача не принадлежит ни одному диспетчеру, вы можете оставить висящие задачи в случае завершения родительского процесса.

Формат сообщений

Ответ, отправленный задачей, будет иметь формат {ref, result}, где ref — ссылка на монитор, хранящаяся в структуре задачи, а result — возвращаемое значение функции задачи.

async_stream(enumerable, fun, options \\ [])

Спецификации

async_stream(Enumerable.t(), (term() -> term()), keyword()) :: Enumerable.t()

Возвращает поток, который выполняет заданную функцию fun параллельно для каждого элемента в enumerable.

Работает так же, как async_stream/5, но с анонимной функцией вместо кортежа модуль-функция-аргументы. fun должна быть анонимной функцией с одним аргументом.

Каждый элемент enumerable передаётся в качестве аргумента заданной функции fun и обрабатывается своей собственной задачей. Задачи будут связаны с текущим процессом, аналогично async/1.

Пример

Подсчитайте символы в каждой строке асинхронно, а затем суммируйте подсчеты с помощью reduce.

iex> strings = ["long string", "longer string", "there are many of these"]
iex> stream = Task.async_stream(strings, fn text -> text |> String.codepoints() |> Enum.count() end)
iex> Enum.reduce(stream, 0, fn {:ok, num}, acc -> num + acc end)
47

См. async_stream/5 для обсуждения, опций и дополнительных примеров.

async_stream(enumerable, module, function_name, args, options \\ [])

Спецификации

async_stream(Enumerable.t(), module(), atom(), [term()], keyword()) ::
  Enumerable.t()

Возвращает поток, где заданная функция (module и function_name) применяется параллельно к каждому элементу в enumerable.

Каждый элемент enumerable будет добавлен в заданную args и обработан своей собственной задачей. Задачи будут связаны с промежуточным процессом, который затем связан с текущим процессом. Это означает, что сбой в задаче завершит текущий процесс, а сбой в текущем процессе завершит все задачи.

При потоковой передаче каждая задача будет выдавать {:ok, value} при успешном выполнении или {:exit, reason} если вызывающий процесс обрабатывает завершения. Порядок результатов зависит от значения опции :ordered.

Уровень параллельности и время выполнения задач можно контролировать с помощью опций (см. раздел "Опции" ниже).

Рассмотрите возможность использования Task.Supervisor.async_stream/6 для запуска задач под управлением диспетчера. Если вы обнаружите, что обрабатываете завершения для обработки завершений внутри асинхронного потока, рассмотрите возможность использования Task.Supervisor.async_stream_nolink/6 для запуска задач, которые не связаны с вызывающим процессом.

Опции

  • :max_concurrency - устанавливает максимальное количество задач, которые могут выполняться одновременно. По умолчанию System.schedulers_online/0.

  • :ordered - определяет, должны ли результаты возвращаться в том же порядке, что и входной поток. Эта опция полезна при работе с большими потоками и позволяет избежать буферизации результатов перед их выдачей. Это также полезно, когда вы используете задачи для побочных эффектов. По умолчанию true.

  • :timeout - максимальное время (в миллисекундах), которое каждая задача может выполняться. По умолчанию 5000.

  • :on_timeout - действия при истечении времени выполнения задачи. Возможные значения:

    • :exit (по умолчанию) — процесс, запустивший задачи, завершается.
    • :kill_task — задача, которая истекла по времени, убивается. Выдаваемое значение для этой задачи — {:exit, :timeout}.

Пример

Давайте создадим поток и затем перечислим его:

stream = Task.async_stream(collection, Mod, :expensive_fun, [])
Enum.to_list(stream)

Параллельность можно увеличить или уменьшить с помощью опции :max_concurrency. Например, если задачи связаны с вводом-выводом, значение можно увеличить:

max_concurrency = System.schedulers_online() * 2
stream = Task.async_stream(collection, Mod, :expensive_fun, [], max_concurrency: max_concurrency)
Enum.to_list(stream)

Если вас не интересуют результаты вычисления, вы можете запустить поток с помощью Stream.run/1. Также установите ordered: false, так как вас не интересует порядок результатов:

stream = Task.async_stream(collection, Mod, :expensive_fun, [], ordered: false)
Stream.run(stream)

await(task, timeout \\ 5000)

Спецификации

await(t(), timeout()) :: term()

Ожидает ответа от задачи и возвращает его.

В случае завершения процесса задачи текущий процесс завершится по той же причине, что и задача.

Таймаут в миллисекундах или :infinity, может быть задан со значением по умолчанию 5000. Если таймаут истекает, текущий процесс завершится. Если процесс задачи связан с текущим процессом, как это происходит при запуске задачи с async, то процесс задачи также завершится. Если процесс задачи обрабатывает завершения или не связан с текущим процессом, то он будет продолжать выполняться.

Эта функция предполагает, что монитор задачи всё ещё активен или сообщение монитора :DOWN находится в очереди сообщений. Если он был отслежен или сообщение уже получено, эта функция будет ждать в течение таймаута, ожидая сообщения.

Эта функция может быть вызвана только один раз для любой заданной задачи. Если вы хотите проверять, завершилась ли длительно выполняющаяся задача, используйте yield/2 вместо этого.

Совместимость с поведением OTP

Не рекомендуется запускать длительно выполняющуюся задачу внутри поведения OTP, такого как GenServer. Вместо этого вы должны сопоставить сообщение, поступающее от задачи, внутри своего обратного вызова GenServer.handle_info/2. Дополнительную информацию о формате сообщения см. в документации для async/1.

Примеры

iex> task = Task.async(fn -> 1 + 1 end)
iex> Task.await(task)
2

child_spec(arg)

Спецификации

child_spec(term()) :: Supervisor.child_spec()

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

arg передаётся в качестве аргумента в Task.start_link/1 в поле :start спецификации.

Дополнительную информацию см. в модуле Supervisor, функции Supervisor.child_spec/2 и типе Supervisor.child_spec/0.

shutdown(task, shutdown \\ 5000)

Характеристики

shutdown(t(), timeout() | :brutal_kill) :: {:ok, term()} | {:exit, term()} | nil

Отсоединяет и завершает задачу, а затем проверяет наличие ответа.

Возвращает {:ok, reply} если ответ получен во время завершения задачи, {:exit, reason} если задача завершилась аномально, в противном случае nil.

Второй аргумент — это либо таймаут, либо :brutal_kill. В случае таймаута, сигнала :shutdown отправляется процессу задачи, и если он не завершается в течение таймаута, он убивается. При использовании :brutal_kill, задача убивается сразу. В случае ненормального завершения задачи (возможно, убитой другим процессом), эта функция завершится по той же причине.

Вызывать эту функцию при завершении вызывающего процесса не требуется, за исключением выхода по причине :normal или если задача перехватывает выходы. Если вызывающий процесс завершается по причине, отличной от :normal, и задача не перехватывает выходы, сигнал завершения вызывающего процесса остановит задачу. Вызывающий процесс может завершиться с причиной :shutdown, чтобы завершить все связанные процессы, включая задачи, которые не перехватывают выходы, без генерации сообщений об ошибках.

Если монитор задачи уже демонизирован или получил ответ, и в очереди сообщений нет ожидающего ответа, функция вернёт {:exit, :noproc} в качестве результата, или причина завершения не может быть определена.

start(fun)

Характеристики

start((() -> any())) :: {:ok, pid()}

Запускает задачу.

fun должна быть анонимной функцией без аргументов.

Это используется только в случае, если задача используется для побочных эффектов (т.е. результат не интересен) и она не должна быть связана с текущим процессом.

start(module, function_name, args)

Характеристики

start(module(), atom(), [term()]) :: {:ok, pid()}

Запускает задачу.

Это используется только в случае, если задача используется для побочных эффектов (т.е. результат не интересен) и она не должна быть связана с текущим процессом.

start_link(fun)

Характеристики

start_link((() -> any())) :: {:ok, pid()}

Запускает процесс, связанный с текущим процессом.

fun должна быть анонимной функцией без аргументов.

Это часто используется для запуска процесса как части дерева контроля.

start_link(module, function_name, args)

Характеристики

start_link(module(), atom(), [term()]) :: {:ok, pid()}

Запускает задачу как часть дерева контроля.

yield(task, timeout \\ 5000)

Характеристики

yield(t(), timeout()) :: {:ok, term()} | {:exit, term()} | nil

Временно блокирует текущий процесс, ожидая ответа от задачи.

Возвращает {:ok, reply} если ответ получен, nil если ответ не получен, или {:exit, reason} если задача уже завершилась. Имейте в виду, что обычно сбой задачи также приводит к завершению процесса, владеющего задачей. Таким образом, эта функция может вернуть {:exit, reason} только если:

  • задача завершилась с причиной :normal
  • она не связана с вызывающим процессом
  • вызывающий процесс перехватывает выходы

Таймаут в миллисекундах или :infinity, может быть задан со значением по умолчанию 5000. Если время истечёт прежде, чем будет получено сообщение от задачи, эта функция вернёт nil, и монитор останется активным. Поэтому yield/2 может быть вызвано несколько раз для той же задачи.

Эта функция предполагает, что монитор задачи по-прежнему активен или сообщение монитора :DOWN находится в очереди сообщений. Если он был демонизирован или сообщение уже получено, эта функция будет ждать в течение таймаута, ожидая сообщения.

Если вы хотите завершить задачу, если она не ответила в течение timeout миллисекунд, вы должны связать это с shutdown/1, как показано ниже:

case Task.yield(task, timeout) || Task.shutdown(task) do
  {:ok, result} ->
    result

  nil ->
    Logger.warn("Failed to get a result in #{timeout}ms")
    nil
end

Это гарантирует, что если задача завершится после истечения timeout но до вызова shutdown/1, вы всё равно получите результат, поскольку shutdown/1 разработан для обработки этого случая и возврата результата.

yield_many(tasks, timeout \\ 5000)

Характеристики

yield_many([t()], timeout()) :: [{t(), {:ok, term()} | {:exit, term()} | nil}]

Получение результатов от нескольких задач в заданном интервале времени.

Эта функция получает список задач и ожидает их ответы в заданном интервале времени. Она возвращает список пар из двух элементов, где первая — это задача, а вторая — полученный результат. Задачи в возвращаемом списке будут в том же порядке, что и в входном аргументе tasks.

Аналогично yield/2, результат каждой задачи будет

  • {:ok, term} если задача успешно сообщила свой результат в заданном интервале времени
  • {:exit, reason} если задача завершилась аномально
  • nil если задача продолжала выполняться после таймаута

Таймаут в миллисекундах или :infinity, может быть задан со значением по умолчанию 5000.

Дополнительную информацию см. в yield/2.

Пример

Task.yield_many/2 позволяет разработчикам запускать несколько задач и получать результаты, полученные в заданном временном интервале. Если объединить его с Task.shutdown/2, это позволит собрать эти результаты и отменить задачи, которые не ответили вовремя.

Рассмотрим пример.

tasks =
  for i <- 1..10 do
    Task.async(fn ->
      Process.sleep(i * 1000)
      i
    end)
  end

tasks_with_results = Task.yield_many(tasks, 5000)

results =
  Enum.map(tasks_with_results, fn {task, res} ->
    # Shut down the tasks that did not reply nor exit
    res || Task.shutdown(task, :brutal_kill)
  end)

# Here we are matching only on {:ok, value} and
# ignoring {:exit, _} (crashed tasks) and `nil` (no replies)
for {:ok, value} <- results do
  IO.inspect(value)
end

В примере выше создаются задачи, которые спят от 1 до 10 секунд и возвращают число секунд, которые они проспали. Если выполнить код сразу, вы увидите числа от 1 до 5, так как это те задачи, которые ответили в заданное время. Все остальные задачи будут завершены с помощью вызова Task.shutdown/2.

© 2012 Plataformatec
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.10.4/Task.html

Spec-Zone.ru

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