NUMA-Blind Agent Placement and the 60% Swarm

Why default scheduling quietly halves throughput on dual-socket boards.

by

A dual-socket server is not one fast computer. It is two computers that happen to share a power supply and a chassis, connected by a link that is an order of magnitude slower than the memory controller sitting next to each CPU. When an agent runtime treats that machine as one flat pool of cores, the kernel scheduler is free to migrate a worker from socket 0 to socket 1 between two adjacent tool calls. The worker keeps running. It just runs against memory that now lives on the wrong side of the interconnect.

That is the whole failure mode. Nothing crashes, nothing logs an error, and the service-level metric that matters most — useful work per unit time — drops. Teams that run agent swarms on dual-socket hardware often describe the same shape: the box looks underutilized, per-core throughput is mediocre, and the problem disappears when the workload is pinned. The usual diagnosis is "the runtime is slow." The actual diagnosis is that the runtime never told the scheduler what it knows about its own topology.

What NUMA actually is

Non-uniform memory access means memory latency depends on which CPU core issues the load. On a two-socket board, each socket has its own memory controller and its own DIMMs. A core reading a local address hits its own controller. A core reading an address owned by the other socket's controller has to traverse the inter-socket fabric — on current server platforms that is a UPI, Infinity Fabric, or similar coherent link.

$ numactl --hardware
available: 2 nodes (0-1)
node 0 cpus: 0 1 2 ... 47
node 0 size: 262144 MB
node 1 cpus: 48 49 ... 95
node 1 size: 262144 MB
node distances:
node   0   1
  0:  10  21
  1:  21  10

The distance matrix is the important part. Local access is scored 10; remote access 21. That ratio is not a benchmark, it is the platform's own declared cost model, and it is usually conservative. Real remote access is often worse once you account for contention on the fabric.

For a workload that streams large tensors through a single thread, this matters less than people assume, because the prefetcher and the memory-level parallelism hide a lot of it. For an agent runtime it matters a great deal, because agent workloads are not streaming. They are pointer-chasing: conversation state, tool schemas, retrieval results, JSON trees that get parsed and re-parsed. Every one of those is a dependent load with poor locality. Remote memory latency shows up directly in the critical path.

Why the default scheduler is NUMA-blind in practice

The Linux CFS scheduler is not ignorant of NUMA; it has had NUMA balancing for years. The problem is that NUMA balancing optimizes for the wrong thing when the workload is a swarm of short-lived, bursty tasks.

Automatic NUMA balancing periodically unmaps pages, takes faults, and migrates data toward the CPU that is touching it. That works well for long-running processes with stable working sets. An agent worker is the opposite: it wakes, does a few milliseconds of work, blocks on a tool call or a network round trip, and wakes again possibly on a different core. The balancer never converges. It generates fault traffic and page migrations for a working set that is about to go cold anyway.

Meanwhile, the load balancer is doing exactly what it was designed to do: spreading runnable tasks evenly across all cores. On a single-socket machine that is optimal. On a dual-socket machine it means the scheduler will happily place half of your workers on socket 1 while their conversation state was allocated by a process that started on socket 0. The first-touch policy means memory lands on whichever node first writes to it, and nothing in the runtime's design ensures that the first writer and the frequent reader are the same socket.

The agent-specific twist

Most service workloads have a request-per-thread or request-per-process model that maps cleanly onto cores. Agent runtimes usually do not. A typical design looks like this:

  • A supervisor process accepts work and maintains a queue.
  • A pool of worker processes or threads pulls tasks.
  • Each task carries a state object: message history, scratchpad, tool definitions, retrieved context.
  • Workers call out to inference, vector search, and graph lookups, each of which has its own memory footprint.
# A representative agent worker loop. Note the state is per-task,
# allocated wherever the worker happens to be running.
def worker(task_queue, state_store):
    while True:
        task = task_queue.get()
        state = state_store.load(task.id)   # first touch decides NUMA node
        while not state.done:
            action = plan(state)            # pointer-heavy, latency bound
            result = execute(action)
            state.append(result)
        state_store.save(state)

The state object is allocated on whichever node the worker is running on at that instant. If the worker later migrates — and with a bursty, blocking workload it will — every subsequent access to that state is potentially remote. The state store itself, if it is a shared in-memory cache, has the same problem at a larger scale: one copy of the data, touched by workers on both sockets, so half of them are always paying the remote penalty.

This is why the symptom is so often "about half the workers are slow." It is not a random half. It is the half that ended up on the far socket from the data.

Fixing it: pin, partition, or replicate

There are three honest strategies. They trade off differently and none of them is free.

1. Pin workers and their state to a socket

The blunt instrument: create one worker pool per NUMA node, bind each pool to that node's CPUs and memory, and route tasks to the pool whose data locality matches.

# One pool per socket, bound to CPUs and memory of that node.
numactl --cpunodebind=0 --membind=0 ./agent-pool --pool-id=0
numactl --cpunodebind=1 --membind=1 ./agent-pool --pool-id=1

This works, and it is the approach most teams land on first. The cost is that you now own the routing problem. A task whose state lives on node 0 must go to pool 0, or you pay a migration. If load is uneven — and agent traffic is bursty, so it will be — one pool idles while the other queues. You have traded a latency problem for a capacity problem, and you have to decide which one you would rather have.

--membind is also stricter than people expect. If node 0 runs out of memory, the allocation fails rather than falling back. For a state-heavy agent runtime that is a real operational hazard; --preferred gives you a soft preference and a fallback, at the cost of occasionally going remote under pressure.

2. Replicate state per socket

If the state is read-mostly during a task, keep a copy on each node and accept the coherence cost at write time. This is the approach that scales best for retrieval-heavy agents, where the expensive part is the vector search and the graph traversal, not the mutation of conversation state.

The trade-off is memory. Two copies of every active session is a real cost on a box where each socket might have a few hundred gigabytes. It also means you need a invalidation path that is correct under concurrent writes, which is where most implementations get subtly wrong. A generation counter per session, checked on read, is usually enough and is far cheaper than a full coherence protocol.

3. Make the runtime topology-aware

This is the option that ages best. Instead of pinning at the process level, have the runtime discover the topology at startup and make placement decisions as part of scheduling.

import os

def numa_nodes():
    nodes = {}
    for entry in os.listdir("/sys/devices/system/node"):
        if not entry.startswith("node"):
            continue
        node_id = int(entry[4:])
        with open(f"/sys/devices/system/node/{entry}/cpulist") as f:
            nodes[node_id] = parse_cpulist(f.read().strip())
    return nodes

# At task admission: choose the node that owns the state,
# then dispatch to a worker bound to that node.

The runtime then routes each task to a worker on the node that holds its state, and only falls back to remote execution when the local pool is saturated. That fallback is explicit and measurable, which is the point: you want the remote accesses to be a decision you made, not an accident the scheduler made for you.

What to watch for in practice

A few things consistently trip people up when they move to topology-aware placement.

Hyperthread siblings are not separate cores. Pinning to a CPU list that includes sibling threads of the same physical core gives you two workers competing for one core's execution resources. The topology files expose this; lscpu -e and the thread_siblings_list entries are the ground truth.

The interconnect is shared. Two workers on the same node both going remote will contend with each other. The distance matrix gives you a static cost; the real cost rises with concurrent remote traffic. This is why "just let it go remote sometimes" degrades faster than a simple model predicts.

Container runtimes hide the topology. If your agents run in containers, the cpuset the container sees may not reflect the host's NUMA layout, and numactl inside the container may be a no-op or may bind to the wrong node. Check what the container actually sees before trusting a pin.

Inference servers have their own topology opinions. A vLLM instance, for example, has its own tensor-parallel and pipeline-parallel placement decisions, and those interact with where your agent workers sit. Putting the agent workers and the inference server on the same socket is often better than spreading them, even if spread looks more balanced on a utilization graph.

The general lesson

Modern infrastructure hides topology behind abstractions that are correct on a single-socket machine and misleading on a multi-socket one. The scheduler sees a pool of cores. The allocator sees a pool of memory. Neither sees the link between them, and neither knows that your agent's state has a home.

The fix is not exotic. It is a startup-time read of /sys/devices/system/node, a placement decision at task admission, and an explicit fallback when the local pool is full. What makes it worth writing down is that the failure mode is invisible: no errors, no crashes, just a box that runs at a fraction of what its silicon can do, and a team that concludes agent runtimes are inherently slow.

They are not. They are just usually placed as if the machine were flat.

#agent-runtime#cpu-affinity#infrastructure#numa#performance#scheduling
Share — X / Twitter · LinkedIn · HN · Email
Damir Radulić
Founder of RiNET. On the Croatian internet since 1996 (Kvarner Net). In Amsterdam now, building autonomous AI infrastructure that runs on Monday morning when nobody's watching — sovereign stacks, agent swarms, LoRA fine-tuning, civic-intelligence platforms.