Spec-Zone.ru › Elixir 1.9

Задача

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

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

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.start/1 или рассмотрите запуск задачи под Task.Supervisor с помощью async_nolink или start_child.

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

Надзорные задачи

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

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

Однако, если вы хотите вызвать определенный модуль, функцию и аргументы или дать процессу задачи имя, вам необходимо определить задачу в собственном модуле:

defmodule MyTask do
  use Task

  def start_link(arg) do
    Task.start_link(__MODULE__, :run, [arg])
  end

  def run(arg) do
    # ...
  end
end

А затем передать его надзорщику:

Supervisor.start_link([
  {MyTask, arg}
], strategy: :one_for_one)

Поскольку эти задачи находятся под надзором и не связаны напрямую с вызывающим процессом, они не могут быть ожиданиями. start_link/1, в отличие от async/1, возвращает {:ok, pid} (который является результатом, ожидаемым надзорщиками).

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

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

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

use Task, restart: :transient

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

use Task, restart: :permanent

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

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

Модуль Task.Supervisor позволяет разработчикам динамически создавать несколько контролируемых задач.

Краткое пример:

{:ok, pid} = Task.Supervisor.start_link()

task =
  Task.Supervisor.async(pid, fn ->
    # Do something
  end)

Task.await(task)

Однако в большинстве случаев вы хотите добавить надзорщик задач в свое дерево надзора:

Supervisor.start_link([
  {Task.Supervisor, name: MyApp.TaskSupervisor}
], strategy: :one_for_one)

Теперь вы можете динамически запускать контролируемые задачи:

Task.Supervisor.start_child(MyApp.TaskSupervisor, fn ->
  # Do something
end)

Или даже использовать шаблон async/await:

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

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

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

Поскольку Elixir предоставляет Task.Supervisor, легко использовать его для динамического запуска задач на разных узлах:

# On the remote node
Task.Supervisor.start_link(name: MyApp.DistSupervisor)

# On the client
supervisor = {MyApp.DistSupervisor, :remote@local}
Task.Supervisor.async(supervisor, MyMod, :my_fun, [arg1, arg2, arg3])

Обратите внимание, что при работе с распределенными задачами следует использовать функцию Task.Supervisor.async/4, которая ожидает явных модуль, функцию и аргументы, вместо Task.Supervisor.async/2, которая работает с анонимными функциями. Это связано с тем, что анонимные функции ожидают, что версия того же модуля будет существовать на всех участвующих узлах. Обратитесь к документации модуля Agent для получения дополнительной информации о распределенных процессах, так как ограничения, описанные там, применяются ко всему экосистеме.

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

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

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

Для отслеживания связи между вашим кодом и задачей мы используем ключ $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.

Резюме

Типы

t()

Тип задачи.

Функции

%Task{}

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

async(fun)

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

async(module, function_name, args)

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

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

Возвращает поток, выполняющий заданную функцию fun параллельно для каждого элемента в enumerable.

async_stream(enumerable, module, function_name, args, options \\ [])

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

await(task, timeout \\ 5000)

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

child_spec(arg)

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

shutdown(task, shutdown \\ 5000)

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

start(fun)

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

start(module, function_name, args)

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

start_link(fun)

Запускает процесс, связанный с текущим процессом.

start_link(module, function_name, args)

Запускает задачу как часть дерева надзора.

yield(task, timeout \\ 5000)

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

yield_many(tasks, timeout \\ 5000)

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

Типы

t()

Спецификации

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

Тип задачи.

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

Функции

%Task{}

Структура Task.

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

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

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

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

async(fun)

Характеристики

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

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

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

Дополнительную информацию о общем использовании async/1 и async/3 см. в документации модуля Task.

См. также async/3.

async(module, function_name, args)

Характеристики

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

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

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

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

Формат сообщения

Ответ, отправленный задачей, будет иметь формат {ref, result}, где ref - ссылка на монитор, хранящаяся в структуре задачи, а result - возвращаемое значение функции задачи.

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

Характеристики

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

Характеристики

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

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

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

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

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

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

Параметры

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

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

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

  • :on_timeout - что делать, когда задача истекает. Возможные значения:

    • :exit (по умолчанию) - процесс, который запустил задачи, завершается.
    • :kill_task - задача, которая истекла, уничтожается. Возвращаемое значение для этой задачи - {:exit, :timeout}.

Пример

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

stream = Task.async_stream(collection, Mod, :expensive_fun, [])
Enum.to_list(stream)

Уровень параллельности можно увеличить или уменьшить, используя опцию :max_concurrency. Например, если задачи связаны с вводом-выводом, значение можно увеличить:

max_concurrency = System.schedulers_online() * 2
stream = Task.async_stream(collection, Mod, :expensive_fun, [], max_concurrency: max_concurrency)
Enum.to_list(stream)

Если вам не нужны результаты вычисления, вы можете запустить поток с помощью Stream.run/1. Также установите ordered: false, так как вам также не нужен порядок результатов:

stream = Task.async_stream(collection, Mod, :expensive_fun, [], ordered: false)
Stream.run(stream)

await(task, timeout \\ 5000)

Характеристики

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

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

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

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

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

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

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

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

Примеры

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

child_spec(arg)

Характеристики

child_spec(term()) :: Supervisor.child_spec()

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

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

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

shutdown(task, shutdown \ 5000)

Характеристики

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

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

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

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

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

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

start(fun)

Характеристики

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

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

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

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

start(module, function_name, args)

Характеристики

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

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

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

start_link(fun)

Характеристики

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

Запускает процесс, связанный с текущим процессом.

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

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

start_link(module, function_name, args)

Характеристики

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

Запускает задачу в рамках дерева надзора.

yield(task, timeout \\ 5000)

Характеристики

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

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

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

  • задача завершилась по причине :normal
  • она не связана с вызывающим процессом
  • вызывающий процесс отслеживает завершения

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

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

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

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

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

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

yield_many(tasks, timeout \\ 5000)

Характеристики

yield_many([t()], timeout()) :: [{t(), {:ok, term()} | {:exit, term()} | nil}]

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

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

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

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

Таймаут, в миллисекундах или :infinity, может быть задан со значением по умолчанию 5000.

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

Пример

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

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

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

tasks_with_results = Task.yield_many(tasks, 5000)

results =
  Enum.map(tasks_with_results, fn {task, res} ->
    # Shut down the tasks that did not reply nor exit
    res || Task.shutdown(task, :brutal_kill)
  end)

# Here we are matching only on {:ok, value} and
# ignoring {:exit, _} (crashed tasks) and `nil` (no replies)
for {:ok, value} <- results do
  IO.inspect(value)
end

В примере выше мы создаём задачи, которые спят от 1 до 10 секунд и возвращают количество секунд, которые они спали. Если выполнить код сразу, вы должны увидеть числа от 1 до 5, поскольку это те задачи, которые ответили в заданный промежуток времени. Все остальные задачи будут остановлены с помощью вызова Task.shutdown/2.

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

Spec-Zone.ru

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