Do we have any plan on extending query lane to His...
# dev
m
Do we have any plan on extending query lane to Historical and realtime ingestion task (i.e. Peons)? Having query lane only on the broker makes it less useful. For example, a single very expensive query (which would run on low lane on broker) could/would still takes up all the processing thread on historical and blocks other query from running on those historical. CC: @Clint Wylie (as author of the query lane on the broker). Thanks!!
g
are you wanting low-priority queries to not use the entire processing thread pool? (maybe restrict to 8 threads on a 32 thread machine)?
the original idea was that on historicals, lanes aren't needed b/c the main resource we're protecting is the processing thread pool, and that's adequately protected by
priority
. a low priority query can use the entire pool, but if a higher priority query comes along, then it will take over the pool
lane
is a concept that IMO doesn't fully make sense on Historicals; its main purpose is to control which queries are accepted for running, but by the time a query gets to a Historical, it's already been accepted (by the Broker that sent it)
it would be weird i think for the Broker to send a query that the Historical then decides not to run
m
@Gian Merlino Thanks for the reply.
are you wanting low-priority queries to not use the entire processing thread pool? (maybe restrict to 8 threads on a 32 thread machine)?
Yes. Exactly.
and that’s adequately protected by
priority
. a low priority query can use the entire pool, but if a higher priority query comes along, then it will take over the pool
That’s not entirely true though as Druid historical does not preempt already running processing threads. For example, if you set query timeout to 5 minutes. Most of your queries are sub-second. In this example, our system only have query A and B. Query A always run for 100ms. Query B is an expensive query that will run for >5minutes (and timeout) was put in low lane. Query B ran run before query A and did have a lane capacity on the low lane. Broker ran Query B. Query B then ran on Historical taking up all available processing threads. Each processing thread will be taken up for 5 minutes. Query A was execute right after query b but now has to wait for 5 minutes for processing threads on historicals. Users of query A then file a ticket because a query that is suppose to run in 100ms is not taking 5minutes/timeout.
lane
is a concept that IMO doesn’t fully make sense on Historicals; its main purpose is to control which queries are accepted for running, but by the time a query gets to a Historical, it’s already been accepted (by the Broker that sent it)
If lane is a concept to protect/reserve resource. Why don’t we do it all the way? i.e. if we want to guarantee that a high priority query will always have resource reserved for it, then I think it make sense to apply it to historical as well.
c
hmm, usually long running queries are long running because they must process a lot of segments on the same historical, not that processing a single segment takes that long
it is true that we don't preempt running threads, but higher priority query segments should be processed before lower priority, even if added to the queue later
m
On a related note, do we have any plan on collecting query stats to determine how expensive a query is? (i.e. looking at number of rows, type of aggregators, column cardinality, if a query is vectorizable or not, etc). This can be use to automatically assign query to lanes. One issue we are having with lanes is on how to assign query to lanes. We first started with Manual prioritization strategy. But then everyone thinks there query is high priority and every query is assigned to high lane (even the very expensive ones). So we tried Threshold prioritization strategy but it is very hard to get the Threshold right. Especially when datasources and query can be mixed workload. A query that is querying few number of segments/time intervals but is doing case search expression, lots of datasketches, high cardinality group by, cannot be vectorized, etc can be more expensive than a query that is querying lots number of segments/time intervals but is simply doing a MAX(x) for example.
@Clint Wylie Do you have any insights on which Threshold for Threshold prioritization strategy works best and/or how to best set the Threshold value? 😵‍💫
c
i don’t have an immediate plans to work on any of that stuff, but I think a lot of things could probably be improved in ways for various levels of effort
yea, the thresholds are quite primitive right now
i don’t think we have easy access to row counts of segments at broker (and depending on filters the actual rows scanned could be a very small portion of the total row count, and filter selectivity would be very hard to know without also knowing a lot more stuff)
but maybe the next level of sophistication of something like the threshold strategy would be like something that looks at segment count and then modifies it in some way based on the actual query
like maybe multiply by
max(1, someQueryCostModifier)
it would still require knowing some stuff about your typical workloads though i would think
but i haven’t thought very much about how a cluster might automatically decide what is heavy to automatically deprioritize vs what is fine, since it would depend both on loaded segment composition and also maybe typical workloads?
a tricky thing to balance here is that deprioritized heavy queries can also become so starved and increase the likelyhood of timing out if there are too many interactive queries because they spend all of their time just getting skipped in line in the queue waiting for segments to be processed
i did some experiments with preempting processing slots once and that was exaggerated even further.. basically had to like periodically increase the low priority tasks or else they never complete because they are always preempted if the cluster is busy enough
i think probably longer term a better behavior would be if there was some way for an interactive query to be like demoted to an async query for being too heavy
but not sure exactly how to best do that either
and we’re probably missing some pieces right now for that to be fully realized
g
> That’s not entirely true though as Druid historical does not preempt already running processing threads. For example, if you set query timeout to 5 minutes. Most of your queries are sub-second. In this example, our system only have query A and B. Query A always run for 100ms. Query B is an expensive query that will run for >5minutes (and timeout) was put in low lane. Query B ran run before query A and did have a lane capacity on the low lane. Broker ran Query B. Query B then ran on Historical taking up all available processing threads. Each processing thread will be taken up for 5 minutes. Query A was execute right after query b but now has to wait for 5 minutes for processing threads on historicals. Users of query A then file a ticket because a query that is suppose to run in 100ms is not taking 5minutes/timeout. This isn't quite how it works with processing threads Each processing thread is held for however long it takes to process one segment, and is released back to the pool after that. When it goes back to the pool, another query can pick it up next, and at that time a higher priority query can take over This generally means that as long as two queries have different priorities, the highest priority one will always "take over" relatively quickly at the Historical level I have only really seen one situation where this doesn't work well— if processing of a single segment takes a very long time for some reason. Sometimes this can happen, although it's rare in my experience for it to be more than a few seconds. So, in this situation, it could take a few seconds for the higher priority query to "take over". But it shouldn't take minutes If it does take minutes something seems "wrong" to me
Anyway, something could definitely be implemented that limits a query from using more than X number of threads on a Historical. It does seem rare that this would actually be useful, but it could happen (it would need to be in the situation where processing a single segment takes a very long time)
Btw, about this problem of not being able to figure out the "correct" prioritization or lane for a query. I've been thinking about that recently too. I think the right thing to do is quotas + dynamic reprioritization. I'm thinking we should do this for Dart (to me Dart is the future 🙂)
Here's something from a draft I'm working on of a proposal. I probably won't post it on Github for a while (still working on it) but seems relevant, so wondering what you think—
It is difficult-to-impossible for either Druid, or the application querying it, to know ahead of time what priority a query should have. The best approach is typically to use very simple heuristics such as “queries that process more segments should be lower priority” or “queries for the ‘download’ feature should be lower priority”. But the amount of time to process a segment varies greatly based on things like value cardinality, filter selectivity, complexity of expressions, etc. This limits the success that people have with prioritization. To improve this, we want to:
• implement a quota system that tracks CPU time used by ‘quota-carrying entities’. Such entities may be users (i.e. identity as determined by an authenticator) or may be groups of users.
• prioritize incoming queries based on recent CPU usage by the relevant ‘quota-carrying entity’.
• reprioritize currently-running queries, which may necessitate canceling and re-running them, if the new priority causes them to switch to a lower lane that is currently full.
b
two things I've missed ever since working w/Oracle DBs: 1. resource quotas (cpu would be great, as would memory usage, etc) and 2. ability to find and kill long-running resource-hog queries. So that's a $0.02 vote for the above idea for dart. 🙂
m
I have only really seen one situation where this doesn’t work well— if processing of a single segment takes a very long time for some reason. Sometimes this can happen, although it’s rare in my experience for it to be more than a few seconds. So, in this situation, it could take a few seconds for the higher priority query to “take over”. But it shouldn’t take minutes
If it does take minutes something seems “wrong” to me
We do have some queries that takes on avg 1-2 minutes to process a single segment (attached metric on P50, P90, P99
query/segment/time
) Maybe the query is the problem but we provide Druid as a platform and allow our end user to ingest whatever they want (they write the ingestionSpec) and query whatever they want (they write the query)…and sometimes those query are not the fastest/most efficient 😞
In this case, the query is something like:
Copy code
SELECT
 DATE_TRUNC('day', __time) "timestamp",
 is_x,
 1000.0 * SUM(CASE WHEN foo IN ('a','b') THEN 0 ELSE CASE WHEN is_x = False THEN bar / 0.9 ELSE baz END END) / SUM(cnt) FILTER (WHERE is_y != 'true') "value"
FROM ds
WHERE __time BETWEEN '2024-01-01T00:00:00.000Z' AND '2025-04-21T23:59:59.999Z'
GROUP BY 1,2 ORDER BY 1,2
LIMIT 2000000
The time range contains about 20,000 segments, each segment about 3.5-5.5 million rows 150MB-250MB. Note that removing/rewriting without the CASE fixes the slowness in the query but the lack of guardrail / resource management on the historical resulted in this query taking over the historical for ~2minutes (blocking everything else).
implement a quota system
This sounds reasonable to me. I guess similar to like Hadoop YARN Fair Scheduler.
a tricky thing to balance here is that deprioritized heavy queries can also become so starved and increase the likelyhood of timing out if there are too many interactive queries because they spend all of their time just getting skipped in line in the queue waiting for segments to be processed
I would argue that starving the slow queries (which usually points to badly written query, query that could be rewritten, optimized, and query that would eventually timeout anyway) is better than making all the other queries slow. One unhappy user is better than 100s of unhappy users haha
It is difficult-to-impossible for either Druid, or the application querying it
Would we be able to use segment metadata stats for this (size, cardinality, type, etc)? Maybe type of aggregators/filters? If a query can be vectorized or not? Or would those be too expensive to compute upfront?
@Gian Merlino > I’m thinking we should do this for Dart (to me Dart is the future 🙂) Do you think in the future, all queries will default to using Dart? Or will there be a smart mechanism for Druid to determine if a query should run using the standard native query engine vs Dart. If Dart is comparable to standard native query engine then maybe the future is to run everything on Dart? Would it be able to handle high QPS (cost efficiently)
g
> Do you think in the future, all queries will default to using Dart? I'd definitely like all queries to default there at some point. It may be some time before that happens. But it is something I'm intending to keep pushing towards. There is some performance work needed on the lightweight-query side Currently, Dart performance is better than native for queries with high-cardinality groupbys, large intermediate resultsets, big subqueries, etc. Native is better for lightweight / very high QPS. I have some benchmarks in my Druid Summit talk from about 11m30s to 14m30s in this video https://imply.io/videos/charting-the-future-of-druid/. I would like to get it so Dart is always better or equivalent.
🙌 2
Would we be able to use segment metadata stats for this (size, cardinality, type, etc)? Maybe type of aggregators/filters? If a query can be vectorized or not?
I think the issue is there's too many variables, especially when a query starts to involve multiple columns. Maybe AI could figure it out? Only half-joking?
We do have some queries that takes on avg 1-2 minutes to process a single segment
🤯 OMG. no wonder you have problems with the current priority system.
Do you actually need these queries to finish? I wonder if a simple solution could be having a timeout at the level of individual-segment processing. i.e., query-level timeout of 5 minutes, but segment-level timeout of 15 seconds. If it takes more than 15 sec to process a single segment then something is probably really inefficient with the query.
m
Do you actually need these queries to finish? I wonder if a simple solution could be having a timeout at the level of individual-segment processing. i.e., query-level timeout of 5 minutes, but segment-level timeout of 15 seconds. If it takes more than 15 sec to process a single segment then something is probably really inefficient with the query.
I am also thinking/agreeing that we probably don’t need these expensive queries to finish. (Similar to what I mentioned at https://apachedruidworkspace.slack.com/archives/C030CMF6B70/p1746142880974949?thread_ts=1745436989.786489&cid=C030CMF6B70). These queries will often timeout on the query-level timeout anyway. These queries are also often indication of something that a user is doing wrong (badly written query). Having these queries fail and those users optimizing/fixing the queries is much better than letting it run (and impacting all the other users/queries on the Druid cluster). Do we already have a timeout at the level of individual-segment processing? I haven’t thought of this solution but now that you mentioned it, I think it makes a lot of sense. I also think it should be pretty simple to configure and use (compare to query lane thresholds tuning, etc).
g
there is not currently a timeout at the individual-segment level
it'd be a new feature. but i think a useful one in the scenario you describe
m
^ @Jesse Tuglu
👀 1
j
@Gian Merlino there's also this old PR, that's a bit related to your recent change (albeit on native instead of Dart).
g
ah, i missed that one back when it was originally raised. it seems like a useful idea
👍 1
j
https://github.com/apache/druid/pull/18148 – I'll touch this up in a bit (not ready for review) but this is to address the above conversation
Wanted to follow-up about scan/window queries @Gian Merlino. Both of these: • Scan • Window queries run on the jetty http threadpool's thread (not on processing threads). Since they run serially for each segment and off the processing pool, standardizing a timeout for these is tough (we'd need to push the work to be on a background thread so the http thread could receive and service the timeout). While these queries do consume CPU/mem resources, they don't explicitly saturate the processing pool (which is both a good and a bad thing). For that reason, the above PR leaves these 2 types alone. Any thoughts?
Wanted to also bring this comment you made up again. What about making this non-fifo behavior "predictable"? e.g somehow having a counter per query id which would provide a total ordering over pending items on the queue? I'd need to think about some way to GC these items, but it's provide a more "predictable" pattern, and would favor queries with smaller # of segments/historical.
g
I think it's OK to leave those 2 types alone, since as you say, they don't use the processing pool, and the purpose of this PR is to protect the processing pool
👍 1
BTW, with MSQ everything uses the processing pool 😉 Dart (MSQ-on-Historicals) is going beta in Druid 34 & we'll be putting more work into it. In many ways it should be more harmonious in its usage of the processing pool. One important thing is it can also give up a thread temporarily in the midst of processing a segment, and in many cases it already does do this
About fifo vs non-fifo, it's been a while since I thought about that and most of the thoughts have left my brain in the meantime
IIRC, the major issue with non-fifo is that latencies can blow up when there's a backlog since we don't clear it in an orderly manner
j
Yeah the better way of solving this is likely just more precise prioritization at the broker-level