पाठ 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.
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}
endBound 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.