Repository navigation
Reduce processing-pool queue-lock contention with a sharded task queue (druid.processing.numThreadPools) #19937
Description
Activity
I'm wondering do you have many "small" segments in your historical nodes? and what is the QPS on your historical nodes?
A ShardedPrioritizedExecutorService composite holding N ordinary PrioritizedExecutorService instances; it delegates per-task work
(execute/submit) to a random shard and fans out lifecycle calls.I think this proposal changes the scheduling semantics of the processing pool significantly.
Today, all queued tasks share one priority queue, so priority ordering is global within the process. With the proposed sharding, priority ordering would become local to each shard. As a result, a lower-priority task could start while a higher-priority task is waiting in another shard.
Similarly, when druid.processing.fifo=true and two queued tasks have the same priority, they are currently dequeued in global FIFO order. With sharding, FIFO ordering would apply only within each shard, so a later task in one shard could be dequeued before an earlier task in another.
We should each an agreement on such semantic changes before discussing more about the proposal.
Thanks @FrankChen021.
I'm wondering do you have many "small" segments in your historical nodes? and what is the QPS on your historical nodes?
Yes — on this cluster we create segments every minute, historical QPS is ~4–5k, and it's almost entirely 1-min and 3-min-interval groupBy queries. Because the segments are small and per-minute, each query fans out over many segments, so the per-segment task rate through the processing pool's single queue is very high — that's the crux: we were running numThreads = num_cores − 1 (55 here), but the single queue lock, not CPU, was the ceiling, so we couldn't actually use the available cores.
Scaling out horizontally does relieve it, but it increases query fan-out and therefore merge overhead on the brokers — so it trades a historical-side bottleneck for a broker-side one, and there's a practical limit to how far we can push it. That's the motivation for this proposal.
On the semantics — you flagged the trade-off accurately. To keep it safe:
Default numThreadPools=1, which is identical to today's behavior (one queue, global priority + global FIFO). Nothing changes unless an operator opts in.
We can document the trade-off explicitly on the config, and call out exactly when it helps — i.e. when numThreads and/or the number of segments scanned per node is high enough that the single queue lock becomes the bottleneck.It's also worth noting that priority/FIFO ordering in the processing pool is already a per-process property, not cluster-global — a query fans out to many data servers, and each historical's pool orders independently, so there is no global ordering across the fleet today. The number of independent ordering domains already grows every time we scale out . NumThreadPools applies the same idea inside a process, so it lets us relieve this contention without adding nodes and the broker-side fan-out cost that comes with them.
(To be precise about the one difference: scaling out keeps each node's single queue globally ordered, just with fewer tasks, whereas numThreadPools relaxes ordering within a process. So it isn't strictly identical — but since ordering is already only process-local, and both approaches simply make each independent queue smaller, this feels like a modest, opt-in step in a direction operators already rely on.)
what is the use case that you create segment every minute? and what is the size of each segment?
what is the use case that you create segment every minute? and what is the size of each segment?
The segment size is 50-100MB.
- We support the metrics scraping use case from this cluster. (near realtime queries with small intervals)
- The QPS is predictable and very high.
- Caching isn't helpful as the time intervals in the queries are non overlapping for the scraping queries
- Historically we have seen reliability issues when these queries were being answered from the indexers with (10mins of segments). So now we have minimal query load on indexers.
- The segments are of 1minute granularity and we handoff every minute. the handoff time is ~2-3s so the queries land on the historicals
Description
Add an option to split the historical/peon processing pool into N independent
pools ("shards"), each with its own
PriorityBlockingQueueand its own lock,instead of the current single queue guarded by one
ReentrantLock. Each task isrouted to a shard at submit time, so each queue lock sees only ~1/N of the
submit/take traffic and ~1/N of the contending worker threads.
New config:
druid.processing.numThreadPools(default1).druid.processing.numThreadsis split as evenly as possible across the pools(remainder to the first pools). Total thread count and processing-buffer sizing
are unchanged, so direct-memory accounting is unaffected.
ThreadLocalRandom— no shared counter, so the routeradds no contention of its own.
numThreadPools = 1(the default) preserves the current single-pool behaviorexactly; no config change is required and no metric names change.
Implementation sketch:
ShardedPrioritizedExecutorServicecomposite holding N ordinaryPrioritizedExecutorServiceinstances; it delegates per-task work(
execute/submit) to a random shard and fans out lifecycle calls.ProcessingPoolStatsinterface (implemented by both the single-pooland sharded executors) so
segment/scan/pendingandsegment/scan/activekeep being emitted — summed across shards in the sharded case.
numThreadPools > 1.Trade-off: priority ordering becomes per-shard rather than global — a
high-priority task in one shard does not preempt work queued in another. For the
high-throughput, effectively-single-priority per-segment workload this pool
serves, that is an acceptable exchange. Deployments that need strict global
priority ordering should keep the default (
numThreadPools = 1).Motivation
On historicals (and any process using the processing pool), every per-segment
scan/merge task is submitted to a single
PrioritizedExecutorServicebacked byone
PriorityBlockingQueue. Because a priority queue is a binary heap, it uses asingle
ReentrantLockfor bothputandtake(unlikeLinkedBlockingQueue'stwo-lock design), so every task pays the lock twice — once when a producer
enqueues it and once when a worker dequeues it.
On workloads with very high task rates — many small segments scanned per query at
high QPS — this single lock becomes the bottleneck:
numThreadsworker threads plus all producers serialize on one lock.AbstractQueuedSynchronizerpark/unpark (futex + context switch) and by cache-line bouncing on the AQS
stateword across cores, which grows worse-than-linearly with the number ofcontending threads.
is spent parked, off-CPU), and it gets worse on larger many-core nodes
because more worker threads contend on the same lock.
In our deployment, once other hotspots were removed, lock profiling on a historical
showed this processing-queue
ReentrantLockas the dominant remaininglock-contention source, with the pool taking tasks at ~170–180k lock
acquisitions/sec through one queue. Reducing the processing-thread count per node
(and scaling out) measurably cut the contention, confirming the single queue —
not CPU or merge buffers — was the limiter. Sharding the queue removes the
bottleneck directly instead of working around it with node capacity, and lets a
single many-core node use its threads without collapsing onto one lock.