Spec-Zone.ru › Elixir 1.4

Задача

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

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

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 можно использовать для остановки задачи.

Задачи под надзором

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

import Supervisor.Spec

children = [
  #
  worker(Task, [fn -> IO.puts "ok" end])
]

Внутренне надзорник вызовет Task.start_link/1.

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

По умолчанию большинство стратегий надзора пытаются перезапустить работника после его завершения независимо от причины. Если вы разработали задачу для нормального завершения (как в примере с IO.puts/2 выше), рассмотрите передачу restart: :transient в опциях Supervisor.Spec.worker/3.

Динамически управляемые задачи

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

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

{:ok, pid} = Task.Supervisor.start_link()
task = Task.Supervisor.async(pid, fn ->
  # Do something
end)
Task.await(task)

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

import Supervisor.Spec

children = [
  supervisor(Task.Supervisor, [[name: MyApp.TaskSupervisor]])
]

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

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 предоставляет надзорник задач, легко использовать его для динамического запуска задач на разных узлах:

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

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

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

Сводка

Типы

t()

Функции

%Task{}

Структура Task

async(fun)

Запускает задачу, к которой необходимо обратиться с помощью await

async(mod, fun, args)

Запускает задачу, к которой необходимо обратиться с помощью await

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

Возвращает поток, который выполняет данную function асинхронно для каждого элемента в enumerable

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

Возвращает поток, который выполняет данную module, function и args асинхронно для каждого элемента в enumerable

await(task, timeout \\ 5000)

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

shutdown(task, shutdown \\ 5000)

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

start(fun)

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

start(mod, fun, args)

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

start_link(fun)

Запускает задачу в рамках дерева надзора

start_link(mod, fun, args)

Запускает задачу в рамках дерева надзора

yield(task, timeout \\ 5000)

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

yield_many(tasks, timeout \\ 5000)

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

Типы

t()

t() :: %Task{owner: term(), pid: term(), ref: term()}

Функции

%Task{} (struct)

Структура Task.

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

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

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

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

async(fun)

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

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

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

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

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

async(mod, fun, args)

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

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

Возвращается структура 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.t()) :: Enumerable.t()

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

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

См. async_stream/5 для обсуждения и примеров.

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

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

Возвращает поток, который выполняет данную module, function и args одновременно для каждого элемента в enumerable.

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

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

Уровень параллелизма может быть контролируется параметром :max_concurrency и по умолчанию равен System.schedulers_online/0. Таймаут также может быть задан как параметр и по умолчанию равен 5000 и по умолчанию равен максимальному времени ожидания без ответа задачи.

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

Параметры

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

Пример

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

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)

await(task, timeout \\ 5000)

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

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

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

Если таймаут истечёт, await выйдет; однако задача будет продолжать выполняться. При выходе вызывающего процесса его сигнал выхода завершит задачу, если она не перехватывает выходы.

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

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

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

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

Примеры

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

shutdown(task, shutdown \\ 5000)

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

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

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

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

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

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

start(fun)

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

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

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

start(mod, fun, args)

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

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

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

start_link(fun)

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

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

start_link(mod, fun, 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
  • он не связан с вызывающим процессом
  • вызывающий процесс перехватывает выходы

Таймаут в миллисекундах может быть задан, по умолчанию он равен 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}]

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

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

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

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

См. 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} ->
  # Shutdown 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.4.5/Task.html

Spec-Zone.ru

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