Spec-Zone.ru › Elixir 1.15

Задача

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

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

task = Task.async(fn -> do_some_work() end)
res = do_some_other_work()
res + Task.await(task)

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

По сравнению с обычными процессами, запущенными с помощью spawn/1, задачи включают мониторинг метаданных и ведение журнала в случае ошибок.

Помимо async/1 и await/2, задачи также могут быть запущены как часть дерева надзора и динамически запущены на удалённых узлах. Мы рассмотрим эти сценарии далее.

async и await

Одно из распространённых применений задач — преобразование последовательного кода в конкурентный код с помощью Task.async/1 при сохранении его семантики. При вызове будет создан новый процесс, связанный и контролируемый вызывающим процессом. После завершения действия задачи сообщение с результатом будет отправлено вызывающему процессу.

Task.await/2 используется для чтения сообщения, отправленного задачей.

Существует две важные вещи, которые следует учитывать при использовании async:

  1. Если вы используете асинхронные задачи, вы должны ожидать ответа, так как они всегда отправляются. Если вы не ожидаете ответа, рассмотрите использование Task.start_link/1, как описано ниже.

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

    Если этого не требуется, вы можете использовать задачи с надзором, описанные далее.

Динамически контролируемые задачи

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

И теперь вы можете использовать async/await, передавая имя надзирателя вместо идентификатора процесса:

Task.Supervisor.async(MyApp.TaskSupervisor, fn ->
  # Do something
end)
|> Task.await()

Мы рекомендуем разработчикам по возможности использовать контролируемые задачи. Контролируемые задачи улучшают видимость количества работающих задач в данный момент и позволяют использовать различные подходы, которые дают вам явный контроль над тем, как обрабатывать результаты, ошибки и таймауты. Вот краткое изложение:

  • Использование Task.Supervisor.start_child/2 позволяет запускать задачу «выстрелить и забыть», когда вас не интересуют её результаты или её успешное или неуспешное завершение.

  • Использование Task.Supervisor.async/2 + Task.await/2 позволяет выполнять задачи параллельно и получать их результат. Если задача завершается неудачно, вызывающий процесс также завершится неудачно.

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

Кроме того, надзиратель гарантирует завершение всех задач в течение настраиваемого периода завершения при завершении работы приложения. Подробную информацию о поддерживаемых операциях см. в модуле Task.Supervisor.

Распределённые задачи

С помощью Task.Supervisor легко динамически запускать задачи на разных узлах:

# On the remote node named :remote@local
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/5, которая ожидает явного модуля, функции и аргументов, вместо Task.Supervisor.async/3, которая работает с анонимными функциями. Это связано с тем, что анонимные функции ожидают, что версия того же модуля существует на всех участвующих узлах. Подробную информацию о распределённых процессах см. в документации модуля Agent, так как указанные там ограничения применяются ко всему экосистеме.

Статически контролируемые задачи

Модуль Task реализует функцию child_spec/1, которая позволяет запускать её непосредственно под обычным Supervisor - а не Task.Supervisor - передавая кортеж с функцией для выполнения:

Supervisor.start_link([
  {Task, fn -> :some_work end}
], strategy: :one_for_one)

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

Если вы не хотите помещать код задачи непосредственно под Supervisor, вы можете обернуть Task в свой собственный модуль, как вы делали бы с GenServer или Agent:

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)

Поскольку эти задачи контролируются и не связаны напрямую с вызывающим процессом, на них нельзя ожидать. По умолчанию функции Task.start/1 и Task.start_link/1 предназначены для задач «выстрелить и забыть», когда вас не интересуют результаты или её успешное или неуспешное завершение.

use Task

Когда вы use Task, модуль Task определит функцию child_spec/1, чтобы ваш модуль мог использоваться как дочерний элемент в дереве надзора.

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

  • :id - идентификатор спецификации дочернего элемента, по умолчанию совпадает с текущим модулем
  • :restart - когда дочерний элемент должен быть перезапущен, по умолчанию :temporary
  • :shutdown - как остановить дочерний элемент, немедленно или, дав ему время для завершения

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

use Task, restart: :transient

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

use Task, restart: :permanent

Дополнительную информацию см. в разделе «Спецификация дочернего элемента» модуля Supervisor. Аннотация @doc , которая непосредственно предшествует use Task , будет добавлена к сгенерированной функции child_spec/1.

Отслеживание предка и вызывающего процесса

Всякий раз, когда вы запускаете новый процесс, Elixir аннотирует родителя этого процесса ключом $ancestors в словаре процесса. Это часто используется для отслеживания иерархии внутри дерева надзора.

Например, мы рекомендуем разработчикам всегда запускать задачи под надзирателем. Это обеспечивает лучшую видимость и позволяет управлять тем, как эти задачи завершаются при выключении узла. Это может выглядеть примерно так Task.Supervisor.start_child(MySupervisor, task_function). Это означает, что, хотя ваш код и вызывает задачу, фактическим предком задачи является надзиратель, так как надзиратель фактически запускает её.

Для отслеживания взаимоотношений между вашим кодом и задачей используется ключ $callers в словаре процесса. Таким образом, предположим, что вызов Task.Supervisor выше, у нас есть:

[your code] -- calls --> [supervisor] ---- spawns --> [task]

Это означает, что мы сохраняем следующие отношения:

[your code]              [supervisor] <-- ancestor -- [task]
    ^                                                  |
    |--------------------- caller ---------------------|

Список вызывающих процессов текущего процесса можно получить из словаря процесса с помощью Process.get(:"$callers"). Это вернёт либо nil , либо список [pid_n, ..., pid2, pid1] с по крайней мере одной записью, где pid_n — это идентификатор процесса, который вызвал текущий процесс, pid2 вызвал pid_n, а pid2 был вызван pid1.

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

Типы

ref()

Непрозрачная ссылка на задачу.

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)

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

await_many(tasks, timeout \\ 5000)

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

child_spec(arg)

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

completed(result)

Запускает задачу, которая немедленно завершается с указанным result.

ignore(task)

Игнорирует существующую задачу.

shutdown(task, shutdown \\ 5000)

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

start(fun)

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

start(module, function_name, args)

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

start_link(fun)

Запускает задачу в составе дерева диспетчеризации с заданным fun.

start_link(module, function, args)

Запускает задачу в составе дерева диспетчеризации с заданным module, function, и args.

yield(task, timeout \\ 5000)

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

yield_many(tasks, opts \\ [])

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

ref()Source

@opaque ref()

Непрозрачная ссылка на задачу.

t()Source

@type t() :: %Task{mfa: mfa(), owner: pid(), pid: pid() | nil, ref: ref()}

Тип задачи.

См. %Task{} для получения информации о каждом поле структуры.

END_OF_DOCUMENT_MARKER

%Task{}Source

Структура Task.

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

  • :mfa - кортеж из трёх элементов, содержащий модуль, имя функции и арность, вызываемые для запуска задачи в async/1 и async/3

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

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

  • :ref - непрозрачный термин, используемый в качестве ссылки монитора задачи

async(fun)Source

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

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

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

Если вы запускаете async, вы обязаны дождаться её завершения. Это делается вызовом Task.await/2 или Task.yield/2, за которым следует Task.shutdown/2 для возвращённой задачи. В качестве альтернативы, если вы запускаете задачу внутри GenServer, то GenServer автоматически будет ждать завершения и вызывать GenServer.handle_info/2 с ответом задачи и связанным :DOWN сообщением.

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

Связывание

Эта функция порождает процесс, который связан и отслеживается вызывающим процессом. Важно, что эта связь прерывает задачу, если родительский процесс завершается. Также она гарантирует, что код до async/await будет иметь те же свойства и после добавления асинхронного вызова. Например, представьте:

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

Метаданные

Задача, созданная этой функцией, хранит :erlang.apply/2 в поле метаданных :mfa, которое используется внутри для применения анонимной функции. Используйте async/3, если вы хотите использовать другую функцию в качестве метаданных.

async(module, function_name, args)Source

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

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

Аналогично async/1, за исключением того, что функция, которую необходимо запустить, задаётся с помощью указанных module, function_name, и args. module, function_name, и его арность хранятся в поле :mfa для целей отражения.

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

@spec 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 \\ [])Source

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

Возвращает поток, где заданная функция (module и function_name) применяется асинхронно к каждому элементу в enumerable.

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

При потоковой передаче каждая задача будет выдавать {:ok, value} при успешном выполнении или {:exit, reason} если вызывающий процесс ловит исключения. Возможны {:exit, {element, reason}} исключений с использованием опции :zip_input_on_exit. Порядок результатов зависит от значения опции :ordered.

Уровень конкурентности и время выполнения задач можно контролировать с помощью опций (см. раздел "Опции" ниже).

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

Опции

  • :max_concurrency — устанавливает максимальное количество задач, выполняемых одновременно. По умолчанию System.schedulers_online/0.

  • :ordered — определяет, должны ли результаты возвращаться в том же порядке, что и входной поток. При упорядоченном выводе Elixir может потребоваться буферизовать результаты, чтобы выводить их в исходном порядке. Установка этого параметра в false отключает буферизацию, но также удаляет упорядоченность. Это также полезно, если вам нужны только побочные эффекты задач. Обратите внимание, что независимо от значения :ordered, задачи будут обрабатываться асинхронно. Если вам нужно обрабатывать элементы в порядке, рассмотрите использование Enum.map/2 или Enum.each/2 вместо этого. По умолчанию true.

  • :timeout — максимальное время (в миллисекундах или :infinity) выполнения каждой задачи. По умолчанию 5000.

  • :on_timeout — действия при истечении времени ожидания задачи. Возможные значения:

    • :exit (по умолчанию) — вызывающий процесс завершается.
    • :kill_task — задача, превысившая лимит времени, убивается. Выдаваемое значение для этой задачи {:exit, :timeout}.
  • :zip_input_on_exit — (с версии v1.14.0) добавляет исходный вход в кортежи :exit. Выдаваемое значение для задачи {:exit, {input, reason}}, где input элемент коллекции, вызвавший завершение обработки. По умолчанию false.

Пример

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

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)

Первые асинхронные задачи, которые завершаются

Вы также можете использовать async_stream/3 для выполнения M задач и поиска N завершенных задач. Например:

[
  &heavy_call_1/0,
  &heavy_call_2/0,
  &heavy_call_3/0
]
|> Task.async_stream(fn fun -> fun.() end, ordered: false, max_concurrency: 3)
|> Stream.filter(&match?({:ok, _}, &1))
|> Enum.take(2)

В приведенном примере мы выполняем три задачи и ждем завершения первых двух. Мы используем Stream.filter/2 для ограничения себя только успешно завершенными задачами и Enum.take/2 для получения N элементов. Важно установить как ordered: false, так и max_concurrency: M, где M — количество задач, чтобы убедиться, что все вызовы выполняются параллельно.

Внимание: неограниченное async + take

Если вы хотите потенциально обработать большое количество элементов и сохранить только часть результатов, вы можете обработать больше элементов, чем нужно. Давайте рассмотрим пример:

1..100
|> Task.async_stream(fn i ->
  Process.sleep(100)
  IO.puts(to_string(i))
end)
|> Enum.take(10)

Запуск примера на машине с 8 ядрами обработает 16 элементов, даже если вам нужны только 10, так как async_stream/3 обрабатывает элементы параллельно. Это происходит потому, что он обрабатывает 8 элементов одновременно. Затем все 8 элементов завершаются примерно в одно и то же время, что приводит к запуску обработки еще 8 элементов. Из этих дополнительных 8 элементов будут использованы только 2, а остальные будут завершены.

В зависимости от задачи вы можете отфильтровать или ограничить количество элементов заранее:

1..100
|> Stream.take(10)
|> Task.async_stream(fn i ->
  Process.sleep(100)
  IO.puts(to_string(i))
end)
|> Enum.to_list()

В других случаях, вероятно, нужно будет настроить :max_concurrency для ограничения количества элементов, которые могут быть обработаны сверх меры, за счет уменьшения конкурентности. Вы также можете установить количество элементов для взятия, кратное :max_concurrency. Например, установив max_concurrency: 5 в приведенном выше примере.

await(task, timeout \\ 5000)Source

@spec await(t(), timeout()) :: term()

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

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

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

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

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

Примеры

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

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

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

GenServer получит два сообщения по handle_info/2:

  • {ref, result} — сообщение ответа, где ref ссылка на монитор, возвращенная task.ref, а result результат задачи

  • {:DOWN, ref, :process, pid, reason} — поскольку все задачи также отслеживаются, вы также получите сообщение :DOWN, отправленное Process.monitor/1. Если вы получите сообщение :DOWN без ответа, это означает, что задача завершилась ошибкой

Также следует учитывать, что задачи, запущенные с помощью Task.async/1, всегда связаны с вызывающими процессами, и вам, возможно, не захочется, чтобы GenServer завершился, если задача завершится ошибкой. Поэтому предпочтительнее использовать Task.Supervisor.async_nolink/3 внутри поведения OTP. Для полноты, вот пример GenServer, который запускает задачи и обрабатывает их результаты:

defmodule GenServerTaskExample do
  use GenServer

  def start_link(opts) do
    GenServer.start_link(__MODULE__, :ok, opts)
  end

  def init(_opts) do
    # We will keep all running tasks in a map
    {:ok, %{tasks: %{}}}
  end

  # Imagine we invoke a task from the GenServer to access a URL...
  def handle_call(:some_message, _from, state) do
    url = ...
    task = Task.Supervisor.async_nolink(MyApp.TaskSupervisor, fn -> fetch_url(url) end)

    # After we start the task, we store its reference and the url it is fetching
    state = put_in(state.tasks[task.ref], url)

    {:reply, :ok, state}
  end

  # If the task succeeds...
  def handle_info({ref, result}, state) do
    # The task succeed so we can demonitor its reference
    Process.demonitor(ref, [:flush])

    {url, state} = pop_in(state.tasks[ref])
    IO.puts "Got #{inspect(result)} for URL #{inspect url}"
    {:noreply, state}
  end

  # If the task fails...
  def handle_info({:DOWN, ref, _, _, reason}, state) do
    {url, state} = pop_in(state.tasks[ref])
    IO.puts "URL #{inspect url} failed with reason #{inspect(reason)}"
    {:noreply, state}
  end
end

После определения сервера вы захотите запустить диспетчера задач выше и GenServer в вашей дереве надзора:

children = [
  {Task.Supervisor, name: MyApp.TaskSupervisor},
  {GenServerTaskExample, name: MyApp.GenServerTaskExample}
]

Supervisor.start_link(children, strategy: :one_for_one)

await_many(tasks, timeout \\ 5000)Source

@spec await_many([t()], timeout()) :: [term()]

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

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

Если какой-либо из процессов задач завершается аварийно, вызывающий процесс завершится по той же причине, что и эта задача.

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

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

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

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

Не рекомендуется await длительные задачи внутри поведения OTP, такого как GenServer. См. await/2 для получения дополнительной информации.

Примеры

iex> tasks = [
...>   Task.async(fn -> 1 + 1 end),
...>   Task.async(fn -> 2 + 3 end)
...> ]
iex> Task.await_many(tasks)
[2, 5]

child_spec(arg)Source

@spec child_spec(term()) :: Supervisor.child_spec()

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

arg передаётся как аргумент в Task.start_link/1 в :start поле спецификации.

Для получения дополнительной информации, см. модуль Supervisor, функцию Supervisor.child_spec/2 и тип Supervisor.child_spec/0.

completed(result)Source

@spec completed(any()) :: t()

Запускает задачу, которая сразу же завершается с указанным result.

В отличие от async/1, эта задача не порождает связанный процесс. Её можно ожидать или вызывать, как и любую другую задачу.

Использование

В некоторых случаях полезно создать задачу «completed», которая представляет задачу, которая уже выполнилась и сгенерировала результат. Например, при обработке данных вы можете определить, что некоторые входные данные неверны, прежде чем передавать их для дальнейшей обработки:

def process(data) do
  tasks =
    for entry <- data do
      if invalid_input?(entry) do
        Task.completed({:error, :invalid_input})
      else
        Task.async(fn -> further_process(entry) end)
      end
    end

  Task.await_many(tasks)
end

Во многих случаях Task.completed/1 можно избежать в пользу непосредственного возврата результата. Вам обычно потребуется этот вариант только при работе со смешанной асинхронностью, когда группа входов будет обрабатываться частично синхронно, а частично асинхронно.

ignore(task)Source

@spec ignore(t()) :: {:ok, term()} | {:exit, term()} | nil

Игнорирует существующую задачу.

Это означает, что задача продолжит выполняться, но она будет отсоединена, и вы больше не сможете её вызывать, ожидать или останавливать.

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

Важно: избегайте использования Task.async/1,3, а затем немедленного игнорирования задачи. Если вы хотите запустить задачи, результаты которых вас не интересуют, используйте Task.Supervisor.start_child/2 вместо этого.

shutdown(task, shutdown \\ 5000)Source

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

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

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

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

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

Если задача не связана с каким-либо процессом, например, задачи, запущенные с помощью Task.completed/1, мы проверяем наличие ответа или ошибки соответствующим образом, но без остановки процесса.

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

start(fun)Source

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

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

fun должна быть анонимной функцией без аргументов.

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

Если текущий узел завершается, узел завершится даже если задача не была завершена. По этой причине рекомендуется использовать Task.Supervisor.start_child/2 вместо этого, что позволяет контролировать время завершения через :shutdown параметр.

start(module, function_name, args)Source

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

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

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

Если текущий узел завершается, узел завершится даже если задача не была завершена. По этой причине рекомендуется использовать Task.Supervisor.start_child/2 вместо этого, что позволяет контролировать время завершения через :shutdown параметр.

start_link(fun)Source

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

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

fun должна быть анонимной функцией без аргументов.

Это используется для запуска статически управляемой задачи в дереве диспетчеризации.

start_link(module, function, args)Source

@spec start_link(module(), atom(), [term()]) :: {:ok, pid()}

Запускает задачу в рамках дерева диспетчеризации с заданными module, function, и args.

Это используется для запуска статически управляемой задачи в дереве диспетчеризации.

yield(task, timeout \\ 5000)Source

@spec yield(t(), timeout()) :: {:ok, term()} | {:exit, term()} | nil

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

Возвращает {:ok, reply} если ответ получен, nil если ответ не получен, или {:exit, reason} если задача завершилась. Обратите внимание, что обычно неудача задачи также приводит к выходу процесса, которому принадлежит задача. Поэтому эта функция может вернуть {:exit, reason} если выполняется хотя бы одно из следующих условий:

  • процесс задачи завершился по причине :normal
  • процесс задачи не связан с отправителем (задача была запущена с помощью Task.Supervisor.async_nolink/2 или Task.Supervisor.async_nolink/4)
  • отправитель обрабатывает завершения процессов

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

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

Если вы хотите завершить задачу, если она не ответит в течение timeout миллисекунд, вам нужно объединить эту функцию с shutdown/1, например:

case Task.yield(task, timeout) || Task.shutdown(task) do
  {:ok, result} ->
    result

  nil ->
    Logger.warning("Failed to get a result in #{timeout}ms")
    nil
end

Если вы хотите проверить задачу, но оставить её работающей после таймаута, вы можете объединить эту функцию с ignore/1, например:

case Task.yield(task, timeout) || Task.ignore(task) do
  {:ok, result} ->
    result

  nil ->
    Logger.warning("Failed to get a result in #{timeout}ms")
    nil
end

Это гарантирует, что если задача завершится после таймаута, но до вызова shutdown/1, вы всё равно получите результат, поскольку shutdown/1 разработан для обработки этого случая и возвращения результата.

yield_many(tasks, opts \\ [])Source

@spec yield_many([t()], timeout()) :: [{t(), {:ok, term()} | {:exit, term()} | nil}]
@spec yield_many([t()],
  timeout: timeout(),
  on_timeout: :nothing | :ignore | :kill_task
) :: [
  {t(), {:ok, term()} | {:exit, term()} | nil}
]

Вызывает множество задач в заданном интервале времени.

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

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

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

См. yield/2 для получения дополнительной информации.

Пример

Task.yield_many/2 позволяет разработчикам запускать несколько задач и получать результаты, полученные в заданный интервал времени. Если объединить его с Task.shutdown/2 (или Task.ignore/1), то можно собрать эти результаты и отменить (или игнорировать) задачи, которые не ответили вовремя.

Давайте посмотрим пример.

tasks =
  for i <- 1..10 do
    Task.async(fn ->
      Process.sleep(i * 1000)
      i
    end)
  end

tasks_with_results = Task.yield_many(tasks, timeout: 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.

Для удобства вы можете добиться аналогичного поведения, задав опцию :on_timeout как :kill_task (или :ignore). См. Task.await_many/2, если вы хотите выйти из вызывающего процесса при таймауте.

Опции

Вторым аргументом является либо таймаут, либо опции, которые по умолчанию:

  • :timeout - максимальное время (в миллисекундах или :infinity) выполнения каждой задачи. По умолчанию 5000.

  • :on_timeout - действие при таймауте задачи. Возможные значения:

    • :nothing - ничего не делать (по умолчанию). На задачи всё ещё можно ожидать, вызывать, игнорировать или завершать позже.
    • :ignore - результаты задачи будут проигнорированы.
    • :kill_task - задача, превысившая таймаут, будет убита.

© 2012 Plataformatec
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.15.4/Task.html

Spec-Zone.ru

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