Задача
Удобства для запуска и ожидания задач.
Задачи — это процессы, предназначенные для выполнения одного конкретного действия на протяжении всего их жизненного цикла, часто с минимальным или отсутствующим взаимодействием с другими процессами. Наиболее распространенный случай использования задач — преобразование последовательного кода в конкурентный код путем асинхронного вычисления значения:
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 необходимо учитывать две важные вещи:
-
Если вы используете асинхронные задачи, вы должны ожидать ответа, поскольку они всегда отправляются. Если вы не ожидаете ответа, используйте
Task.start_link/1, подробно описанную ниже. -
Асинхронные задачи связывают вызывающий процесс и запущенный процесс. Это означает, что если вызывающий процесс завершается аварийно, задача также завершится аварийно, и наоборот. Это сделано намеренно: если процесс, предназначенный для получения результата, больше не существует, нет смысла завершать вычисление.
Если этого не требуется, используйте
Task.start/1или рассмотрите возможность запуска задачи подTask.Supervisorс использованиемasync_nolinkилиstart_child.
Task.yield/2 — альтернатива await/2, где вызывающий процесс временно блокируется, ожидая ответа или сбоя задачи. Если результат не поступит в течение таймаута, его можно снова вызвать позже. Это позволяет проверять результат задачи несколько раз. Если ответ не поступит в течение желаемого времени, Task.shutdown/2 может быть использована для остановки задачи.
Задачи под надзором
Также можно запустить задачу под надзором. Это часто делается путем определения задачи в собственном модуле:
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])
Поскольку эти задачи находятся под надзором и не связаны напрямую с вызывающим процессом, к ним нельзя обратиться с ожиданием. Обратите внимание, что 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{} (структура)
Структура Task.
Она содержит следующие поля:
-
:pid— PID процесса задачи;nilесли задача не использует процесс задачи -
:ref— ссылка на монитор задачи -
:owner— PID процесса, который запустил задачу
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 так же, как вы бы это сделали, если бы асинхронного вызова не было. Например, для возвращения результатов {: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.
Каждый элемент 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 \\ [])
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)
Возвращает спецификацию для запуска задачи под контролем контролирующего процесса.
См. 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.6.6/Task.html