task (ex_stdlib v0.3.0)

View Source

Conveniences for spawning and awaiting tasks.

Tasks are processes meant to execute one particular action throughout their lifetime, often with little or no communication with other processes. The most common use case for tasks is to convert sequential code into concurrent code by computing a value asynchronously:

   Task = task:async(fun() -> do_some_work() end),
   Res = do_some_other_work(),
   Res + task:await(Task).

Tasks spawned with async/1 can be awaited on by their caller process (and only their caller) as shown in the example above. They are implemented by spawning a process that sends a message to the caller once the given computation is performed.

async and await

When invoked, async/1 creates a new process that is linked to and monitored by the caller. Once the task action finishes, a message is sent to the caller with the result.

await/2 is used to read the message sent by the task.

There are two important things to consider when using async:

1. If you are using async tasks, you **must await** a reply as they are *always* sent. If you are not expecting a reply, consider using start_link/1 as detailed below.

2. Async tasks link the caller and the spawned process. This means that, if the caller crashes, the task will crash too and vice-versa. This is on purpose: if the process meant to receive the result no longer exists, there is no purpose in completing the computation.

Unlike Elixir, an exception raised inside the task function does not crash the task process. It is sent back to the caller instead, where await/2 re-raises it and yield/2 returns {exit, Reason}.

Replies are sent to a monitor alias, so once a task has been awaited, ignored or shut down, late replies are dropped instead of piling up in the caller's mailbox.

Tasks are processes

Tasks are processes and so data will need to be completely copied to them. Consider extracting only the necessary data before creating the task, or move the data loading altogether to the task.

Summary

Functions

Starts a task that must be awaited on.

Starts a task that must be awaited on with the given module, function, and arguments.

Runs Fun concurrently on each element of List with default options.

Runs Fun concurrently on each element of List.

Like async_stream/5 with default options.

Runs Module:Function(Elem, Args...) concurrently on each element of List.

Awaits a task reply and returns the result, with a timeout of 5000 ms.

Awaits a task reply and returns the result.

Awaits replies from multiple tasks, with a timeout of 5000 ms.

Awaits replies from multiple tasks and returns them in the order of the given tasks.

Returns a child specification to start a task under a supervisor.

Returns a task that is already completed with the given result.

Ignores an existing task.

Unlinks and shuts down the task, with a timeout of 5000 ms.

Unlinks and shuts down the task, and then checks for a reply.

Starts a task with the given function.

Starts a task with the given module, function, and arguments.

Starts a task as part of a supervision tree with the given function.

Starts a task as part of a supervision tree with the given module, function, and arguments.

Yields to a task, with a timeout of 5000 ms.

Temporarily blocks the caller waiting for a task reply.

Yields to multiple tasks, with a timeout of 5000 ms.

Yields to multiple tasks in the given time interval.

Types

stream_option/0

-type stream_option() ::
          {max_concurrency, pos_integer()} |
          {ordered, boolean()} |
          {timeout, timeout()} |
          {on_timeout, exit | kill_task} |
          {zip_input_on_exit, boolean()}.

task/0

-opaque task()

task_fun/0

-type task_fun() :: fun(() -> term()).

yield_many_option/0

-type yield_many_option() ::
          {limit, pos_integer()} | {timeout, timeout()} | {on_timeout, nothing | ignore | kill_task}.

yield_result/0

-type yield_result() :: {ok, term()} | {exit, term()} | nil.

Functions

async(Fun)

-spec async(task_fun()) -> task().

Starts a task that must be awaited on.

This function spawns a new process linked to and monitored by the caller that will execute the given function. A task record is returned containing the information needed to await the result.

Example

   Task = task:async(fun() -> 1 + 1 end),
   2 = task:await(Task).

async(Module, Function, Args)

-spec async(module(), atom(), [term()]) -> task().

Starts a task that must be awaited on with the given module, function, and arguments.

Similar to async/1 except the function to be started is specified by the given Module, Function, and Args. The Module, Function, and its arity are stored as a tuple in the mfa field for reflection purposes.

async_stream(List, Fun)

-spec async_stream(list(), fun((term()) -> term())) -> [{ok, term()} | {exit, term()}].

Runs Fun concurrently on each element of List with default options.

async_stream(List, Fun, Opts)

-spec async_stream(list(), fun((term()) -> term()), [stream_option()]) ->
                      [{ok, term()} | {exit, term()}].

Runs Fun concurrently on each element of List.

Each element is processed by its own task, linked to and monitored by the caller. Unlike Elixir, which returns a lazy stream, this function runs eagerly and returns the list of results: {ok, Value} for each successful task and {exit, Reason} for each task that raised or exited ({exit, {Element, Reason}} with zip_input_on_exit).

Options:

  • max_concurrency - the maximum number of tasks running at the same time. Defaults to erlang:system_info(schedulers_online).
  • ordered - whether results are returned in input order. Defaults to true.
  • timeout - the maximum time, in milliseconds or infinity, each task is allowed to run. Defaults to 5000.
  • on_timeout - exit (default) makes the caller exit with {timeout, {task, async_stream, [Timeout]}} after killing all running tasks; kill_task kills only the task that timed out and reports {exit, timeout} for it.
  • zip_input_on_exit - include the input element in exit results. Defaults to false.

async_stream(List, Module, Function, Args)

-spec async_stream(list(), module(), atom(), [term()]) -> [{ok, term()} | {exit, term()}].

Like async_stream/5 with default options.

async_stream(List, Module, Function, Args, Opts)

-spec async_stream(list(), module(), atom(), [term()], [stream_option()]) ->
                      [{ok, term()} | {exit, term()}].

Runs Module:Function(Elem, Args...) concurrently on each element of List.

Each element is prepended to Args. See async_stream/3 for options and the returned value.

await(Task)

-spec await(task()) -> term().

Awaits a task reply and returns the result, with a timeout of 5000 ms.

await(Task, Timeout)

-spec await(task(), timeout()) -> term().

Awaits a task reply and returns the result.

If the task raised an exception, it is re-raised in the caller with the same class, reason and stacktrace. If the task process exits, the caller exits with {Reason, {task, await, [Task, Timeout]}}. If the timeout is exceeded, the caller exits with {timeout, {task, await, [Task, Timeout]}}.

This function can only be called by the process that started the task.

await_many(Tasks)

-spec await_many([task()]) -> [term()].

Awaits replies from multiple tasks, with a timeout of 5000 ms.

await_many(Tasks, Timeout)

-spec await_many([task()], timeout()) -> [term()].

Awaits replies from multiple tasks and returns them in the order of the given tasks.

The timeout applies to all tasks together. If any task raises, the exception is re-raised in the caller; if any task process exits or the timeout is exceeded, the caller exits. In both cases the remaining tasks are demonitored, but not shut down.

child_spec(Fun)

-spec child_spec(task_fun() | {module(), atom(), [term()]}) -> supervisor:child_spec().

Returns a child specification to start a task under a supervisor.

Arg is either a zero-arity function or a {Module, Function, Args} tuple. The task is started with start_link and is temporary.

completed(Result)

-spec completed(term()) -> task().

Returns a task that is already completed with the given result.

It can be awaited or yielded like any other task. This is useful for APIs that work with both synchronous and asynchronous operations.

ignore(Task)

-spec ignore(task()) -> yield_result().

Ignores an existing task.

The task continues running, but it is unlinked and its reply will be discarded. Returns {ok, Reply} or {exit, Reason} if the task had already completed, nil otherwise.

shutdown(Task)

-spec shutdown(task()) -> yield_result().

Unlinks and shuts down the task, with a timeout of 5000 ms.

shutdown(Task, How)

-spec shutdown(task(), timeout() | brutal_kill) -> yield_result().

Unlinks and shuts down the task, and then checks for a reply.

Returns {ok, Reply} if the reply was received while shutting down the task, {exit, Reason} if the task raised or died, otherwise nil.

The second argument is either a timeout or brutal_kill. With a timeout, a shutdown exit signal is sent to the task and, if it does not exit within the timeout, it is killed. With brutal_kill the task is killed straight away.

start(Fun)

-spec start(task_fun()) -> {ok, pid()}.

Starts a task with the given function.

The task is not linked to the caller. This should only be used when the task is used for side-effects (like I/O) and you have no interest in its results nor if it completes successfully.

start(Module, Function, Args)

-spec start(module(), atom(), [term()]) -> {ok, pid()}.

Starts a task with the given module, function, and arguments.

start_link(Fun)

-spec start_link(task_fun()) -> {ok, pid()}.

Starts a task as part of a supervision tree with the given function.

The task is linked to the caller and its result cannot be awaited. It's meant for fire-and-forget operations under a supervisor.

start_link(Module, Function, Args)

-spec start_link(module(), atom(), [term()]) -> {ok, pid()}.

Starts a task as part of a supervision tree with the given module, function, and arguments.

yield(Task)

-spec yield(task()) -> yield_result().

Yields to a task, with a timeout of 5000 ms.

yield(Task, Timeout)

-spec yield(task(), timeout()) -> yield_result().

Temporarily blocks the caller waiting for a task reply.

Returns {ok, Reply} if the reply is received, nil if no reply arrived within the timeout, and {exit, Reason} if the task raised or its process already exited. On nil the task keeps running and can be yielded again, awaited, ignored or shut down.

yield_many(Tasks)

-spec yield_many([task()]) -> [{task(), yield_result()}].

Yields to multiple tasks, with a timeout of 5000 ms.

yield_many(Tasks, Timeout)

-spec yield_many([task()], timeout() | [yield_many_option()]) -> [{task(), yield_result()}].

Yields to multiple tasks in the given time interval.

Returns a list of {Task, Result} tuples in the order of the given tasks, where Result is as in yield/2. The second argument is either a timeout or a list of options:

  • limit - stop waiting once this many tasks have replied or exited. Defaults to the number of tasks.
  • timeout - the maximum time to wait. Defaults to 5000.
  • on_timeout - what to do with the tasks that did not reply when the timeout is reached: nothing (default), ignore or kill_task.