A federation layer serves 60 analysts across four sources: the lakehouse, a Postgres replica, a managed warehouse and a SaaS REST connector. At 11:40 the REST connector's upstream begins answering in 40 seconds rather than 200 ms. Nothing errors and nothing is down. Walk through what happens to the other analysts over the next ten minutes.
Show the full answer Hide the answer
Minute by minute
11:40. Queries touching the REST catalogue stop returning. Each one holds a worker slot per split, and a split waiting on a slow source holds that slot for the whole 40 seconds. Nothing is logged, because a slow response is a successful response.
11:42. Little's Law does the damage. At two REST-touching queries a second, concurrency demand is arrival rate times service time: 2 × 0.2 s is 0.4 slots at normal speed, and 2 × 40 s is 80 slots. Against a pool of 200 worker threads, those queries now hold roughly 40% of the cluster.
11:45. Analysts whose queries never touch the REST source start queueing. This is the part people do not predict: admission is first-in-first-out against one shared queue, so a lakehouse query that would run in 3 seconds waits behind stalled work it has no relationship with.
11:48. Memory, not slots, produces the first errors. Each stalled split holds its partially built hash tables and exchange buffers, so the cluster's memory pool fills and the engine starts spilling or killing queries with an out-of-memory message that names an innocent query. The on-call engineer is now investigating the wrong thing.
What the user sees
Analysts see a mix of "still running" and memory failures on queries that worked an hour ago, distributed across all four sources. The one source actually at fault is the one nobody mentions, because its queries have not failed, they have merely not finished.
What stops it
- A per-catalogue concurrency cap. Limit the REST catalogue to, say, 20 concurrent splits. Its queries queue among themselves and everyone else is untouched. This is a bulkhead, and it is one configuration line.
- A per-source timeout well below the query timeout. If the source normally answers in 200 ms, a 5-second cap converts a stall into an error, and an error is something the engine and the human can act on.
- Resource groups with a memory share per group, so a stalled group cannot consume the pool that other groups need.
- Not more workers. Doubling the pool to 400 threads is absorbed by the same source in about a minute, at double the cost.
What would have to be true for it to self-heal
Only a fall in offered load to that source, which will not happen: analysts whose queries hang cancel and re-run, so the arrival rate against the slow source rises during the incident. A federated query engine has no natural damping, which is why the concurrency cap has to be static and set in advance rather than reactive.
When this is the wrong answer
If the slow source is the one that matters rather than a side connector, isolation only converts stalling into rejection, and analysts get fast failures instead of slow results. At that point the honest move is to stop federating it and land a copy on a schedule the source can sustain. Federation is the right pattern when the remote source is small, occasional and not on the critical path for most queries; the moment it becomes the majority of the workload, the absence of a copy is the problem, not the concurrency setting.
Little's Law dates from 1961 and settles this argument faster than a benchmark: latency times arrival rate is concurrency, and concurrency you have not capped is concurrency you have given away.