skip to content

In a search engine sharded across 50 index shards, why does jumping to results page 1,000 cost far more than page 1?

level: middleimportance: should knowfreq 45%

answer

  1. any shard can own the page
  2. offset plus size per shard
  3. multiply by shard count
  4. resume after a sort key
  5. unique tiebreaker in the key

basics

~20 s

Any shard could supply results 9,991-10,000, so each shard must return its own top 10,000, and the coordinator merges 500,000 entries to keep 10. Cursor paging avoids this by resuming after the last hit's sort key.

solid answer

~50 s

With offset paging, page `p` of size `s` needs results after `offset = (p − 1)·s`. The coordinator cannot know which shard holds hit number 9,991, so every shard returns its top `offset + s` hits and the coordinator merges `S·(offset + s)` of them. On 50 shards at 10 per page, page 1 merges 500 entries and page 1,000 merges 500,000. Memory, CPU and network all grow linearly with both depth and shard count, which is why systems cap the offset. **Cursor (keyset) paging** sorts by a total order, such as score and then a unique document ID, and hands back the last hit's sort values. Each shard then returns only its next `s` hits after that key, so each page costs `S·s` entries at any depth. The trade-off is that users cannot jump straight to page N.

code

pseudocode · 7 lines
pseudocode
function nextPage(query, size, after):     // after = null for page 1
  lists = parallel for shard in shards:
            shard.topK(query, size, strictlyAfter = after)
  hits = mergeByKey(lists)                 // key = (score desc, docId asc)
  page = first size entries of hits
  cursor = page.size == size ? sortKey(last of page) : null
  return (page, cursor)

go deeper

for a junior

Recall that deeper pages cost more in a sharded search, and that the usual remedies are a depth cap and cursor paging, which resumes after the last hit.

for a middle

Explain why each shard must return offset plus size hits and compute the coordinator's merge size for a given page and shard count. Then describe how a cursor with a unique tiebreaker keeps each page's cost constant.

for a senior

Show production judgment: spot deep-paging scrapers in the load profile, set a result-window cap, route bulk consumers to a cursor or export path, and handle result shifts when the index refreshes between pages.

for a principal

Decide per product surface whether random page access is worth its linear cost. Weigh snapshot retention for stable paging against its storage cost.

## Offset pagination over shards With **offset paging**, the client asks for page `p` of size `s`, which means the results after `offset = (p − 1)·s`. Even on a single index this is wasteful: the engine must rank `offset + s` hits only to throw away the first `offset`. Over `S` shards the waste multiplies. The coordinator cannot know in advance which shard holds result number 9,991, and any shard might supply every hit on that page. So each shard returns its own top `offset + s` hits, and the coordinator merges `S·(offset + s)` entries and keeps the last `s`. ## The arithmetic Assume 50 shards and 10 hits per page: | Page | Hits each shard returns | Entries the coordinator merges | |---|---|---| | 1 | 10 | 500 | | 100 | 1,000 | 50,000 | | 1,000 | 10,000 | 500,000 | Showing 10 results on page 1,000 means ranking and shipping half a million candidates. Per-shard heaps, network payloads and coordinator memory all grow linearly with depth **and** with shard count. A few crawlers or scripts walking deep pages can put more load on the cluster than thousands of users reading page 1. ## Why systems cap the depth Most search services protect themselves with one or more of these: - a **maximum result window** on `offset + s`, with deeper requests rejected; - a UI that shows only the first N pages; - a separate cursor or export path for bulk consumers. A cap is a safety valve, not a fix. It removes the cost by removing the access. ## Cursor (keyset) paging Cursor paging replaces the numeric offset with a position in a **total order**: 1. Sort by a key that is unique for every hit. For relevance ranking that is `(score descending, document ID ascending)`. 2. Return the page together with a **cursor**, which holds the sort values of the page's last hit. 3. For the next page, the client sends that cursor back. 4. Each shard returns only its next `s` hits that sort strictly after the cursor. 5. The coordinator merges `S·s` entries, which is 500 in the example above, whatever the depth. ```pseudocode function nextPage(query, size, after): lists = parallel for shard in shards: shard.topK(query, size, strictlyAfter = after) hits = mergeByKey(lists) page = first size entries of hits cursor = page.size == size ? sortKey(last of page) : null return (page, cursor) ``` Each shard still has to find where the cursor falls in its ordering, but its heap holds only `s` candidates instead of `offset + s`. The network and coordinator costs stay flat. ## Consistency while paging The index can change between two page requests: documents are added or deleted, and scores shift. - With **offset paging**, one new document ranked above the current page pushes every later result down by one, so a hit is shown twice or skipped. - With **cursor paging**, a new document that sorts before the cursor is never shown, and one that sorts after it appears in its proper place. Nothing is repeated. - For fully stable paging, some systems pin a **snapshot** of the index for the session. The price is keeping old index files alive until the snapshot expires. The unique tiebreaker matters. If the cursor held only the score, every hit whose score equals the cursor's would be either repeated or skipped at the page boundary. ## Choosing an approach | Need | Approach | |---|---| | Interactive UI, first few pages | offset paging with a depth cap | | Infinite scroll or a Next button only | cursor paging | | Bulk export of every match | cursor over a pinned snapshot, or read from the source of truth | | Jump to an arbitrary page N | offset within the cap; not supported beyond it | Cursor paging gives up random access in exchange for a constant cost per page. Few users go past the first pages, so that is usually the right trade.

  • In cursor paging over a search index, why must the sort key include a unique document ID?
    Relevance scores often tie. If the cursor held only the last score, the next request would have to choose between strictly lower scores, which skips any tied hits not yet shown, and lower-or-equal scores, which repeats the tied hits already shown. Adding a unique document ID as a tiebreaker makes the order total, so every hit has exactly one position and each page starts right after the previous one.
  • How would you export all 20 million documents matching a query from a sharded search index?
    Don't page through them with offsets: the depth cap blocks it, and the merge cost would grow without bound. Use a cursor over a pinned snapshot so each batch costs the same and results don't shift. Because an export needs every match and no ranking, sorting by document ID is enough. If the index is only a copy of a source of truth, reading from that source is often simpler and cheaper.

saying these in an interview costs you the question

  • Each shard only needs to return the 10 hits for the requested page.
  • Page 1,000 costs the same as page 1 because only 10 hits come back.
  • Cursor paging lets users jump directly to any page number.
  • Sorting by score alone is a safe cursor key.
  • Caching page 1 solves the cost of deep pagination.