Spec-Zone.ru › Elixir 1.17

Исходный код Задача

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

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

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

Задачи — это процессы

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

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

В приведенном выше коде копируются все large_data, что может быть ресурсоёмким в зависимости от размера данных. Есть два способа решения этой проблемы.

Во-первых, если вам нужен доступ только к части large_data, рассмотрите возможность извлечения её перед запуском задачи:

large_data = fetch_large_data()
subset_data = large_data.some_field
task = Task.async(fn -> do_some_work(subset_data) end)

В качестве альтернативы, если вы можете перенести загрузку данных полностью в задачу, это может быть ещё лучше:

task = Task.async(fn ->
  large_data = fetch_large_data()
  do_some_work(large_data)
end)

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

Модуль 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, передавая имя надсмотрщика вместо PID:

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 легко динамически запускать задачи на разных узлах:

# First on the remote node named :remote@local
Task.Supervisor.start_link(name: MyApp.DistSupervisor)

# Then on the local client node
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 — PID, вызвавший текущий процесс, pid2 вызвал pid_n, а pid2 был вызван pid1.

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

Краткое описание

Типы

async_stream_option()

Параметры, передаваемые функциям async_stream.

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

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

Типы

async_stream_option()Source

@type async_stream_option() ::
  {:max_concurrency, pos_integer()}
  | {:ordered, boolean()}
  | {:timeout, timeout()}
  | {:on_timeout, :exit | :kill_task}
  | {:zip_input_on_exit, boolean()}

Параметры, передаваемые функциям async_stream.

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()), [async_stream_option()]) ::
  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()], [async_stream_option()]) ::
  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 - (с версии 1.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)

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

Внимание: не связанные асинхронные + 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

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

yield_many(tasks, opts \\ [])Source

@spec yield_many([t()], timeout()) :: [{t(), {:ok, term()} | {:exit, term()} | nil}]
@spec yield_many([t()],
  limit: pos_integer(),
  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, если вам нужно завершить процесс вызывающей программы при истечении времени ожидания.

Опции

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

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

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

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

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

Скачать версию ePub

Создано с помощью ExDoc (v0.34.1) для языка программирования Elixir

© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.17.2/Task.html

Spec-Zone.ru

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