Spec-Zone.ru › Elixir 1.7

Задача

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

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

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 function ... end}
])

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

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}
])

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

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

  • :id - идентификатор спецификации дочернего элемента, по умолчанию — текущий модуль
  • :start - как запустить дочерний процесс (по умолчанию вызывается __MODULE__.start_link/1)
  • :restart - когда следует перезапустить дочерний элемент, по умолчанию :temporary
  • :shutdown - как остановить дочерний элемент

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

use Task, restart: :transient

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

use Task, restart: :permanent

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

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

Модуль 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}
])

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

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{}

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

async(fun)

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

async(mod, fun, args)

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

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

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

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

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

await(task, timeout \\ 5000)

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

child_spec(arg)

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

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)

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

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

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

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

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

async(fun)

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

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

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

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

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

async(mod, fun, args)

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

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

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

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

Связывание

Эта функция запускает процесс, связанный и контролируемый вызывающим процессом. Важность связи заключается в том, что она прерывает задачу, если родительский процесс завершается аварийно. Она также гарантирует, что код до 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 \\ []) (с версии 1.4.0)

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

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

Каждый 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, args, options \\ []) (с версии 1.4.0)

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

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

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

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

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

Наконец, рассмотрите использование 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)

await(task, timeout \\ 5000)

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

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

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

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

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

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

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

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

Примеры

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

child_spec(arg) (с версии 1.5.0)

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

См. Supervisor.

shutdown(task, shutdown \\ 5000)

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

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

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

Второй аргумент — это либо таймаут, либо :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.7.4/Task.html

Spec-Zone.ru

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