Concurrency, Threads & Async
Two Problems, Two Tools
"Concurrency" is not one problem. "Make this computation faster" and "wait for ten thousand sockets at once" have different shapes and different costs, so the first decision is which one you actually have.
CPU-Bound versus I/O-Bound
| Question | Threads | Async / await |
|---|---|---|
| What is waiting? | The CPU is busy | The network or the disk is busy |
| Cores used | Several, truly in parallel | One; tasks interleave at await points |
| Cost per unit of work | An OS thread with its own stack | A closure and a future object |
| How state is shared | Discouraged — send messages instead | Not shared: one thread, no preemption |
| Typical use | Parsing, image processing, worker pools | HTTP clients and servers, sockets, timers |
An async task that calls a CPU-heavy routine blocks every other task in that event loop, and a thread that spends its life waiting on a socket wastes one megabyte of stack to do nothing. Choose the model from the bottleneck, not from fashion.
Enabling Thread Support
Thread support is a compile-time capability, not a runtime library you load. The system documentation for channels states it plainly: "To activate thread support compile with the --threads:on command line switch." Current 2.x releases ship a config/nim.cfg that already sets threads:on, so on a stock install threads work out of the box; a trimmed configuration or an explicit --threads:off removes them, and code that depends on threads should say so:
# Guard self-contained modules with a compile-time check. The standard library
# uses exactly this pattern: threadpool.nim raises an error when threads are off.
when not compileOption("threads"):
{.error: "This module requires --threads:on.".}
# The check is evaluated by the compiler, so the offending build fails early
# with a readable message instead of producing a broken binary.
echo "threads are enabled" # printed only when the guard above passes
Parallel Threads
A Nim thread is an OS thread with its own stack. You create it with createThread, hand it one argument, and join it when you need its results. The procedure that runs on the thread carries the {.thread.} pragma, which is what makes the compiler check that it is safe to start there.
Figure 1 — the main thread spawns workers and communicates only through channels; each thread keeps its own heap.
createThread, Thread[T] and joinThread
var results: array[4, int] # each thread writes exactly one disjoint slot
proc worker(id: int) {.thread.} =
## The 'thread' pragma marks a routine that may start on a new thread.
## Such a routine must not raise: an unhandled exception would kill the
## process, so errors are reported back through data, not exceptions.
var total = 0
for i in 1 .. (id + 1) * 1000:
total += i
results[id] = total # no lock needed: no two threads touch one slot
var threads: array[4, Thread[int]]
for id in 0 .. 3:
createThread(threads[id], worker, id) # returns immediately
for id in 0 .. 3:
joinThread(threads[id]) # block until that thread has finished
echo results # [500500, 2001000, 4501500, 8002000]
Two rules cover most mistakes. First, the argument passed to createThread must not contain garbage-collected data unless you hand over a deepCopy; a string or a seq created on the main thread belongs to that thread's heap. Second, the main thread must joinThread (or joinThreads) before it exits, otherwise the program terminates while workers are still running.
Shared Counters: Locks and Atomics
Deep copying every message is fine for coarse work; for a single shared counter it is wasteful. std/atomics gives lock-free operations for machine words, and std/locks gives mutual exclusion for larger structures.
import std/[locks, atomics]
var hits: Atomic[int] # lock-free, but only for scalars
var total: int # compound state: needs the lock
var guard: Lock
initLock(guard)
proc tally() {.thread.} =
for i in 1 .. 100:
discard hits.fetchAdd(1) # atomic increment: safe from any thread
withLock guard: # 'withLock' releases the lock even if the body
total += i # raises — a bare acquire/release would not
var worker: Thread[void]
createThread(worker, tally)
joinThread(worker)
echo hits.load() # 100 — an atomic read, not a plain load
echo total # 5050
deinitLock(guard) # release the OS resource behind the lock
Atomic operations are sequentially consistent by default, which is the easiest order to reason about; weaker orders are an optimisation for code you can measure. For anything larger than a machine word, or for an invariant spanning several fields, take a lock and keep the critical section short — a lock held across I/O serialises every thread that needs it.
Message Passing with Channels
Nim's answer to shared mutable state is a Channel[T]: a thread-safe queue where every message is deeply copied, so the receiver owns its data and no lifetime is shared. The channel documentation states the three constraints to design around — "objects passed through channels will be deeply copied", "Channels cannot be passed between threads. Use globals or pass them by ptr", and that they are "designed for the Thread type" and unstable with spawn.
open, send, recv and close
var jobs: Channel[string] # module level: a shared, thread-safe heap
var done: Channel[int]
jobs.open(maxItems = 0) # 0 = unbounded queue; N blocks send when full
done.open()
proc producer() {.thread.} =
for i in 1 .. 3:
jobs.send("task " & $i) # deep copy, then enqueue
jobs.send("") # sentinel: an empty message ends the stream
proc consumer() {.thread.} =
var handled = 0
while true:
let job = jobs.recv() # blocks until a message arrives
if job.len == 0: break
inc handled
done.send(handled) # send the result back to the main thread
var producers: array[1, Thread[void]]
var consumers: array[1, Thread[void]]
createThread(producers[0], producer)
createThread(consumers[0], consumer)
joinThreads(producers)
joinThreads(consumers)
echo done.recv() # 3
jobs.close() # free the queue's resources
done.close()
Use maxItems for back pressure: with maxItems = 64 a fast producer blocks once 64 messages are waiting, which bounds memory without extra coordination. When waiting is unacceptable, the non-blocking pair is trySend and tryRecv, where the receiver reports availability instead of blocking:
# Non-blocking receive: the caller keeps working while the queue is empty.
let tried = jobs.tryRecv()
if tried.dataAvailable:
echo "got: ", tried.msg # 'msg' is only meaningful when dataAvailable
else:
echo "nothing yet" # ... do something useful instead of waiting
peek reports how many messages are queued (-1 once the channel is closed), and its own documentation warns it "is dangerous to use as it encourages races" — treat it as a diagnostic, not as a synchronisation tool.
Higher-Level Parallelism: Taskpools
Raw channels are the primitive; production code usually wants a scheduler on top of them. The standard library's std/threadpool module is documented as deprecated — and its documentation points at the community packages taskpools, malebolgia and weave.
Taskpools and FlowVar
import std/taskpools # installed with 'nimble install taskpools'
proc square(x: int): int =
## Runs on a worker thread. A task takes its input and returns its output by
## value, so no two tasks ever share mutable state.
x * x
let pool = Taskpool.new(numThreads = 4) # one worker per core is a fine default
var futures: seq[FlowVar[int]]
for value in 1 .. 8:
futures.add pool.spawn square(value) # queued: it may not have started yet
var total = 0
for future in futures:
total += sync(future) # block until this task's result exists
pool.syncAll() # never shut down with tasks in flight
pool.shutdown()
echo total # 204 — 1² + 2² + ... + 8²
A FlowVar[T] is a handle to a result that does not exist yet: sync blocks for that one value, while isReady asks without blocking. Collect the handles in a sequence and sync them in order — the tasks still run in parallel, and the code stays linear. Both this and std/threadpool require threads to be enabled.
Asynchronous I/O
Threads cost an OS stack and a context switch each; when the work is waiting on the network or on disk, that cost buys nothing. asyncdispatch runs thousands of coroutines on one thread by suspending a routine at every await instead of blocking it.
async, await and Future
import std/asyncdispatch
proc answer(id: int): Future[string] {.async.} =
## 'async' turns the body into a state machine. The returned Future completes
## when the body finishes; awaiting it yields control back to the event loop.
await sleepAsync(10) # a non-blocking pause, not a thread sleep
result = "worker-" & $id
proc main() {.async.} =
# 'asyncCheck' starts a void future and re-raises its failure on the loop;
# 'discard' would drop the exception silently.
asyncCheck answer(1)
let a = await answer(2) # sequential composition
let b = await answer(3)
echo a, " ", b # worker-2 worker-3
waitFor main() # runs the event loop until 'main' completes
Concurrency comes from starting futures before awaiting them: let f = answer(3) schedules the work, and the await you write later only claims the result. Two idioms complete the picture: FutureVar[T] (defined as distinct Future[T]) is a future you complete yourself — the manual way to bridge a callback-based API into async — and all/join from std/asyncfutures awaits a collection at once.
The ordering rule to remember is that the event loop is started exactly once, by waitFor at the top of the program; an async routine never calls it internally. Mixing the two worlds — a thread blocking inside an async routine, or an await inside a {.thread.} proc — is the single most common source of hangs in Nim code.