Visitar URL original
Selection with ORDER BY on a sorted column projects every output column for the whole match set · Issue #19761 · apache/pinot · GitHub
Skip to content

Selection with ORDER BY on a sorted column projects every output column for the whole match set #19761

Description

@arunkumarucet

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:

  1. In the sortedColumnsPrefixSize > 0 branch, when there are output expressions beyond the order-by ones, project only the order-by expressions.
  2. Have LinearSelectionOrderByOperator emit rows as [sortKeys..., docId].
  3. 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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions