Spec-Zone.ru › Elixir 1.3

Задача

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

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

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

Наблюдаемые задачи

Также можно запустить задачу под наблюдением с помощью start_link/1 и start_link/3:

Task.start_link(fn -> IO.puts "ok" end)

Эти задачи могут быть установлены в вашем дереве наблюдения следующим образом:

import Supervisor.Spec

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

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

По умолчанию большинство стратегий наблюдения будут пытаться перезапустить работника после его выхода независимо от причины. Если вы спроектировали задачу на нормальное завершение (как в примере с IO.puts/2 выше), рассмотрите возможность передачи restart: :transient в параметрах к 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])

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

Резюме

Типы

t()

Функции

%Task{}

Структура Task

async(fun)

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

async(mod, fun, args)

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

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{} (структура)

Структура Task.

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

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

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

  • :owner - идентификатор процесса, запустившего задачу

async(fun)

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

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

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

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

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

async(mod, fun, 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 так же, как вы бы это сделали, если бы не имели вызова async. Например, для возврата {:ok, val} | :error результатов или в более сложных случаях, используя try/rescue. Другими словами, асинхронная задача должна рассматриваться как расширение процесса, а не как механизм изоляции его от всех ошибок.

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

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

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

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

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

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

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

await(task, timeout \\ 5000)

await(t, timeout) :: term | no_return

Ожидание ответа задачи.

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

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

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

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

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

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

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 для завершения всех связанных процессов, включая задачи, которые не обрабатывают завершения, без генерации каких-либо сообщений в журнал.

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

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 находится в очереди сообщений. Если мониторинг был отключен или сообщение уже получено, функция будет ждать в течение тайм-аута, ожидая сообщения.

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 ->
      :timer.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.3.4/Task.html

Spec-Zone.ru

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