Summary
SELECT <wide projection> ... ORDER BY <sorted column> DESC LIMIT k reads every selected column for every document in the match set, then keeps k rows. Cost scales with match-set size times projection width, not with LIMIT.
Pinot already avoids this on the unsorted order-by path — SelectionPlanNode projects only the order-by expressions there and SelectionOrderByOperator fetches the rest for the top-k doc ids. The sorted-column path, which is otherwise the faster one, does not.
Measurements
Apache Pinot 1.6.0-SNAPSHOT, single node, 100M rows of OpenTelemetry logs, 24 segments, Lucene text index on Body. Timestamp has column.Timestamp.isSorted = true. Same filter throughout, varying only the projection.
| Order-by column |
Projected columns |
numDocsScanned |
numEntriesScannedPostFilter |
Returns |
Timestamp (sorted) |
1 |
1,778,209 |
1,778,209 |
100 rows |
Timestamp (sorted) |
4 |
1,778,209 |
7,112,836 |
100 rows |
SeverityNumber (unsorted) |
1 |
2,317,628 |
2,317,628 |
100 rows |
SeverityNumber (unsorted) |
4 |
2,317,628 |
2,320,028 |
100 rows |
On the sorted path, 7,112,836 = 4 x 1,778,209 exactly — all four columns are read for every scanned document.
On the unsorted path, 2,320,028 = 2,317,628 + 2,400, i.e. the sort column for every document plus 100 rows x 24 segments for the late fetch. That is the behaviour the sorted path should have.
Latency, same query, varying only the projection:
| Projected |
Server time |
Timestamp |
16 ms |
+ SeverityText (dictionary) |
22 ms |
+ ServiceName (dictionary) |
27 ms |
+ TraceId (RAW, 32 chars) |
44 ms |
+ Body (RAW, long log lines) |
108 ms |
The penalty is worst for RAW columns, which is the common shape for a log body. Dropping the ORDER BY entirely costs 21 ms with the same four columns, so sorting is not the cost.
The two optimizations are orthogonal
The sorted path is doing something valuable: numSegmentsMatched was 6 of 24, so it prunes most segments. Within the matched segments it then scans the whole match set with the full projection. So this is not a case of choosing between the linear scan and late materialization — the fix is to combine them, and routing these queries to SelectionOrderByOperator instead would trade one win for another.
Where it is
SelectionPlanNode#run(), the sortedColumnsPrefixSize > 0 branch, calls getSortedByProject(expressions, ...) with the full output expression list. The LinearSelectionOrderByOperator subclasses then build complete rows from each block. SelectionPartiallyOrderedByDescOperation#fetch has a comment describing the related constraint:
// Ideally we would use a descending cursor, but we don't actually have them
// ...
// The only alternative we have right now is to retrieve the last LIMIT elements from each block
Suggested fix
Apply the mechanism that already exists in SelectionOrderByOperator (roughly lines 264-340) to the sorted path:
- In the
sortedColumnsPrefixSize > 0 branch, when there are output expressions beyond the order-by ones, project only the order-by expressions.
- Have
LinearSelectionOrderByOperator emit rows as [sortKeys..., docId].
- After the top-k is selected, fetch the remaining expressions for those doc ids with
BitmapDocIdSetOperator.ascending(docIds, numRows) -> ProjectionOperator -> TransformOperator, as SelectionOrderByOperator does, including the null-bitmap and data-schema handling.
Expected on the measurement above: 7,112,836 -> about 1,778,209 post-filter entries, with Body read ~2,400 times instead of ~1.78M.
Non-trivial parts: _columnContexts is currently resolved from the project operator, which would no longer contain the output expressions; createDataSchema and both subclasses' fetch() need to move with it.
Reproduction
-- wide projection, sorted order-by column: entriesPostFilter = numDocsScanned x numColumns
SELECT "Timestamp", ServiceName, SeverityText, Body
FROM otel_logs WHERE TEXT_MATCH(Body, 'error')
ORDER BY "Timestamp" DESC LIMIT 100;
-- same shape on an unsorted order-by column: entriesPostFilter = numDocsScanned + (limit x segments)
SELECT "Timestamp", ServiceName, SeverityText, Body
FROM otel_logs WHERE TEXT_MATCH(Body, 'error')
ORDER BY SeverityNumber DESC LIMIT 100;
Compare numEntriesScannedPostFilter in the two responses.
Why it shows up
This is the dominant cost of "show me the N most recent matching log lines", which is the central query shape in observability workloads, and it gets worse as the log body column gets wider.
Summary
SELECT <wide projection> ... ORDER BY <sorted column> DESC LIMIT kreads every selected column for every document in the match set, then keepskrows. Cost scales with match-set size times projection width, not withLIMIT.Pinot already avoids this on the unsorted order-by path —
SelectionPlanNodeprojects only the order-by expressions there andSelectionOrderByOperatorfetches the rest for the top-k doc ids. The sorted-column path, which is otherwise the faster one, does not.Measurements
Apache Pinot 1.6.0-SNAPSHOT, single node, 100M rows of OpenTelemetry logs, 24 segments, Lucene text index on
Body.Timestamphascolumn.Timestamp.isSorted = true. Same filter throughout, varying only the projection.numDocsScannednumEntriesScannedPostFilterTimestamp(sorted)Timestamp(sorted)SeverityNumber(unsorted)SeverityNumber(unsorted)On the sorted path,
7,112,836 = 4 x 1,778,209exactly — all four columns are read for every scanned document.On the unsorted path,
2,320,028 = 2,317,628 + 2,400, i.e. the sort column for every document plus100 rows x 24 segmentsfor the late fetch. That is the behaviour the sorted path should have.Latency, same query, varying only the projection:
Timestamp+ SeverityText(dictionary)+ ServiceName(dictionary)+ TraceId(RAW, 32 chars)+ Body(RAW, long log lines)The penalty is worst for RAW columns, which is the common shape for a log body. Dropping the
ORDER BYentirely costs 21 ms with the same four columns, so sorting is not the cost.The two optimizations are orthogonal
The sorted path is doing something valuable:
numSegmentsMatchedwas 6 of 24, so it prunes most segments. Within the matched segments it then scans the whole match set with the full projection. So this is not a case of choosing between the linear scan and late materialization — the fix is to combine them, and routing these queries toSelectionOrderByOperatorinstead would trade one win for another.Where it is
SelectionPlanNode#run(), thesortedColumnsPrefixSize > 0branch, callsgetSortedByProject(expressions, ...)with the full output expression list. TheLinearSelectionOrderByOperatorsubclasses then build complete rows from each block.SelectionPartiallyOrderedByDescOperation#fetchhas a comment describing the related constraint:Suggested fix
Apply the mechanism that already exists in
SelectionOrderByOperator(roughly lines 264-340) to the sorted path:sortedColumnsPrefixSize > 0branch, when there are output expressions beyond the order-by ones, project only the order-by expressions.LinearSelectionOrderByOperatoremit rows as[sortKeys..., docId].BitmapDocIdSetOperator.ascending(docIds, numRows)->ProjectionOperator->TransformOperator, asSelectionOrderByOperatordoes, including the null-bitmap and data-schema handling.Expected on the measurement above: 7,112,836 -> about 1,778,209 post-filter entries, with
Bodyread ~2,400 times instead of ~1.78M.Non-trivial parts:
_columnContextsis currently resolved from the project operator, which would no longer contain the output expressions;createDataSchemaand both subclasses'fetch()need to move with it.Reproduction
Compare
numEntriesScannedPostFilterin the two responses.Why it shows up
This is the dominant cost of "show me the N most recent matching log lines", which is the central query shape in observability workloads, and it gets worse as the log body column gets wider.