Summary
AssignmentGrid distributes agents across every live node without any notion of how much a node can actually host. When a node dies, its share is redistributed over the survivors — which can push them past their own limits, killing them in turn. There is no floor to that process, so a cluster under memory pressure cannot converge, and adding replicas makes it worse rather than better.
We hit this on production last weekend: 34 hours of intermittent unavailability, restart counters up to 334, and no self-recovery.
What already works
To be clear about what I'm not asking for: EventSubscriptionAgentFamily.EvaluateAssignmentsAsync already does the right thing for sharded stores. Multi-database stores get DistributeByGroupAffinity(SchemeName, DatabaseKeyOf, …), keeping all of a shard database's per-tenant agents together on one node so pools scale with databases rather than nodes × databases (#486 / marten#4806). That behaviour is exactly right for our topology and it worked.
The gap
Neither distributor bounds what a node takes on.
AssignmentGrid.Distribution.cs:
var spread = (double)agents.Count / _nodes.Count;
var minimum = (int)Math.Floor(spread);
var maximum = (int)Math.Ceiling(spread);
and group affinity places groups "largest-first onto the least-loaded node". Both balance, neither caps. The implicit assumption is that any live node can host whatever share it is handed.
That assumption breaks whenever an agent's cost is not uniform over time.
What happened
Topology: Marten with a sharded store, 512 shard databases, per-tenant event partitioning. Six projections × 1.463 tenants = 8.792 agents across 6 nodes (one pod per node, enforced by pod anti-affinity).
A release bumped two projection versions, so those read models rebuilt from event 0 — about 381M events. We measured that an agent working off a backlog holds roughly 4.4 MB, against roughly 1 MB for an agent that is caught up. Same agents, same code, four times the memory, purely as a function of having work to do.
At ~1,465 agents per node that is ~6.4 GB during catch-up, against a .NET GC hard limit of 9,216 MiB. Tight but survivable — until one node went over and died. Its ~1,465 agents were then spread over five survivors, taking each of them to ~1,760 agents, over the limit. Then four. Then three.
Two properties made it unrecoverable:
- Adding capacity does not help. Distribution is over live nodes, so new pods join the same doomed rebalance and die before they can relieve anyone. Measured: scaling 6 → 12 replicas raised agents-per-surviving-pod from 1,196 to 1,290.
- Every restart triggers another rebalance. With pods failing every few minutes the grid never settles, and each rebalance is itself the expensive moment.
What finally recovered it was replacing all pods simultaneously (maxUnavailable: 100%) so the grid distributed once across six live nodes instead of six times across a shrinking set. That is a Kubernetes-side workaround for a distribution-side problem.
Proposal
An optional capacity ceiling per node:
// DurabilitySettings
public int? MaxAgentsPerNode { get; set; } = null; // null = today's behaviour
advertised alongside the node's capabilities, and respected by both DistributeEvenly and DistributeByGroupAffinity: when no node can take an agent (or, under group affinity, a group) without exceeding its ceiling, leave it unassigned rather than overload a node. Unassigned agents get picked up on a later evaluation once capacity appears — which it does, as soon as an operator adds a node or the backlog that made agents expensive is worked off.
That converts a cascading failure into visible, bounded degradation: some shards stop progressing, the rest keep running, and the cluster stays up. For our case that is strictly better — a projection falling behind is an internal problem, a fleet that cannot stay up is a customer-facing one.
Two smaller things that would have helped independently:
- Expose the assigned-agent count per node as a metric. We had to infer it from log lines; there was no way to see "this node is carrying 1,760 agents" while it was happening.
- Consider a hysteresis on redistribution. A node that just died is very likely to be replaced within seconds; immediately redistributing its full share is what turned one failure into six.
Willing to help
Happy to prototype this against a fork if the direction is agreeable — we have a reproduction with real scale (8,792 agents, 512 databases) and the production measurements above to validate against.
Summary
AssignmentGriddistributes agents across every live node without any notion of how much a node can actually host. When a node dies, its share is redistributed over the survivors — which can push them past their own limits, killing them in turn. There is no floor to that process, so a cluster under memory pressure cannot converge, and adding replicas makes it worse rather than better.We hit this on production last weekend: 34 hours of intermittent unavailability, restart counters up to 334, and no self-recovery.
What already works
To be clear about what I'm not asking for:
EventSubscriptionAgentFamily.EvaluateAssignmentsAsyncalready does the right thing for sharded stores. Multi-database stores getDistributeByGroupAffinity(SchemeName, DatabaseKeyOf, …), keeping all of a shard database's per-tenant agents together on one node so pools scale with databases rather than nodes × databases (#486 / marten#4806). That behaviour is exactly right for our topology and it worked.The gap
Neither distributor bounds what a node takes on.
AssignmentGrid.Distribution.cs:and group affinity places groups "largest-first onto the least-loaded node". Both balance, neither caps. The implicit assumption is that any live node can host whatever share it is handed.
That assumption breaks whenever an agent's cost is not uniform over time.
What happened
Topology: Marten with a sharded store, 512 shard databases, per-tenant event partitioning. Six projections × 1.463 tenants = 8.792 agents across 6 nodes (one pod per node, enforced by pod anti-affinity).
A release bumped two projection versions, so those read models rebuilt from event 0 — about 381M events. We measured that an agent working off a backlog holds roughly 4.4 MB, against roughly 1 MB for an agent that is caught up. Same agents, same code, four times the memory, purely as a function of having work to do.
At ~1,465 agents per node that is ~6.4 GB during catch-up, against a .NET GC hard limit of 9,216 MiB. Tight but survivable — until one node went over and died. Its ~1,465 agents were then spread over five survivors, taking each of them to ~1,760 agents, over the limit. Then four. Then three.
Two properties made it unrecoverable:
What finally recovered it was replacing all pods simultaneously (
maxUnavailable: 100%) so the grid distributed once across six live nodes instead of six times across a shrinking set. That is a Kubernetes-side workaround for a distribution-side problem.Proposal
An optional capacity ceiling per node:
advertised alongside the node's capabilities, and respected by both
DistributeEvenlyandDistributeByGroupAffinity: when no node can take an agent (or, under group affinity, a group) without exceeding its ceiling, leave it unassigned rather than overload a node. Unassigned agents get picked up on a later evaluation once capacity appears — which it does, as soon as an operator adds a node or the backlog that made agents expensive is worked off.That converts a cascading failure into visible, bounded degradation: some shards stop progressing, the rest keep running, and the cluster stays up. For our case that is strictly better — a projection falling behind is an internal problem, a fleet that cannot stay up is a customer-facing one.
Two smaller things that would have helped independently:
Willing to help
Happy to prototype this against a fork if the direction is agreeable — we have a reproduction with real scale (8,792 agents, 512 databases) and the production measurements above to validate against.