Spec-Zone.ru › Elixir 1.14

Задача

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

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

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

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

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

async и await

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

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

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

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

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

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

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

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

# 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 определяет функцию 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.

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

Типы

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, timeout \\ 5000)

Уступает право множеству задач в заданный интервал времени.

Типы

t()Source

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

Тип Задачи.

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

Функции

%Task{}Source

Структура Задачи.

Содержит следующие поля:

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

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

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

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

async(fun)Source

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

Запускает задачу, требующую ожидания.

fun должна быть анонимной функцией нулевой арности. Эта функция запускает процесс, связанный и отслеживаемый процессом вызывающего кода. Возвращается структура Task, содержащая соответствующую информацию. Разработчики должны в конечном итоге вызвать Task.await/2 или Task.yield/2, а затем Task.shutdown/2 для возвращенной задачи.

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

Связывание

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

x = heavy_fun()
y = some_fun()
x + y

Теперь вы хотите сделать heavy_fun() асинхронным:

x = Task.async(&heavy_fun/0)
y = some_fun()
Task.await(x) + y

Как и прежде, если heavy_fun/0 завершается с ошибкой, вся вычисления завершатся ошибкой, включая вызывающий процесс. Если вы не хотите, чтобы задача завершилась с ошибкой, вы должны изменить код heavy_fun/0 таким же образом, как вы бы это сделали, если бы у вас не было вызова async. Например, чтобы вернуть {:ok, val} | :error результаты или, в более крайних случаях, с помощью try/rescue. Иными словами, асинхронная задача должна рассматриваться как расширение вызывающего процесса, а не как механизм его изоляции от всех ошибок.

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

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

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

    Более того, даже при перехвате выходов, вызов await всё ещё завершится ошибкой, если задача завершилась без возвращения результата.

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

Метаданные

Задача, созданная с помощью этой функции, хранит :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} если вызывающий процесс перехватывает завершения. Возможно настроить перехват завершений с помощью опции :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)

В примере выше мы выполняем три задачи и ждем завершения первых 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 cancel the monitoring and discard the DOWN message
    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

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

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

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

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

Требуется Erlang/OTP 24+.

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

Если монитор задачи уже был демонизирован или получил сообщение, и в очереди сообщений нет ожидающего ответа, эта функция вернёт {: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.warn("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.warn("Failed to get a result in #{timeout}ms")
    nil
end

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

yield_many(tasks, timeout \\ 5000)Source

@spec 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 (или 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, 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.14.1/Task.html

Spec-Zone.ru

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