पाठ 13 / 25

Tasks and Concurrent Work

Run concurrent work with Task, async_stream and supervised tasks.

Simple concurrency with Task

Task wraps a process for a single unit of work. Task.async(fn -> ... end) starts work and returns a task; Task.await(task, timeout) waits for its result (default timeout 5 seconds). The caller is linked to the task, so a crash propagates. Task.await_many waits for several tasks, and Task.yield returns nil on timeout instead of crashing. For processing collections concurrently, Task.async_stream/3 is the workhorse: it runs a function over each element with a bounded max_concurrency (defaulting to the number of schedulers), preserving order and supporting timeout and on_timeout: :kill_task. Use Task.Supervisor.async_nolink when the caller must survive a failing task, such as a web request calling an unreliable API, and Task.Supervisor.start_child for fire-and-forget work. For CPU-bound work, the BEAM spreads processes across all cores automatically. For durable background jobs that must survive restarts, use a job library such as Oban, which stores jobs in PostgreSQL.

Fan-out with async_stream

A list of inputs is processed by a bounded number of concurrent tasks, results in order.

A list of input boxes fanning out to four parallel worker circles, then fanning back in to an ordered list of result boxes.
Figure 5.1 — Bounded concurrency with Task.async_stream.

Concurrent price lookups

async/await for a few calls, async_stream for many.

defmodule Shop.Prices do
  def fetch(sku) do
    Process.sleep(Enum.random(50..200))          # simulate a slow API
    {sku, 4_950}
  end

  def dashboard(customer_id) do
    orders_task = Task.async(fn -> fetch_orders(customer_id) end)
    profile_task = Task.async(fn -> fetch_profile(customer_id) end)
    [orders, profile] = Task.await_many([orders_task, profile_task], 2_000)   # run in parallel
    %{orders: orders, profile: profile}
  end

  def price_all(skus) do
    skus
    |> Task.async_stream(&fetch/1, max_concurrency: 10, timeout: 1_000, on_timeout: :kill_task)
    |> Enum.flat_map(fn
      {:ok, {sku, price}} -> [{sku, price}]
      {:exit, :timeout} -> []                         # skip slow lookups
    end)
    |> Map.new()
  end

  def notify_safely(order_id) do
    # the caller is not linked, so a crash here does not take the caller down
    Task.Supervisor.async_nolink(Shop.TaskSupervisor, fn -> send_email(order_id) end)
  end

  defp fetch_orders(_id), do: [:o1, :o2]
  defp fetch_profile(_id), do: %{name: "Asha"}
  defp send_email(order_id), do: {:sent, order_id}
end

Bound your concurrency

Spawning a task per item is easy on the BEAM, but the services you call are not infinitely scalable. max_concurrency on async_stream protects downstream APIs and databases from overload.

त्वरित जाँच: What does max_concurrency in Task.async_stream control?

  • How many items are processed concurrently at once
  • The total number of items
  • The number of CPU cores
  • The retry count
Answer

How many items are processed concurrently at once — It bounds the number of tasks running at the same time.