Paged Streaming

Paged streaming provides cursor-based pagination with automatic connection release between pages. Unlike queryStream which holds a single connection open for the entire stream, paged streaming fetches data in batches and releases the connection after each batch. This makes it ideal for:

These examples share a seedEvents helper that creates the table and inserts five rows, so each example below can populate its own data before running its paged query.

pagedStream - Cursor-Based Pagination

Use pagedStream to get a stream of Page[A, K] objects, each containing a batch of items and a cursor for checkpointing:

xa.run(for
  _ <- seedEvents
  // Fetch in pages of 2
  pages <- Query[PagedEvent].all
    .pagedStream(_.id, pageSize = 2)
    .runCollect
yield pages.map(p => s"Page ${p.pageNumber}: ${p.items.size} items, cursor=${p.cursor}"))
  .either
Right(Chunk(Page 0: 2 items, cursor=Some(2),Page 1: 2 items, cursor=Some(4),Page 2: 1 items, cursor=Some(5)))

Each Page[A, K] contains:

  • items: Chunk[A] - The rows in this page

  • cursor: Option[K] - Last key on this page (use with startAfter to resume)

  • pageNumber: Int - 0-indexed page number

The stream ends after a partial page (fewer than pageSize items) or an empty fetch.

Checkpointing for Resumable Processing

The cursor in each page enables resumable processing - save the cursor, and if processing fails, resume from where you left off:

xa.run(for
  _ <- seedEvents
  // Simulate processing with checkpoint storage
  checkpoint <- Ref.make(Option.empty[Int])

  // Process pages and save cursor after each
  _ <- Query[PagedEvent].all
    .pagedStream(_.id, pageSize = 2)
    .tap(page => checkpoint.set(page.cursor))
    .runDrain

  // Get final checkpoint
  finalCursor <- checkpoint.get
yield s"Final cursor: $finalCursor")
  .either
Right(Final cursor: Some(10))

Resuming from a Checkpoint

Use startAfter to resume processing from a saved cursor:

xa.run(for
  _ <- seedEvents
  // First, process first page only
  firstPage <- Query[PagedEvent].all
    .pagedStream(_.id, pageSize = 2)
    .runHead

  // Save the cursor
  savedCursor = firstPage.flatMap(_.cursor)

  // Later, resume from the checkpoint
  remaining <- Query[PagedEvent].all
    .pagedStream(_.id, pageSize = 2, startAfter = savedCursor)
    .runCollect
yield (s"Saved cursor: $savedCursor", s"Remaining pages: ${remaining.size}"))
  .either
Right((Saved cursor: Some(12),Remaining pages: 2))

seekingStream - Row-by-Row with Batched Fetching

If you don't need page metadata, use seekingStream to get individual items while still releasing connections between batches:

xa.run(for
  _      <- seedEvents
  result <- Query[PagedEvent].all
    .seekingStream(_.id, batchSize = 2)
    .runCollect
    .map(items => s"Got ${items.size} items")
yield result)
  .either
Right(Got 5 items)

Flattening Pages to Items

Use the .items extension to flatten a paged stream to individual items:

xa.run(for
  _      <- seedEvents
  result <- Query[PagedEvent].all
    .pagedStream(_.id, pageSize = 2)
    .items // Flatten to individual items
    .runCollect
    .map(items => s"Got ${items.size} items")
yield result)
  .either
Right(Got 5 items)

Items with Checkpoint Callback

Use .itemsWithCheckpoint to process items individually while receiving a callback after each page completes:

xa.run(for
  _           <- seedEvents
  checkpoints <- Ref.make(Chunk.empty[Option[Int]])

  items <- Query[PagedEvent].all
    .pagedStream(_.id, pageSize = 2)
    .itemsWithCheckpoint(cursor => checkpoints.update(_ :+ cursor))
    .runCollect

  recorded <- checkpoints.get
yield (s"Processed ${items.size} items", s"Checkpoints: ${recorded.toList}"))
  .either
Right((Processed 5 items,Checkpoints: List(Some(27), Some(29), Some(30))))

pagedStream vs queryStream

FeaturequeryStreampagedStream / seekingStream
Connection usageHeld for entire streamReleased between pages
Memory per pageSingle row bufferFull page buffered
Checkpoint supportNoYes (cursor in each page)
Best forReal-time, low-latencyLong-running, resumable jobs
Early terminationConnection releasedConnection already released

Combining with WHERE Clauses

Paged streaming works with all query builder methods:

xa.run(for
  _      <- seedEvents
  result <- Query[PagedEvent]
    .where(_.processed)
    .eq(false)
    .pagedStream(_.id, pageSize = 10)
    .items
    .runCollect
    .map(items => s"Unprocessed: ${items.size}")
yield result)
  .either
Right(Unprocessed: 5)

Resource Safety

Paged streams release the database connection after fetching each page. This prevents connection pool exhaustion during long-running operations. The program below runs a full checkpointed pipeline against the database:

(xa.run(for
  _ <- ddl.createTable[LargeRow](ifNotExists = true)
  _ <- ddl.truncateTable[LargeRow]()
  _ <- ZIO.foreachDiscard(1 to 5)(i => dml.insert(LargeRow(-1, s"row-$i")))
  // Process rows in pages of 2, checkpointing after each page; the connection
  // is released between pages so a long job never exhausts the pool.
  _ <- Query[LargeRow].all
    .pagedStream(_.id, pageSize = 2)
    .itemsWithCheckpoint(cursor => Console.printLine(s"Checkpoint: $cursor").orDie)
    .foreach(processRow)
yield ()) *> Console.printLine("Processed all rows")).either
Right(())