Задача
Удобства для запуска и ожидания задач.
Задачи — это процессы, предназначенные для выполнения одного конкретного действия на протяжении всего своего жизненного цикла, часто с минимальным или нулевым взаимодействием с другими процессами. Наиболее распространенный случай использования задач — преобразование последовательного кода в конкуретнтный код путем асинхронного вычисления значения:
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 можно использовать для остановки задачи.
Наблюдаемые задачи
Также можно запустить задачу под управлением наблюдателя. Модуль 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- идентификатор спецификации дочернего элемента, по умолчанию совпадает с текущим модулем -
:start- способ запуска дочернего процесса (по умолчанию вызывается__MODULE__.start_link/1) -
: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, так как описанные там ограничения применяются ко всему экосистеме.
Резюме
Типы
- 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)
Ожидает ответа от нескольких задач в течение заданного интервала времени.
Типы
t()
t() :: %Task{owner: pid() | nil, pid: pid() | nil, ref: reference() | nil} Тип задачи.
См. %Task{} для получения информации о каждом поле структуры.
Функции
%Task{}
(структура)Структура задачи.
Она содержит следующие поля:
-
:pid- PID процесса задачи;nilесли задача не использует процесс задачи -
:ref- ссылка на монитор задачи -
:owner- PID процесса, который запустил задачу
async(fun)
async((() -> any())) :: t()
Запускает задачу, для которой требуется ожидание.
fun должна быть анонимной функцией без аргументов. Эта функция запускает процесс, связанный и контролируемый вызывающим процессом. Возвращается структура Task , содержащая соответствующую информацию.
Дополнительную информацию о общем использовании async/1 и async/3 см. в документации модуля Task.
См. также 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 \\ [])
(since 1.4.0)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 \\ [])
(since 1.4.0)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
Не рекомендуется await долго выполняющуюся задачу внутри поведения OTP, такого как GenServer. Вместо этого вы должны сопоставить сообщение, полученное от задачи, внутри вашего обратного вызова GenServer.handle_info/2. Для получения дополнительной информации о формате сообщения см. документацию по async/1.
Примеры
iex> task = Task.async(fn -> 1 + 1 end) iex> Task.await(task) 2
child_spec(arg)
(since 1.5.0)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.8.2/Task.html