Skip to content

Commit 8e30197

Browse files
authored
fix: Stop dataset iterators from skipping items when unwind is used (#1059)
## What was wrong `iterate_items` advanced the dataset offset by `max(scanned_rows, len(items))`. With `unwind`, a page returns more items than the rows it scanned, so the offset jumped past rows the next request never read. At the default `chunk_size`, a 5000-row dataset with a 3x unwind yielded about a third of it. ## The fix The offset now advances by at most the number of rows the call asked for. The items endpoint applies `offset` and `limit` to the rows first and shapes the result afterwards, and it passes no `maxLimit` (unlike the collection list routes), so a page never covers more rows than the `limit` it was sent. Capping the advance there keeps an unwound page from running past the window it actually read. The `max()` stays underneath the cap on purpose. `x-apify-pagination-count` comes from the dataset's `itemCount`, which the API increments through a throttled write while the items themselves stream fresh. Following the header alone, which is what [apify-client-js#1044](apify/apify-client-js#1044) does and what #1058 suggested, would truncate `iterate_items` right after `push_items` for anything over one page. The JS iterator bounds itself by `total` and so is exposed to that lag either way; this one consults neither, which is why the two clients end up with different fixes. `DatasetItemsPage.count` keeps its current value of `max(header, len(items))`. That is not the no-op it looks like: it is what makes `count` usable right after a push, and two integration assertions depend on it. Its docstring changes, along with the `limit` docs on `iterate_items` and `chunk_size` on both twins, since all of them promised items where the field and the parameters have always counted scanned rows. One case stays imperfect. On a final short page with `unwind`, the advance is still the full `limit` rather than the rows the page covered, so rows appended by a concurrent push after that page can be missed. Reaching them needs the raw header, which `list_items` folds into `count`. It is strictly better than before, and the docstring says so rather than claiming the overshoot is free. ## Behavior change worth a changelog line `iterate_items(limit=N)` with `unwind` and a `chunk_size` below `N` now walks all N rows instead of stopping after the first page, so it can yield several times more items for the same `limit`. That is the fix working, but it is user-visible. ## Tests - `unwind` reached the pagination tests for the first time. The fake API now splits each row into three items after the row window is picked, exactly as the transform stream does, and three cases iterate through it. - The fake API used to cap every page at 1000 items "mirroring the real API", which is only true of the collection endpoints. It now applies the requested limit verbatim on the items endpoint, which is what let a `chunk_size` above that cap be covered at all. - A page whose `count` lags behind its items has a test of its own. Nothing pinned that before, and it is the invariant the `max()` rests on. - `test_dataset_iterate_items_unwound` covers the whole thing against the live API. It has not run locally, so CI is its first real execution. ## One thing found while verifying `get_cursor_iterator` justified its termination rule with filters that "can drop every item on a page while a live cursor still points at more data". The API rules that out: the key-value store returns the last key of the page as the next cursor, and the request queue returns a cursor only once a page came back full, so neither can hand back a cursor that outlives its page. The behavior is fine and unchanged; the docstring now says why it holds. Closes #1058 *✍️ Drafted by Claude Code*
1 parent f52190d commit 8e30197

5 files changed

Lines changed: 210 additions & 46 deletions

File tree

‎docs/02_concepts/08_pagination.mdx‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,8 @@ Most methods named `list` or `list_something` in the Apify client return a page
2525
- `count` - The number of items in the current page.
2626
- `limit` - The maximum number of items per page.
2727

28+
On <ApiLink to="class/DatasetItemsPage">`DatasetItemsPage`</ApiLink>, `count` reports how many dataset rows the API scanned for the page. Filters drop items from a page and `unwind` multiplies them, so a page can hold fewer or more items than `count`.
29+
2830
Some methods paginate differently. For example, <ApiLink to="class/RequestQueueClient#list_requests">`RequestQueueClient.list_requests`</ApiLink> returns a cursor-based <ApiLink to="class/ListOfRequests">`ListOfRequests`</ApiLink> without the `total`, `offset`, and `count` fields. To fetch the next page, pass its `next_cursor` value back as the `cursor` parameter. Other examples include `list_keys` and `list_head`. Regardless, the primary results are always stored under the `items` field, and the `limit` field can be used to control the number of results returned.
2931

3032
The following example shows how to fetch all items from a dataset using pagination:

‎src/apify_client/_pagination.py‎

Lines changed: 54 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,10 @@
1919
class HasItems(Protocol[T]):
2020
"""Structural contract for a single page of results from a paginated API endpoint.
2121
22-
Implementations must expose `items`. They may optionally expose `count` - the number of items scanned by the API for
23-
this page, which can exceed `len(items)` when filters drop items from the response. The iterator helpers consult
24-
`count` opportunistically via `getattr` for offset bookkeeping and fall back to `len(items)` when it is absent.
22+
Implementations must expose `items`. They may optionally expose `count` - the number of rows the API scanned to
23+
produce this page, which `len(items)` can land below (filters drop items) or above (`unwind` splits one row into
24+
several items). The iterator helpers consult `count` opportunistically via `getattr` for offset bookkeeping and
25+
fall back to `len(items)` when it is absent.
2526
"""
2627

2728
items: list[T]
@@ -38,34 +39,34 @@ def get_items_iterator(
3839
3940
The `callback` is invoked lazily to fetch each page from the API. It must accept `limit` and `offset` keyword
4041
arguments and return an object whose `items` attribute is a list. If the object also exposes a `count` attribute, it
41-
is used for offset bookkeeping (the Apify API's `count` reflects items scanned, which can exceed items returned when
42-
filters are applied).
42+
is used for offset bookkeeping - `_page_scanned_rows` describes how the next offset is derived.
4343
44-
Iteration stops when a page scans no items (`count` is `0`, or `items` is empty when `count` is absent) or when the
45-
user-requested `limit` is reached. A page can scan items while returning none - filters like `clean` drop items from
46-
`items` but still count toward `count` - so terminating on scanned rather than returned items keeps the iterator
47-
advancing across fully-filtered pages. The `total` field is intentionally not consulted, because it can change
48-
between calls.
44+
Iteration stops when a page scans no rows or when the user-requested `limit` is reached. A page can scan rows while
45+
returning no items - filters like `clean` drop items from `items` but still count toward `count` - so terminating on
46+
scanned rather than returned rows keeps the iterator advancing across fully-filtered pages. The `total` field is
47+
intentionally not consulted, because it can change between calls.
4948
5049
Args:
5150
callback: Function returning a single page of items.
52-
limit: Maximum total number of items to yield across all pages. `None` or `0` means no limit.
51+
limit: Maximum total number of rows scanned across all pages. On the dataset items endpoint `unwind` can
52+
turn one row into several items, so more items than this can be yielded. `None` or `0` means no limit.
5353
offset: Starting offset for the first page.
54-
chunk_size: Maximum number of items requested per API call. `None` or `0` lets the API decide.
54+
chunk_size: Per-page cap, sent to the API as its `limit`. `None` or `0` lets the API decide.
5555
"""
5656
effective_chunk = chunk_size or 0
5757
initial_offset = offset or 0
5858
initial_limit = limit or 0
5959
fetched_items = 0
6060

6161
while True:
62+
page_limit = _next_page_limit(initial_limit, fetched_items, effective_chunk)
6263
current_page = callback(
63-
limit=_next_page_limit(initial_limit, fetched_items, effective_chunk),
64+
limit=page_limit,
6465
offset=initial_offset + fetched_items,
6566
)
6667
yield from current_page.items
6768

68-
page_scanned = max(getattr(current_page, 'count', 0), len(current_page.items))
69+
page_scanned = _page_scanned_rows(current_page, page_limit)
6970
fetched_items += page_scanned
7071

7172
if not page_scanned or (initial_limit and fetched_items >= initial_limit):
@@ -89,14 +90,15 @@ async def get_items_iterator_async(
8990
fetched_items = 0
9091

9192
while True:
93+
page_limit = _next_page_limit(initial_limit, fetched_items, effective_chunk)
9294
current_page = await callback(
93-
limit=_next_page_limit(initial_limit, fetched_items, effective_chunk),
95+
limit=page_limit,
9496
offset=initial_offset + fetched_items,
9597
)
9698
for item in current_page.items:
9799
yield item
98100

99-
page_scanned = max(getattr(current_page, 'count', 0), len(current_page.items))
101+
page_scanned = _page_scanned_rows(current_page, page_limit)
100102
fetched_items += page_scanned
101103

102104
if not page_scanned or (initial_limit and fetched_items >= initial_limit):
@@ -126,20 +128,30 @@ def get_cursor_iterator(
126128
limit: int | None = None,
127129
chunk_size: int | None = None,
128130
) -> Iterator[KeyValueStoreKey] | Iterator[Request]:
129-
"""Yield individual items from cursor-paginated API responses.
131+
"""Yield individual items from a cursor-paginated API response.
132+
133+
This iterator supports the two API responses that use cursor pagination. `ListOfKeys` is used for key-value store
134+
keys, while `ListOfRequests` is used for request queue requests.
135+
136+
Pagination continues until either:
137+
138+
- the API returns no next cursor, or
139+
- the requested `limit` is reached.
140+
141+
An empty page does not explicitly stop the iteration. In practice, both supported endpoints return a next cursor
142+
only when the current page contains items, so an empty page always has a `None` cursor and naturally ends the
143+
iteration.
130144
131-
Cursor pagination is restricted to the two API responses that expose it: `ListOfKeys` (for key-value store keys) and
132-
`ListOfRequests` (for request queue requests). Iteration ends when the next cursor is `None` or the user-requested
133-
`limit` is reached. Emptiness alone does not stop iteration: server-side filters (such as the request-queue state
134-
`filter`) can drop every item on a page while a live cursor still points at more data, so termination relies on the
135-
cursor, not on whether a page returned items. Unlike offset responses, cursor responses expose no scanned-item
136-
`count`, so `count` cannot be used to detect a fully-filtered page here.
145+
The endpoints determine the next cursor differently:
146+
147+
- For key-value store keys, the cursor is the last key returned on the current page.
148+
- For request queue requests, a cursor is returned only when the current page is full.
137149
138150
Args:
139-
callback: Function returning a single page of items. Receives `cursor` and `limit` kwargs.
140-
cursor: Value of the cursor for the first request, or `None` to start from the beginning.
151+
callback: Function that returns one page of items and accepts `cursor` and `limit` keyword arguments.
152+
cursor: Cursor to use for the first request. If `None`, iteration starts from the beginning.
141153
limit: Maximum total number of items to yield across all pages.
142-
chunk_size: Maximum number of items requested per API call.
154+
chunk_size: Maximum number of items to request in a single API call.
143155
"""
144156
effective_chunk = chunk_size or 0
145157
initial_limit = limit or 0
@@ -218,3 +230,19 @@ def _next_page_limit(initial_limit: int, fetched_items: int, effective_chunk: in
218230
if not effective_chunk:
219231
return remaining
220232
return min(remaining, effective_chunk)
233+
234+
235+
def _page_scanned_rows(page: HasItems[T], requested_limit: int) -> int:
236+
"""Compute how far the offset advances past `page`, in dataset rows.
237+
238+
Neither reported number is right on its own. `count` follows the rows the API scanned, but it is derived from a
239+
dataset's item count, which is incremented by a throttled write and so lags a fresh push. `len(items)` counts the
240+
items the API shaped out of those rows: filters (`clean`, `skip_empty`, `skip_hidden`) drop some, and `unwind`
241+
splits one row into several. The larger of the two absorbs a `count` that lags behind the items returned, and
242+
capping it at the rows the call asked for keeps an unwound page from advancing past rows the next call would then
243+
never read. The cap is a valid bound because the endpoint applies the `limit` it is sent verbatim; on a page
244+
covering fewer rows than that, the advance can still overshoot into rows a concurrent push appends afterwards. A
245+
`requested_limit` of `0` means the call sent no limit, leaving the advance unbounded.
246+
"""
247+
scanned_rows = max(getattr(page, 'count', 0), len(page.items))
248+
return min(scanned_rows, requested_limit) if requested_limit else scanned_rows

‎src/apify_client/_resource_clients/dataset.py‎

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ class DatasetItemsPage:
4040
"""The offset of the first item in this page."""
4141

4242
count: int
43-
"""Number of items in this page."""
43+
"""Number of dataset rows the API scanned for this page, or the number of items returned when that is larger."""
4444

4545
limit: int
4646
"""The limit that was used for this request."""
@@ -204,8 +204,8 @@ def list_items(
204204
items=items,
205205
total=int(response.headers['x-apify-pagination-total']),
206206
offset=int(response.headers['x-apify-pagination-offset']),
207-
# x-apify-pagination-count returns count of processed items, not count of returned items
208-
# This makes difference when items were filtered using hidden/empty
207+
# The header counts the rows the API scanned, which `unwind` and a lagging dataset item count can
208+
# both leave below the number of items returned.
209209
count=max(int(response.headers['x-apify-pagination-count']), len(items)),
210210
# API returns 999999999999 when no limit is used
211211
limit=int(response.headers['x-apify-pagination-limit']),
@@ -237,7 +237,8 @@ def iterate_items(
237237
238238
Args:
239239
offset: Number of items that should be skipped at the start. The default value is 0.
240-
limit: Maximum number of items to return. By default there is no limit.
240+
limit: Maximum number of dataset rows to scan. Fewer items are yielded when filters drop some, more
241+
when `unwind` splits a row into several. By default there is no limit.
241242
desc: By default, results are returned in the same order as they were stored. To reverse the order,
242243
set this parameter to True.
243244
clean: If True, returns only non-empty items and skips hidden fields (i.e. fields starting with
@@ -260,7 +261,7 @@ def iterate_items(
260261
skip_hidden: If True, then hidden fields are skipped from the output, i.e. fields starting with
261262
the # character.
262263
signature: Signature used to access the items.
263-
chunk_size: Maximum number of items requested per API call when iterating across pages.
264+
chunk_size: Maximum number of dataset rows requested per API call when iterating across pages.
264265
timeout: Timeout for the API HTTP request.
265266
266267
Yields:
@@ -763,8 +764,8 @@ async def list_items(
763764
items=items,
764765
total=int(response.headers['x-apify-pagination-total']),
765766
offset=int(response.headers['x-apify-pagination-offset']),
766-
# x-apify-pagination-count returns count of processed items, not count of returned items
767-
# This makes difference when items were filtered using hidden/empty
767+
# The header counts the rows the API scanned, which `unwind` and a lagging dataset item count can
768+
# both leave below the number of items returned.
768769
count=max(int(response.headers['x-apify-pagination-count']), len(items)),
769770
# API returns 999999999999 when no limit is used
770771
limit=int(response.headers['x-apify-pagination-limit']),
@@ -796,7 +797,8 @@ def iterate_items(
796797
797798
Args:
798799
offset: Number of items that should be skipped at the start. The default value is 0.
799-
limit: Maximum number of items to return. By default there is no limit.
800+
limit: Maximum number of dataset rows to scan. Fewer items are yielded when filters drop some, more
801+
when `unwind` splits a row into several. By default there is no limit.
800802
desc: By default, results are returned in the same order as they were stored. To reverse the order,
801803
set this parameter to True.
802804
clean: If True, returns only non-empty items and skips hidden fields (i.e. fields starting with
@@ -819,7 +821,7 @@ def iterate_items(
819821
skip_hidden: If True, then hidden fields are skipped from the output, i.e. fields starting with
820822
the # character.
821823
signature: Signature used to access the items.
822-
chunk_size: Maximum number of items requested per API call when iterating across pages.
824+
chunk_size: Maximum number of dataset rows requested per API call when iterating across pages.
823825
timeout: Timeout for the API HTTP request.
824826
825827
Yields:

‎tests/integration/test_dataset.py‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -596,6 +596,47 @@ async def get_items() -> DatasetItemsPage:
596596
await maybe_await(dataset_client.delete())
597597

598598

599+
async def test_dataset_iterate_items_unwound(client: ApifyClient | ApifyClientAsync, *, is_async: bool) -> None:
600+
"""Test iterate_items with `unwind`, where a page carries more items than the rows it scanned."""
601+
dataset_name = get_random_resource_name('dataset')
602+
created_dataset = await maybe_await(client.datasets().get_or_create(name=dataset_name))
603+
assert isinstance(created_dataset, Dataset)
604+
dataset_client = client.dataset(created_dataset.id)
605+
606+
try:
607+
items_to_push = [{'idx': i, 'parts': [{'part': p} for p in range(3)]} for i in range(12)]
608+
await maybe_await(dataset_client.push_items(items_to_push))
609+
610+
# Poll until all 12 rows are visible (eventual consistency) so the chunked iteration sees every page
611+
async def get_items() -> DatasetItemsPage:
612+
page = await maybe_await(dataset_client.list_items(limit=12))
613+
assert isinstance(page, DatasetItemsPage)
614+
return page
615+
616+
await poll_until_condition(get_items, lambda page: len(page.items) == 12)
617+
618+
# chunk_size=5 caps a page at 5 rows, which `unwind` expands into 15 items
619+
iterator = dataset_client.iterate_items(unwind=['parts'], chunk_size=5)
620+
collected: list[dict] = []
621+
if is_async:
622+
assert isinstance(iterator, AsyncIterator)
623+
async for item in iterator:
624+
assert isinstance(item, dict)
625+
collected.append(item)
626+
else:
627+
assert isinstance(iterator, Iterator)
628+
for item in iterator:
629+
assert isinstance(item, dict)
630+
collected.append(item)
631+
632+
# Every part of every row arrives exactly once: no page is skipped and none is read twice.
633+
assert sorted((item['idx'], item['part']) for item in collected) == [
634+
(idx, part) for idx in range(12) for part in range(3)
635+
]
636+
finally:
637+
await maybe_await(dataset_client.delete())
638+
639+
599640
async def test_dataset_iterate_items_with_fields(client: ApifyClient | ApifyClientAsync, *, is_async: bool) -> None:
600641
"""Test iterate_items with `fields` filter."""
601642
dataset_name = get_random_resource_name('dataset')

0 commit comments

Comments
 (0)