Исходный код Задача
Удобства для запуска и ожидания задач.
Задачи — это процессы, предназначенные для выполнения конкретного действия на протяжении всего своего жизненного цикла, часто с минимальным или нулевым взаимодействием с другими процессами. Наиболее распространенный случай использования задач — преобразование последовательного кода в конкурентный код путём асинхронного вычисления значения:
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:
Если вы используете асинхронные задачи, вы обязаны ожидать ответа, так как они всегда отправляются. Если вы не ожидаете ответа, рассмотрите использование
Task.start_link/1, как описано ниже.Асинхронные задачи связывают вызывающий процесс и запущенный процесс. Это означает, что если вызывающий процесс завершится с ошибкой, задача также завершится с ошибкой, и наоборот. Это сделано намеренно: если процесса, предназначенного для получения результата, больше нет, нет смысла завершать вычисление. Если этого не нужно, следует использовать задачи с надзором, описанные в последующем разделе.
Задачи — это процессы
Задачи — это процессы, поэтому данные нужно будет полностью скопировать в них. Рассмотрим следующий пример кода:
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, передавая имя надзирателя вместо идентификатора процесса:
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 — идентификатор процесса, который вызвал текущий процесс, pid2 вызвал pid_n, а pid2 был вызван pid1.
Если задача завершается с ошибкой, поле вызывающих процессов включается в состав сообщения журнала как метаданные под ключом :callers.
Краткое описание
Типы
- async_stream_option()
Параметры, задаваемые для функций
async_stream.- ref()
Непрозрачная ссылка на задачу.
- t()
Тип задачи.
Функции
- %Task{}
Структура 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{} для получения информации о каждом поле структуры.
Функции
%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 имеет те же свойства после добавления вызова async. Например, представьте, что у вас есть это:
x = heavy_fun() y = some_fun() x + y
Теперь вы хотите сделать heavy_fun() асинхронным:
x = Task.async(&heavy_fun/0) y = some_fun() Task.await(x) + y
Как и прежде, если heavy_fun/0 завершится ошибкой, вся вычисление завершится ошибкой, включая вызывающий процесс. Если вы не хотите, чтобы задача завершилась ошибкой, вы должны изменить код heavy_fun/0 таким же образом, как вы бы это сделали, если бы не было асинхронного вызова. Например, чтобы вернуть {:ok, val} | :error результаты или, в более экстремальных случаях, используя try/rescue. Другими словами, асинхронная задача должна рассматриваться как расширение вызывающего процесса, а не как механизм для изоляции его от всех ошибок.
Если вы не хотите связывать вызывающий процесс с задачей, то вы должны использовать контролируемую задачу с Task.Supervisor и вызвать Task.Supervisor.async_nolink/2.
В любом случае, избегайте следующих действий:
-
Установку
:trap_exitвtrue- отлов завершений следует использовать только в особых случаях, так как это сделает ваш процесс неуязвимым не только к завершениям задачи, но и к любым другим процессам.Более того, даже при отлове завершений вызов
awaitвсё равно завершится ошибкой, если задача завершилась без возвращения своего результата. Отключение связи процесса задачи, запущенного с
async/await. Если вы отключите связь процессов, и задача не принадлежит ни одному контролёру, вы можете оставить висящие задачи в случае завершения вызывающего процесса.
Метаданные
Созданная с помощью этой функции задача сохраняет :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.sum_by(stream, fn {:ok, num} -> num 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} , если вызывающий процесс перехватывает завершения. Можно задать перехват завершений с помощью опции :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 - количество задач, чтобы убедиться, что все вызовы выполняются конкурирующим образом.
Внимание: неограниченный асинхронный + take
Если вы хотите потенциально обработать большое количество элементов и сохранить только часть результатов, вы можете обработать больше элементов, чем нужно. Посмотрим пример:
1..100 |> Task.async_stream(fn i -> Process.sleep(100) IO.puts(to_string(i)) end) |> Enum.take(10)
Запуск приведённого выше примера на машине с 8 ядрами обработает 16 элементов, даже если вы хотите только 10, так как async_stream/3 обрабатывает элементы конкурирующим образом. Это происходит потому, что он обработает 8 элементов сразу. Затем все 8 элементов завершатся примерно в одно и то же время, вызывая запуск 8 дополнительных элементов для обработки. Из этих дополнительных 8 будет использовано только 2, а остальные будут завершены.
В зависимости от проблемы вы можете отфильтровать или ограничить количество элементов на первом этапе:
1..100 |> Stream.take(10) |> Task.async_stream(fn i -> Process.sleep(100) IO.puts(to_string(i)) end) |> Enum.to_list()
В других случаях, вы вероятно захотите настроить :max_concurrency для ограничения того, сколько элементов может быть переработано ценой снижения конкурентности. Вы также можете установить количество извлекаемых элементов, кратным :max_concurrency. Например, установив max_concurrency: 5 в приведённом выше примере.
await(task, timeout \\ 5000)Source
@spec await(t(), timeout()) :: term()
Ожидает ответа от задачи и возвращает его.
В случае смерти процесса задачи вызывающий процесс завершится с той же причиной, что и задача.
Таймаут, в миллисекундах или :infinity, может быть задан со значением по умолчанию 5000. Если таймаут истечёт, вызывающий процесс завершится. Если процесс задачи связан с вызывающим процессом, как это происходит при запуске задачи с использованием async, то процесс задачи также завершится. Если процесс задачи перехватывает завершения или не связан с вызывающим процессом, он продолжит выполняться.
Эта функция предполагает, что монитор задачи всё ещё активен или сообщение монитора :DOWN находится в очереди сообщений. Если мониторинг был отменён или сообщение уже получено, эта функция будет ожидать в течение таймаута, ожидая сообщения.
Эта функция может быть вызвана только один раз для любой данной задачи. Если вам нужно многократно проверять, завершило ли долго выполняющаяся задача вычисления, используйте yield/2 вместо этого.
Примеры
iex> task = Task.async(fn -> 1 + 1 end) iex> Task.await(task) 2
Совместимость с поведением OTP
Не рекомендуется await долго выполняющуюся задачу внутри поведения OTP, такого как GenServer. Вместо этого вы должны сопоставлять сообщение, пришедшее от задачи внутри вашего обратного вызова GenServer.handle_info/2.
GenServer получит два сообщения по handle_info/2:
{ref, result}- сообщение ответа, гдеref- ссылка на монитор, возвращённаяtask.ref, аresult- результат задачи{:DOWN, ref, :process, pid, reason}- так как все задачи также отслеживаются, вы также получите сообщение:DOWN, доставленноеProcess.monitor/1. Если вы получите сообщение:DOWNбез ответа, это означает, что задача завершилась ошибкой
Ещё одно соображение, которое следует учитывать, заключается в том, что задачи, запущенные с помощью Task.async/1, всегда связаны со своими вызывающими сторонами, и вы, возможно, не захотите, чтобы GenServer завершился ошибкой, если задача завершится ошибкой. Поэтому предпочтительнее использовать Task.Supervisor.async_nolink/3 внутри поведения OTP. Для полноты картины вот пример GenServer, запускающего задачи и обрабатывающего их результаты:
defmodule GenServerTaskExample do
use GenServer
def start_link(opts) do
GenServer.start_link(__MODULE__, :ok, opts)
end
def init(_opts) do
# We will keep all running tasks in a map
{:ok, %{tasks: %{}}}
end
# Imagine we invoke a task from the GenServer to access a URL...
def handle_call(:some_message, _from, state) do
url = ...
task = Task.Supervisor.async_nolink(MyApp.TaskSupervisor, fn -> fetch_url(url) end)
# After we start the task, we store its reference and the url it is fetching
state = put_in(state.tasks[task.ref], url)
{:reply, :ok, state}
end
# If the task succeeds...
def handle_info({ref, result}, state) do
# The task succeed so we can demonitor its reference
Process.demonitor(ref, [:flush])
{url, state} = pop_in(state.tasks[ref])
IO.puts("Got #{inspect(result)} for URL #{inspect url}")
{:noreply, state}
end
# If the task fails...
def handle_info({:DOWN, ref, _, _, reason}, state) do
{url, state} = pop_in(state.tasks[ref])
IO.puts("URL #{inspect url} failed with reason #{inspect(reason)}")
{:noreply, state}
end
end
После определения сервера вы захотите запустить надзирателя задач выше и GenServer в вашей дереве надзора:
children = [
{Task.Supervisor, name: MyApp.TaskSupervisor},
{GenServerTaskExample, name: MyApp.GenServerTaskExample}
]
Supervisor.start_link(children, strategy: :one_for_one) await_many(tasks, timeout \\ 5000)Source
@spec await_many([t()], timeout()) :: [term()]
Ожидает ответы от нескольких задач и возвращает их.
Эта функция получает список задач и ожидает их ответы в заданном интервале времени. Она возвращает список результатов в том же порядке, что и задачи, предоставленные в аргументе tasks.
Если любой из процессов задач завершится ошибкой, вызывающий процесс завершится с той же причиной, что и эта задача.
Таймаут, в миллисекундах или :infinity, может быть задан со значением по умолчанию 5000. Если таймаут истечёт, вызывающий процесс завершится. Любой процесс задачи, связанный с вызывающим процессом (что происходит при запуске задачи с помощью async), также завершится. Любые процессы задач, перехватывающие завершения или не связанные с вызывающим процессом, будут продолжать выполняться.
Эта функция предполагает, что мониторы задач всё ещё активны или сообщение монитора :DOWN находится в очереди сообщений. Если какой-либо из мониторингов был отменён или сообщение уже получено, эта функция будет ожидать в течение всего таймаута.
Эта функция может быть вызвана только один раз для любой данной задачи. Если вам нужно многократно проверять, завершило ли долго выполняющаяся задача вычисления, используйте yield_many/2 вместо этого.
Совместимость с поведением OTP
Не рекомендуется await долго выполняющиеся задачи внутри поведения OTP, такого как GenServer. Смотрите await/2 для получения дополнительной информации.
Примеры
iex> tasks = [ ...> Task.async(fn -> 1 + 1 end), ...> Task.async(fn -> 2 + 3 end) ...> ] iex> Task.await_many(tasks) [2, 5]
child_spec(arg)Source
@spec child_spec(term()) :: Supervisor.child_spec()
Возвращает спецификацию для запуска задачи под надзором.
arg передаётся в качестве аргумента функции Task.start_link/1 в поле :start спецификации.
Для получения дополнительной информации обратитесь к модулю Supervisor, функции Supervisor.child_spec/2 и типу Supervisor.child_spec/0.
completed(result)Source
@spec completed(any()) :: t()
Запускает задачу, которая немедленно завершается с указанным result.
В отличие от async/1, эта задача не создаёт связанный процесс. К ней можно обратиться или передать её, как и к любой другой задаче.
Использование
В некоторых случаях полезно создать задачу «completed», представляющую задачу, которая уже выполнилась и сгенерировала результат. Например, при обработке данных вы можете определить, что некоторые входные данные недействительны, прежде чем отправлять их на дальнейшую обработку:
def process(data) do
tasks =
for entry <- data do
if invalid_input?(entry) do
Task.completed({:error, :invalid_input})
else
Task.async(fn -> further_process(entry) end)
end
end
Task.await_many(tasks)
end
Во многих случаях можно избежать использования Task.completed/1, просто вернув результат напрямую. Вам обычно понадобится этот вариант только при работе с смешанной асинхронностью, когда группа входных данных будет обрабатываться частично синхронно, а частично асинхронно.
ignore(task)Source
@spec ignore(t()) :: {:ok, term()} | {:exit, term()} | nil Игнорирует существующую задачу.
Это означает, что задача будет продолжать выполняться, но она будет разъединена, и вы больше не сможете её передавать, ждать или останавливать.
Возвращает {:ok, reply} , если ответ получен до игнорирования задачи, {:exit, reason} если задача завершилась до игнорирования, в противном случае nil.
Важно: избегайте использования Task.async/1,3, а затем немедленно игнорирования задачи. Если вы хотите запустить задачи, результаты которых вам не важны, используйте Task.Supervisor.start_child/2 вместо этого.
shutdown(task, shutdown \\ 5000)Source
@spec shutdown(t(), timeout() | :brutal_kill) :: {:ok, term()} | {:exit, term()} | nil Разъединяет и завершает задачу, а затем проверяет наличие ответа.
Возвращает {:ok, reply} , если ответ получен во время завершения задачи, {:exit, reason} если задача завершилась неудачно, в противном случае nil. После завершения вы больше не можете ждать или передавать её.
Второй аргумент — это либо таймаут, либо :brutal_kill. В случае таймаута задаче отправляется сигнал выхода :shutdown, и если она не завершается в течение таймаута, она убивается. При использовании :brutal_kill задача убивается сразу. В случае аварийного завершения задачи (возможно, убитой другим процессом) эта функция завершится с той же причиной.
Вызов этой функции не требуется при завершении вызывающего процесса, если не происходит выход с причиной :normal или если задача обрабатывает выходы. Если вызывающий процесс завершается с причиной, отличной от :normal, и задача не обрабатывает выходы, сигнал выхода вызывающего процесса остановит задачу. Вызывающий процесс может завершиться с причиной :shutdown для остановки всех связанных процессов, включая задачи, которые не обрабатывают выходы, без создания сообщений об ошибках.
Если к задаче нет привязанного процесса, например, задач, запущенных с помощью Task.completed/1, мы проверяем наличие ответа или ошибки соответствующим образом, но без остановки процесса.
Если монитор задачи уже отслеживается или получен, и в очереди сообщений нет ожидаемого ответа, эта функция вернёт {:exit, :noproc}, так как причина выхода не может быть определена.
start(fun)Source
@spec start((-> any())) :: {:ok, pid()} Запускает задачу.
fun должно быть безымянным нулеарным функцией.
Это следует использовать только в том случае, если задача используется для побочных эффектов (например, ввода/вывода), и вы не заинтересованы в её результатах или в том, завершилась ли она успешно.
Если текущий узел остановлен, узел завершится, даже если задача не была завершена. По этой причине рекомендуется использовать Task.Supervisor.start_child/2 вместо этого, что позволяет управлять временем завершения с помощью параметра :shutdown.
start(module, function_name, args)Source
@spec start(module(), atom(), [term()]) :: {:ok, pid()} Запускает задачу.
Это следует использовать только в том случае, если задача используется для побочных эффектов (например, ввода/вывода), и вы не заинтересованы в её результатах или в том, завершилась ли она успешно.
Если текущий узел остановлен, узел завершится, даже если задача не была завершена. По этой причине рекомендуется использовать Task.Supervisor.start_child/2 вместо этого, что позволяет управлять временем завершения с помощью параметра :shutdown.
start_link(fun)Source
@spec start_link((-> any())) :: {:ok, pid()} Запускает задачу в рамках дерева надзора с заданным fun.
fun должно быть безымянным нулеарным функцией.
Используется для запуска статически контролируемой задачи в рамках дерева надзора.
start_link(module, function, args)Source
@spec start_link(module(), atom(), [term()]) :: {:ok, pid()} Запускает задачу в рамках дерева надзора с заданным module, function, и args.
Используется для запуска статически контролируемой задачи в рамках дерева надзора.
yield(task, timeout \\ 5000)Source
@spec yield(t(), timeout()) :: {:ok, term()} | {:exit, term()} | nil Временно блокирует вызывающий процесс, ожидая ответа от задачи.
Возвращает {:ok, reply} если ответ получен, nil если ответ не пришёл, или {:exit, reason} если задача уже завершилась. Имейте в виду, что обычно сбой задачи также приводит к завершению процесса, владеющего задачей. Поэтому эта функция может вернуть {:exit, reason} если выполняются хотя бы одно из следующих условий:
- процесс задачи завершился с причиной
:normal - задача не связана с вызывающим процессом (задача была запущена с помощью
Task.Supervisor.async_nolink/2илиTask.Supervisor.async_nolink/4) - вызывающий процесс обрабатывает выходы
Таймаут, в миллисекундах или :infinity, может быть задан со значением по умолчанию 5000. Если время истечёт до получения сообщения от задачи, эта функция вернёт nil и монитор останется активным. Поэтому yield/2 можно вызывать несколько раз для одной и той же задачи.
Эта функция предполагает, что монитор задачи по-прежнему активен или сообщение монитора :DOWN находится в очереди сообщений. Если монитор был отменён или сообщение уже получено, эта функция будет ожидать в течение таймаута, ожидая сообщения.
Если вы хотите остановить задачу, если она не ответит в течение timeout миллисекунд, вы должны объединить её с shutdown/1, как показано ниже:
case Task.yield(task, timeout) || Task.shutdown(task) do
{:ok, result} ->
result
nil ->
Logger.warning("Failed to get a result in #{timeout}ms")
nil
end
Если вы хотите проверить задачу, но оставить её работающей после таймаута, вы можете объединить её с ignore/1, как показано ниже:
case Task.yield(task, timeout) || Task.ignore(task) do
{:ok, result} ->
result
nil ->
Logger.warning("Failed to get a result in #{timeout}ms")
nil
end
Это гарантирует, что если задача завершится после таймаута, но до вызова shutdown/1, вы всё равно получите результат, поскольку shutdown/1 предназначена для обработки этого случая и возвращает результат.
yield_many(tasks, opts \\ [])Source
@spec yield_many([t()], timeout()) :: [{t(), {:ok, term()} | {:exit, term()} | nil}] @spec yield_many([t()],
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— задача, превысившая таймаут, будет убита.
-
© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.18.1/Task.html