From c9007e52a64c194d594aa178ed8cbb823e89562e Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 29 Sep 2026 12:51:06 +0200 Subject: [PATCH 1/8] feat: Add live iteration over a run's dataset items --- src/apify_client/_resource_clients/run.py | 200 +++++++++++- tests/unit/test_run_iterate_dataset_items.py | 308 +++++++++++++++++++ 2 files changed, 507 insertions(+), 1 deletion(-) create mode 100644 tests/unit/test_run_iterate_dataset_items.py diff --git a/src/apify_client/_resource_clients/run.py b/src/apify_client/_resource_clients/run.py index 7ddd6b30..0d0ae873 100644 --- a/src/apify_client/_resource_clients/run.py +++ b/src/apify_client/_resource_clients/run.py @@ -1,5 +1,6 @@ from __future__ import annotations +import asyncio import json import random import string @@ -10,7 +11,8 @@ from apify_client._docs import docs_group from apify_client._logging import create_redirect_logger from apify_client._models import Run, RunResponse -from apify_client._resource_clients._resource_client import ResourceClient, ResourceClientAsync +from apify_client._pagination import DEFAULT_CHUNK_SIZE, _page_scanned_rows +from apify_client._resource_clients._resource_client import _TERMINAL_STATUSES, ResourceClient, ResourceClientAsync from apify_client._status_message_watcher import StatusMessageWatcher, StatusMessageWatcherAsync from apify_client._streamed_log import StreamedLog, StreamedLogAsync from apify_client._utils.encoding import encode_key_value_store_record_value @@ -19,6 +21,7 @@ if TYPE_CHECKING: import logging + from collections.abc import AsyncIterator, Iterator from decimal import Decimal from apify_client._literals import GeneralAccess @@ -32,6 +35,7 @@ RequestQueueClient, RequestQueueClientAsync, ) + from apify_client._resource_clients.dataset import DatasetItemsPage from apify_client.types import Timeout @@ -466,6 +470,102 @@ def get_status_message_watcher( return StatusMessageWatcher(run_client=self, to_logger=to_logger, check_period=check_period) + def iterate_dataset_items( + self, + *, + offset: int | None = None, + limit: int | None = None, + clean: bool | None = None, + fields: list[str] | None = None, + omit: list[str] | None = None, + unwind: list[str] | None = None, + skip_empty: bool | None = None, + skip_hidden: bool | None = None, + chunk_size: int | None = None, + poll_interval: timedelta = timedelta(seconds=5), + timeout: Timeout = 'long', + ) -> Iterator[dict]: + """Iterate over the items of the run's default dataset while the run is still producing them. + + While the run has not finished, the dataset is polled every `poll_interval` and the rows below its + `item_count` are yielded. Each page is requested with a `limit` that ends at `item_count`, so it covers + exactly the rows it asks for, whatever the filters or `unwind` do to the items. `item_count` lags a few + seconds behind the pushed items, so once the run reaches a terminal status, the rows past it are read until + a page comes back shorter than requested, and the iterator returns. + + https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items + + Args: + offset: Number of items that should be skipped at the start. The default value is 0. + limit: Maximum number of dataset rows to scan. Fewer items are yielded when filters drop some, more + when `unwind` splits a row into several. By default there is no limit. + clean: If True, returns only non-empty items and skips hidden fields (i.e. fields starting with + the # character). The clean parameter is just a shortcut for skip_hidden=True and skip_empty=True + parameters. + fields: A list of fields which should be picked from the items, only these fields will remain in + the resulting record objects. + omit: A list of fields which should be omitted from the items. + unwind: A list of fields which should be unwound, in order which they should be processed. Each field + should be either an array or an object. If the field is an array then every element of the array + will become a separate record and merged with parent object. If the unwound field is an object then + it is merged with the parent object. + skip_empty: If True, then empty items are skipped from the output. + skip_hidden: If True, then hidden fields are skipped from the output, i.e. fields starting with + the # character. + chunk_size: Maximum number of dataset rows requested per API call. + poll_interval: How long to wait between polls while the run has not finished. + timeout: Timeout for each API HTTP request. + + Yields: + An item from the dataset. + """ + dataset_client = self.dataset() + page_size = chunk_size or DEFAULT_CHUNK_SIZE + position = offset or 0 + end = position + limit if limit else None + + def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: + return dataset_client.list_items( + offset=page_offset, + limit=page_limit, + clean=clean, + fields=fields, + omit=omit, + unwind=unwind, + skip_empty=skip_empty, + skip_hidden=skip_hidden, + timeout=timeout, + ) + + while True: + run = self.get(timeout=timeout) + is_finished = run is None or run.status in _TERMINAL_STATUSES + dataset = dataset_client.get(timeout=timeout) + item_count = dataset.item_count if dataset else 0 + if end is not None: + item_count = min(item_count, end) + + while position < item_count: + page_limit = min(page_size, item_count - position) + page = list_page(position, page_limit) + yield from page.items + position += page_limit + + if end is not None and position >= end: + return + if is_finished: + break + time.sleep(to_seconds(poll_interval)) + + while True: + page_limit = min(page_size, end - position) if end is not None else page_size + page = list_page(position, page_limit) + yield from page.items + scanned_rows = _page_scanned_rows(page, page_limit) + position += scanned_rows + if scanned_rows < page_limit or (end is not None and position >= end): + return + @docs_group('Resource clients') class RunClientAsync(ResourceClientAsync): @@ -897,3 +997,101 @@ async def get_status_message_watcher( to_logger = create_redirect_logger(f'apify.{name}') return StatusMessageWatcherAsync(run_client=self, to_logger=to_logger, check_period=check_period) + + async def iterate_dataset_items( + self, + *, + offset: int | None = None, + limit: int | None = None, + clean: bool | None = None, + fields: list[str] | None = None, + omit: list[str] | None = None, + unwind: list[str] | None = None, + skip_empty: bool | None = None, + skip_hidden: bool | None = None, + chunk_size: int | None = None, + poll_interval: timedelta = timedelta(seconds=5), + timeout: Timeout = 'long', + ) -> AsyncIterator[dict]: + """Iterate over the items of the run's default dataset while the run is still producing them. + + While the run has not finished, the dataset is polled every `poll_interval` and the rows below its + `item_count` are yielded. Each page is requested with a `limit` that ends at `item_count`, so it covers + exactly the rows it asks for, whatever the filters or `unwind` do to the items. `item_count` lags a few + seconds behind the pushed items, so once the run reaches a terminal status, the rows past it are read until + a page comes back shorter than requested, and the iterator returns. + + https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items + + Args: + offset: Number of items that should be skipped at the start. The default value is 0. + limit: Maximum number of dataset rows to scan. Fewer items are yielded when filters drop some, more + when `unwind` splits a row into several. By default there is no limit. + clean: If True, returns only non-empty items and skips hidden fields (i.e. fields starting with + the # character). The clean parameter is just a shortcut for skip_hidden=True and skip_empty=True + parameters. + fields: A list of fields which should be picked from the items, only these fields will remain in + the resulting record objects. + omit: A list of fields which should be omitted from the items. + unwind: A list of fields which should be unwound, in order which they should be processed. Each field + should be either an array or an object. If the field is an array then every element of the array + will become a separate record and merged with parent object. If the unwound field is an object then + it is merged with the parent object. + skip_empty: If True, then empty items are skipped from the output. + skip_hidden: If True, then hidden fields are skipped from the output, i.e. fields starting with + the # character. + chunk_size: Maximum number of dataset rows requested per API call. + poll_interval: How long to wait between polls while the run has not finished. + timeout: Timeout for each API HTTP request. + + Yields: + An item from the dataset. + """ + dataset_client = self.dataset() + page_size = chunk_size or DEFAULT_CHUNK_SIZE + position = offset or 0 + end = position + limit if limit else None + + async def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: + return await dataset_client.list_items( + offset=page_offset, + limit=page_limit, + clean=clean, + fields=fields, + omit=omit, + unwind=unwind, + skip_empty=skip_empty, + skip_hidden=skip_hidden, + timeout=timeout, + ) + + while True: + run = await self.get(timeout=timeout) + is_finished = run is None or run.status in _TERMINAL_STATUSES + dataset = await dataset_client.get(timeout=timeout) + item_count = dataset.item_count if dataset else 0 + if end is not None: + item_count = min(item_count, end) + + while position < item_count: + page_limit = min(page_size, item_count - position) + page = await list_page(position, page_limit) + for item in page.items: + yield item + position += page_limit + + if end is not None and position >= end: + return + if is_finished: + break + await asyncio.sleep(to_seconds(poll_interval)) + + while True: + page_limit = min(page_size, end - position) if end is not None else page_size + page = await list_page(position, page_limit) + for item in page.items: + yield item + scanned_rows = _page_scanned_rows(page, page_limit) + position += scanned_rows + if scanned_rows < page_limit or (end is not None and position >= end): + return diff --git a/tests/unit/test_run_iterate_dataset_items.py b/tests/unit/test_run_iterate_dataset_items.py new file mode 100644 index 00000000..797d0da6 --- /dev/null +++ b/tests/unit/test_run_iterate_dataset_items.py @@ -0,0 +1,308 @@ +from __future__ import annotations + +import json +from dataclasses import dataclass +from datetime import timedelta +from typing import TYPE_CHECKING, Any +from unittest.mock import AsyncMock, Mock, call + +import pytest +from werkzeug import Response + +from apify_client import ApifyClient, ApifyClientAsync + +if TYPE_CHECKING: + from pytest_httpserver import HTTPServer + from werkzeug import Request + + from apify_client._literals import ActorJobStatus + +pytestmark = pytest.mark.usefixtures('http_client_classes') + +RUN_ID = 'test-run-id' +RUN_PATH = f'/v2/actor-runs/{RUN_ID}' +UNWIND_PARTS = 3 +NO_WAIT = timedelta(0) + + +@dataclass +class Step: + """State of the run that the fake API switches to on one run status read.""" + + pushed_rows: int + """Total number of rows in the dataset, readable right away.""" + + item_count: int + """The dataset's `itemCount`, which lags behind the pushed rows.""" + + status: ActorJobStatus + + +@dataclass +class FakeRunApi: + """Fake run, run dataset and dataset items endpoints, advanced by one `Step` per run status read.""" + + steps: list[Step] + step_index: int = -1 + + @property + def step(self) -> Step: + return self.steps[max(self.step_index, 0)] + + def register(self, httpserver: HTTPServer) -> None: + httpserver.expect_request(RUN_PATH, method='GET').respond_with_handler(self.handle_run) + httpserver.expect_request(f'{RUN_PATH}/dataset', method='GET').respond_with_handler(self.handle_dataset) + httpserver.expect_request(f'{RUN_PATH}/dataset/items', method='GET').respond_with_handler(self.handle_items) + + def handle_run(self, _request: Request) -> Response: + self.step_index = min(self.step_index + 1, len(self.steps) - 1) + run = { + 'id': RUN_ID, + 'actId': 'test-actor-id', + 'userId': 'test-user-id', + 'startedAt': '2019-11-30T07:34:24.202Z', + 'status': self.step.status, + 'meta': {'origin': 'WEB'}, + 'stats': {'restartCount': 0, 'resurrectCount': 0, 'computeUnits': 0.1}, + 'options': {'build': 'latest', 'timeoutSecs': 300, 'memoryMbytes': 1024, 'diskMbytes': 2048}, + 'buildId': 'test-build-id', + 'defaultKeyValueStoreId': 'test-kvs-id', + 'defaultDatasetId': 'test-dataset-id', + 'defaultRequestQueueId': 'test-rq-id', + 'containerUrl': 'https://test.runs.apify.net', + } + return Response(json.dumps({'data': run}), status=200, mimetype='application/json') + + def handle_dataset(self, _request: Request) -> Response: + dataset = { + 'id': 'test-dataset-id', + 'userId': 'test-user-id', + 'createdAt': '2019-12-12T07:34:14.202Z', + 'modifiedAt': '2019-12-13T08:36:13.202Z', + 'accessedAt': '2019-12-14T08:36:13.202Z', + 'itemCount': self.step.item_count, + 'cleanItemCount': self.step.item_count, + 'consoleUrl': 'https://console.apify.com/storage/datasets/test-dataset-id', + } + return Response(json.dumps({'data': dataset}), status=200, mimetype='application/json') + + def handle_items(self, request: Request) -> Response: + """Scan exactly `limit` rows from `offset`, like the real endpoint, and derive the count from `itemCount`.""" + offset = int(request.args.get('offset', 0)) + limit = int(request.args.get('limit', 0)) or 999_999_999_999 + rows = range(offset, min(offset + limit, self.step.pushed_rows)) + items = shape_items(rows, clean=request.args.get('clean') == 'true', unwind=bool(request.args.get('unwind'))) + headers = { + 'x-apify-pagination-total': str(self.step.item_count), + 'x-apify-pagination-offset': str(offset), + 'x-apify-pagination-count': str(max(min(self.step.item_count - offset, limit), 0)), + 'x-apify-pagination-limit': str(limit), + 'x-apify-pagination-desc': 'false', + } + return Response(json.dumps(items), status=200, headers=headers, mimetype='application/json') + + +def shape_items(rows: range, *, clean: bool = False, unwind: bool = False) -> list[dict[str, Any]]: + """Turn dataset rows into items: `clean` drops every odd row, `unwind` splits a row into `UNWIND_PARTS` items.""" + kept_rows = [row for row in rows if not (clean and row % 2)] + if unwind: + return [{'row': row, 'part': part} for row in kept_rows for part in range(UNWIND_PARTS)] + return [{'row': row} for row in kept_rows] + + +LAGGING_RUN_STEPS = [ + Step(pushed_rows=3, item_count=0, status='READY'), + Step(pushed_rows=40, item_count=3, status='RUNNING'), + Step(pushed_rows=52, item_count=25, status='RUNNING'), + # The run has finished, but `itemCount` still lags more than one page behind the readable rows. + Step(pushed_rows=75, item_count=52, status='SUCCEEDED'), +] + +# The same run with `itemCount` caught up by the time it finished. +CAUGHT_UP_RUN_STEPS = [*LAGGING_RUN_STEPS[:-1], Step(pushed_rows=52, item_count=52, status='SUCCEEDED')] + + +def test_iterate_dataset_items_yields_every_row_once_sync(httpserver: HTTPServer) -> None: + """Rows pushed across polls and past a lagging item count after the run finished are all yielded, in order.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + + items = list(client.run(RUN_ID).iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)) + + assert items == shape_items(range(75)) + + +async def test_iterate_dataset_items_yields_every_row_once_async(httpserver: HTTPServer) -> None: + """Rows pushed across polls and past a lagging item count after the run finished are all yielded, in order.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + + items = [item async for item in client.run(RUN_ID).iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)] + + assert items == shape_items(range(75)) + + +@pytest.mark.parametrize( + 'shaping', + [ + pytest.param({'clean': True}, id='clean drops items'), + pytest.param({'unwind': True}, id='unwind multiplies items'), + ], +) +def test_iterate_dataset_items_with_shaped_items_sync(httpserver: HTTPServer, shaping: dict[str, bool]) -> None: + """Filters and `unwind` change the item count per page without duplicating or skipping rows while the run runs.""" + api = FakeRunApi(CAUGHT_UP_RUN_STEPS) + api.register(httpserver) + client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + + items = list( + client.run(RUN_ID).iterate_dataset_items( + clean=shaping.get('clean'), + unwind=['parts'] if shaping.get('unwind') else None, + chunk_size=10, + poll_interval=NO_WAIT, + ) + ) + + assert items == shape_items(range(52), **shaping) + + +@pytest.mark.parametrize( + 'shaping', + [ + pytest.param({'clean': True}, id='clean drops items'), + pytest.param({'unwind': True}, id='unwind multiplies items'), + ], +) +async def test_iterate_dataset_items_with_shaped_items_async(httpserver: HTTPServer, shaping: dict[str, bool]) -> None: + """Filters and `unwind` change the item count per page without duplicating or skipping rows while the run runs.""" + api = FakeRunApi(CAUGHT_UP_RUN_STEPS) + api.register(httpserver) + client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + + items = [ + item + async for item in client.run(RUN_ID).iterate_dataset_items( + clean=shaping.get('clean'), + unwind=['parts'] if shaping.get('unwind') else None, + chunk_size=10, + poll_interval=NO_WAIT, + ) + ] + + assert items == shape_items(range(52), **shaping) + + +@pytest.mark.parametrize( + 'status', + [ + pytest.param('ABORTING', id='aborting'), + pytest.param('TIMING-OUT', id='timing out'), + ], +) +def test_iterate_dataset_items_keeps_polling_until_terminal_sync( + httpserver: HTTPServer, status: ActorJobStatus +) -> None: + """A run that is aborting or timing out can still push items, so polling goes on until a terminal status.""" + api = FakeRunApi( + [ + Step(pushed_rows=5, item_count=5, status=status), + Step(pushed_rows=8, item_count=8, status='ABORTED'), + ] + ) + api.register(httpserver) + client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + + items = list(client.run(RUN_ID).iterate_dataset_items(poll_interval=NO_WAIT)) + + assert items == shape_items(range(8)) + + +@pytest.mark.parametrize( + 'status', + [ + pytest.param('ABORTING', id='aborting'), + pytest.param('TIMING-OUT', id='timing out'), + ], +) +async def test_iterate_dataset_items_keeps_polling_until_terminal_async( + httpserver: HTTPServer, status: ActorJobStatus +) -> None: + """A run that is aborting or timing out can still push items, so polling goes on until a terminal status.""" + api = FakeRunApi( + [ + Step(pushed_rows=5, item_count=5, status=status), + Step(pushed_rows=8, item_count=8, status='ABORTED'), + ] + ) + api.register(httpserver) + client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + + items = [item async for item in client.run(RUN_ID).iterate_dataset_items(poll_interval=NO_WAIT)] + + assert items == shape_items(range(8)) + + +def test_iterate_dataset_items_respects_offset_and_limit_sync(httpserver: HTTPServer) -> None: + """Iteration starts at `offset` and stops once `limit` rows are scanned, without waiting for the run to finish.""" + api = FakeRunApi( + [Step(pushed_rows=4, item_count=4, status='RUNNING'), Step(pushed_rows=20, item_count=20, status='RUNNING')] + ) + api.register(httpserver) + client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + + items = list(client.run(RUN_ID).iterate_dataset_items(offset=2, limit=7, chunk_size=3, poll_interval=NO_WAIT)) + + assert items == shape_items(range(2, 9)) + assert api.step.status == 'RUNNING' + + +async def test_iterate_dataset_items_respects_offset_and_limit_async(httpserver: HTTPServer) -> None: + """Iteration starts at `offset` and stops once `limit` rows are scanned, without waiting for the run to finish.""" + api = FakeRunApi( + [Step(pushed_rows=4, item_count=4, status='RUNNING'), Step(pushed_rows=20, item_count=20, status='RUNNING')] + ) + api.register(httpserver) + client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + + items = [ + item + async for item in client.run(RUN_ID).iterate_dataset_items( + offset=2, limit=7, chunk_size=3, poll_interval=NO_WAIT + ) + ] + + assert items == shape_items(range(2, 9)) + assert api.step.status == 'RUNNING' + + +def test_iterate_dataset_items_sleeps_between_polls_sync( + httpserver: HTTPServer, monkeypatch: pytest.MonkeyPatch +) -> None: + """The iterator waits `poll_interval` after each poll of an unfinished run and not after the final one.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + sleep = Mock() + monkeypatch.setattr('apify_client._resource_clients.run.time.sleep', sleep) + + list(client.run(RUN_ID).iterate_dataset_items(poll_interval=timedelta(seconds=2))) + + assert sleep.call_args_list == [call(2.0)] * 3 + + +async def test_iterate_dataset_items_sleeps_between_polls_async( + httpserver: HTTPServer, monkeypatch: pytest.MonkeyPatch +) -> None: + """The iterator waits `poll_interval` after each poll of an unfinished run and not after the final one.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) + sleep = AsyncMock() + monkeypatch.setattr('apify_client._resource_clients.run.asyncio.sleep', sleep) + + [item async for item in client.run(RUN_ID).iterate_dataset_items(poll_interval=timedelta(seconds=2))] + + assert sleep.call_args_list == [call(2.0)] * 3 From 60398406cc04cfcd4f48715421af0e4291a4b07f Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 29 Sep 2026 13:09:58 +0200 Subject: [PATCH 2/8] fix: Keep reading rows past a lagging item count until none are left --- src/apify_client/_resource_clients/run.py | 34 +++-- tests/unit/test_run_iterate_dataset_items.py | 128 ++++++++++++------- 2 files changed, 104 insertions(+), 58 deletions(-) diff --git a/src/apify_client/_resource_clients/run.py b/src/apify_client/_resource_clients/run.py index 0d0ae873..1c32b70d 100644 --- a/src/apify_client/_resource_clients/run.py +++ b/src/apify_client/_resource_clients/run.py @@ -11,7 +11,7 @@ from apify_client._docs import docs_group from apify_client._logging import create_redirect_logger from apify_client._models import Run, RunResponse -from apify_client._pagination import DEFAULT_CHUNK_SIZE, _page_scanned_rows +from apify_client._pagination import DEFAULT_CHUNK_SIZE from apify_client._resource_clients._resource_client import _TERMINAL_STATUSES, ResourceClient, ResourceClientAsync from apify_client._status_message_watcher import StatusMessageWatcher, StatusMessageWatcherAsync from apify_client._streamed_log import StreamedLog, StreamedLogAsync @@ -490,8 +490,8 @@ def iterate_dataset_items( While the run has not finished, the dataset is polled every `poll_interval` and the rows below its `item_count` are yielded. Each page is requested with a `limit` that ends at `item_count`, so it covers exactly the rows it asks for, whatever the filters or `unwind` do to the items. `item_count` lags a few - seconds behind the pushed items, so once the run reaches a terminal status, the rows past it are read until - a page comes back shorter than requested, and the iterator returns. + seconds behind the pushed items, so once the run reaches a terminal status, the rows past it are read a page + at a time until none are left, and the iterator returns. https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items @@ -561,9 +561,15 @@ def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: page_limit = min(page_size, end - position) if end is not None else page_size page = list_page(position, page_limit) yield from page.items - scanned_rows = _page_scanned_rows(page, page_limit) - position += scanned_rows - if scanned_rows < page_limit or (end is not None and position >= end): + # Only an empty page marks the end, as filters can shorten a full one. A page that `clean`, `skip_empty` or + # `unwind` emptied past a lagging `item_count` reports no scanned rows either, so a plain read checks. + if not page.count and ( + not (clean or skip_empty or unwind) + or not dataset_client.list_items(offset=position, limit=1, timeout=timeout).items + ): + return + position += page_limit + if end is not None and position >= end: return @@ -1018,8 +1024,8 @@ async def iterate_dataset_items( While the run has not finished, the dataset is polled every `poll_interval` and the rows below its `item_count` are yielded. Each page is requested with a `limit` that ends at `item_count`, so it covers exactly the rows it asks for, whatever the filters or `unwind` do to the items. `item_count` lags a few - seconds behind the pushed items, so once the run reaches a terminal status, the rows past it are read until - a page comes back shorter than requested, and the iterator returns. + seconds behind the pushed items, so once the run reaches a terminal status, the rows past it are read a page + at a time until none are left, and the iterator returns. https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items @@ -1091,7 +1097,13 @@ async def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: page = await list_page(position, page_limit) for item in page.items: yield item - scanned_rows = _page_scanned_rows(page, page_limit) - position += scanned_rows - if scanned_rows < page_limit or (end is not None and position >= end): + # Only an empty page marks the end, as filters can shorten a full one. A page that `clean`, `skip_empty` or + # `unwind` emptied past a lagging `item_count` reports no scanned rows either, so a plain read checks. + if not page.count and ( + not (clean or skip_empty or unwind) + or not (await dataset_client.list_items(offset=position, limit=1, timeout=timeout)).items + ): + return + position += page_limit + if end is not None and position >= end: return diff --git a/tests/unit/test_run_iterate_dataset_items.py b/tests/unit/test_run_iterate_dataset_items.py index 797d0da6..bd9b4f07 100644 --- a/tests/unit/test_run_iterate_dataset_items.py +++ b/tests/unit/test_run_iterate_dataset_items.py @@ -9,12 +9,11 @@ import pytest from werkzeug import Response -from apify_client import ApifyClient, ApifyClientAsync - if TYPE_CHECKING: from pytest_httpserver import HTTPServer from werkzeug import Request + from apify_client import ApifyClient, ApifyClientAsync from apify_client._literals import ActorJobStatus pytestmark = pytest.mark.usefixtures('http_client_classes') @@ -103,10 +102,13 @@ def handle_items(self, request: Request) -> Response: def shape_items(rows: range, *, clean: bool = False, unwind: bool = False) -> list[dict[str, Any]]: - """Turn dataset rows into items: `clean` drops every odd row, `unwind` splits a row into `UNWIND_PARTS` items.""" + """Turn dataset rows into items: `clean` drops every odd row, `unwind` splits a row into `UNWIND_PARTS` items. + + Under `unwind`, every third row holds an empty array and so unwinds into no items at all. + """ kept_rows = [row for row in rows if not (clean and row % 2)] if unwind: - return [{'row': row, 'part': part} for row in kept_rows for part in range(UNWIND_PARTS)] + return [{'row': row, 'part': part} for row in kept_rows if row % 3 != 2 for part in range(UNWIND_PARTS)] return [{'row': row} for row in kept_rows] @@ -118,81 +120,90 @@ def shape_items(rows: range, *, clean: bool = False, unwind: bool = False) -> li Step(pushed_rows=75, item_count=52, status='SUCCEEDED'), ] -# The same run with `itemCount` caught up by the time it finished. -CAUGHT_UP_RUN_STEPS = [*LAGGING_RUN_STEPS[:-1], Step(pushed_rows=52, item_count=52, status='SUCCEEDED')] - -def test_iterate_dataset_items_yields_every_row_once_sync(httpserver: HTTPServer) -> None: +def test_iterate_dataset_items_yields_every_row_once_sync(httpserver: HTTPServer, sync_client: ApifyClient) -> None: """Rows pushed across polls and past a lagging item count after the run finished are all yielded, in order.""" api = FakeRunApi(LAGGING_RUN_STEPS) api.register(httpserver) - client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) - items = list(client.run(RUN_ID).iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)) + items = list(sync_client.run(RUN_ID).iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)) assert items == shape_items(range(75)) -async def test_iterate_dataset_items_yields_every_row_once_async(httpserver: HTTPServer) -> None: +async def test_iterate_dataset_items_yields_every_row_once_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: """Rows pushed across polls and past a lagging item count after the run finished are all yielded, in order.""" api = FakeRunApi(LAGGING_RUN_STEPS) api.register(httpserver) - client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) - items = [item async for item in client.run(RUN_ID).iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)] + items = [ + item async for item in async_client.run(RUN_ID).iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT) + ] assert items == shape_items(range(75)) +CHUNK_SIZES = [ + pytest.param(10, id='partly filtered pages'), + pytest.param(1, id='fully filtered pages'), +] + + +@pytest.mark.parametrize('chunk_size', CHUNK_SIZES) @pytest.mark.parametrize( 'shaping', [ pytest.param({'clean': True}, id='clean drops items'), - pytest.param({'unwind': True}, id='unwind multiplies items'), + pytest.param({'unwind': True}, id='unwind multiplies or drops items'), ], ) -def test_iterate_dataset_items_with_shaped_items_sync(httpserver: HTTPServer, shaping: dict[str, bool]) -> None: - """Filters and `unwind` change the item count per page without duplicating or skipping rows while the run runs.""" - api = FakeRunApi(CAUGHT_UP_RUN_STEPS) +def test_iterate_dataset_items_with_shaped_items_sync( + httpserver: HTTPServer, sync_client: ApifyClient, shaping: dict[str, bool], chunk_size: int +) -> None: + """Filters and `unwind` reshape or empty pages without duplicating or skipping rows, also past a lagging count.""" + api = FakeRunApi(LAGGING_RUN_STEPS) api.register(httpserver) - client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) items = list( - client.run(RUN_ID).iterate_dataset_items( + sync_client.run(RUN_ID).iterate_dataset_items( clean=shaping.get('clean'), unwind=['parts'] if shaping.get('unwind') else None, - chunk_size=10, + chunk_size=chunk_size, poll_interval=NO_WAIT, ) ) - assert items == shape_items(range(52), **shaping) + assert items == shape_items(range(75), **shaping) +@pytest.mark.parametrize('chunk_size', CHUNK_SIZES) @pytest.mark.parametrize( 'shaping', [ pytest.param({'clean': True}, id='clean drops items'), - pytest.param({'unwind': True}, id='unwind multiplies items'), + pytest.param({'unwind': True}, id='unwind multiplies or drops items'), ], ) -async def test_iterate_dataset_items_with_shaped_items_async(httpserver: HTTPServer, shaping: dict[str, bool]) -> None: - """Filters and `unwind` change the item count per page without duplicating or skipping rows while the run runs.""" - api = FakeRunApi(CAUGHT_UP_RUN_STEPS) +async def test_iterate_dataset_items_with_shaped_items_async( + httpserver: HTTPServer, async_client: ApifyClientAsync, shaping: dict[str, bool], chunk_size: int +) -> None: + """Filters and `unwind` reshape or empty pages without duplicating or skipping rows, also past a lagging count.""" + api = FakeRunApi(LAGGING_RUN_STEPS) api.register(httpserver) - client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) items = [ item - async for item in client.run(RUN_ID).iterate_dataset_items( + async for item in async_client.run(RUN_ID).iterate_dataset_items( clean=shaping.get('clean'), unwind=['parts'] if shaping.get('unwind') else None, - chunk_size=10, + chunk_size=chunk_size, poll_interval=NO_WAIT, ) ] - assert items == shape_items(range(52), **shaping) + assert items == shape_items(range(75), **shaping) @pytest.mark.parametrize( @@ -203,7 +214,7 @@ async def test_iterate_dataset_items_with_shaped_items_async(httpserver: HTTPSer ], ) def test_iterate_dataset_items_keeps_polling_until_terminal_sync( - httpserver: HTTPServer, status: ActorJobStatus + httpserver: HTTPServer, sync_client: ApifyClient, status: ActorJobStatus ) -> None: """A run that is aborting or timing out can still push items, so polling goes on until a terminal status.""" api = FakeRunApi( @@ -213,9 +224,8 @@ def test_iterate_dataset_items_keeps_polling_until_terminal_sync( ] ) api.register(httpserver) - client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) - items = list(client.run(RUN_ID).iterate_dataset_items(poll_interval=NO_WAIT)) + items = list(sync_client.run(RUN_ID).iterate_dataset_items(poll_interval=NO_WAIT)) assert items == shape_items(range(8)) @@ -228,7 +238,7 @@ def test_iterate_dataset_items_keeps_polling_until_terminal_sync( ], ) async def test_iterate_dataset_items_keeps_polling_until_terminal_async( - httpserver: HTTPServer, status: ActorJobStatus + httpserver: HTTPServer, async_client: ApifyClientAsync, status: ActorJobStatus ) -> None: """A run that is aborting or timing out can still push items, so polling goes on until a terminal status.""" api = FakeRunApi( @@ -238,38 +248,37 @@ async def test_iterate_dataset_items_keeps_polling_until_terminal_async( ] ) api.register(httpserver) - client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) - items = [item async for item in client.run(RUN_ID).iterate_dataset_items(poll_interval=NO_WAIT)] + items = [item async for item in async_client.run(RUN_ID).iterate_dataset_items(poll_interval=NO_WAIT)] assert items == shape_items(range(8)) -def test_iterate_dataset_items_respects_offset_and_limit_sync(httpserver: HTTPServer) -> None: +def test_iterate_dataset_items_respects_offset_and_limit_sync(httpserver: HTTPServer, sync_client: ApifyClient) -> None: """Iteration starts at `offset` and stops once `limit` rows are scanned, without waiting for the run to finish.""" api = FakeRunApi( [Step(pushed_rows=4, item_count=4, status='RUNNING'), Step(pushed_rows=20, item_count=20, status='RUNNING')] ) api.register(httpserver) - client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) - items = list(client.run(RUN_ID).iterate_dataset_items(offset=2, limit=7, chunk_size=3, poll_interval=NO_WAIT)) + items = list(sync_client.run(RUN_ID).iterate_dataset_items(offset=2, limit=7, chunk_size=3, poll_interval=NO_WAIT)) assert items == shape_items(range(2, 9)) assert api.step.status == 'RUNNING' -async def test_iterate_dataset_items_respects_offset_and_limit_async(httpserver: HTTPServer) -> None: +async def test_iterate_dataset_items_respects_offset_and_limit_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: """Iteration starts at `offset` and stops once `limit` rows are scanned, without waiting for the run to finish.""" api = FakeRunApi( [Step(pushed_rows=4, item_count=4, status='RUNNING'), Step(pushed_rows=20, item_count=20, status='RUNNING')] ) api.register(httpserver) - client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) items = [ item - async for item in client.run(RUN_ID).iterate_dataset_items( + async for item in async_client.run(RUN_ID).iterate_dataset_items( offset=2, limit=7, chunk_size=3, poll_interval=NO_WAIT ) ] @@ -278,31 +287,56 @@ async def test_iterate_dataset_items_respects_offset_and_limit_async(httpserver: assert api.step.status == 'RUNNING' +def test_iterate_dataset_items_limit_ends_reading_past_item_count_sync( + httpserver: HTTPServer, sync_client: ApifyClient +) -> None: + """A `limit` beyond the lagging item count of a finished run stops the reads past it at the limit.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + items = list(sync_client.run(RUN_ID).iterate_dataset_items(limit=60, chunk_size=10, poll_interval=NO_WAIT)) + + assert items == shape_items(range(60)) + + +async def test_iterate_dataset_items_limit_ends_reading_past_item_count_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: + """A `limit` beyond the lagging item count of a finished run stops the reads past it at the limit.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + items = [ + item + async for item in async_client.run(RUN_ID).iterate_dataset_items(limit=60, chunk_size=10, poll_interval=NO_WAIT) + ] + + assert items == shape_items(range(60)) + + def test_iterate_dataset_items_sleeps_between_polls_sync( - httpserver: HTTPServer, monkeypatch: pytest.MonkeyPatch + httpserver: HTTPServer, sync_client: ApifyClient, monkeypatch: pytest.MonkeyPatch ) -> None: """The iterator waits `poll_interval` after each poll of an unfinished run and not after the final one.""" api = FakeRunApi(LAGGING_RUN_STEPS) api.register(httpserver) - client = ApifyClient(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) sleep = Mock() monkeypatch.setattr('apify_client._resource_clients.run.time.sleep', sleep) - list(client.run(RUN_ID).iterate_dataset_items(poll_interval=timedelta(seconds=2))) + list(sync_client.run(RUN_ID).iterate_dataset_items(poll_interval=timedelta(seconds=2))) assert sleep.call_args_list == [call(2.0)] * 3 async def test_iterate_dataset_items_sleeps_between_polls_async( - httpserver: HTTPServer, monkeypatch: pytest.MonkeyPatch + httpserver: HTTPServer, async_client: ApifyClientAsync, monkeypatch: pytest.MonkeyPatch ) -> None: """The iterator waits `poll_interval` after each poll of an unfinished run and not after the final one.""" api = FakeRunApi(LAGGING_RUN_STEPS) api.register(httpserver) - client = ApifyClientAsync(token='test-token', api_url=httpserver.url_for('/').removesuffix('/')) sleep = AsyncMock() monkeypatch.setattr('apify_client._resource_clients.run.asyncio.sleep', sleep) - [item async for item in client.run(RUN_ID).iterate_dataset_items(poll_interval=timedelta(seconds=2))] + [item async for item in async_client.run(RUN_ID).iterate_dataset_items(poll_interval=timedelta(seconds=2))] assert sleep.call_args_list == [call(2.0)] * 3 From 0dcca7f3fbb6d9e5844b2ec1859bd2ceefd9ecd8 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 29 Sep 2026 13:10:00 +0200 Subject: [PATCH 3/8] docs: List iterate_dataset_items among the convenience methods --- docs/02_concepts/07_convenience_methods.mdx | 1 + 1 file changed, 1 insertion(+) diff --git a/docs/02_concepts/07_convenience_methods.mdx b/docs/02_concepts/07_convenience_methods.mdx index 24a25593..87c913c8 100644 --- a/docs/02_concepts/07_convenience_methods.mdx +++ b/docs/02_concepts/07_convenience_methods.mdx @@ -18,6 +18,7 @@ The Apify client provides several convenience methods to handle actions that the - `ActorClient.call` - Starts an Actor and waits for it to finish, handling network timeouts internally. Waits indefinitely by default, or up to the specified `wait_duration`. - `ActorClient.start` - Starts an Actor and immediately returns the Run object without waiting for it to finish. - `RunClient.wait_for_finish` - Waits for an already-started run to reach a terminal status. +- `RunClient.iterate_dataset_items` - Yields the items of the run's default dataset as the run pushes them, and returns once the run has finished and every item is read. Additionally, storage-related resources offer flexible options for data retrieval: From 2c2688e3b9754e431214a0effddd44eb7c620ba0 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 29 Sep 2026 13:59:43 +0200 Subject: [PATCH 4/8] test: Cover skip_empty and option forwarding in iterate_dataset_items --- tests/unit/test_run_iterate_dataset_items.py | 54 ++++++++++++++++++-- 1 file changed, 49 insertions(+), 5 deletions(-) diff --git a/tests/unit/test_run_iterate_dataset_items.py b/tests/unit/test_run_iterate_dataset_items.py index bd9b4f07..558224ca 100644 --- a/tests/unit/test_run_iterate_dataset_items.py +++ b/tests/unit/test_run_iterate_dataset_items.py @@ -1,7 +1,7 @@ from __future__ import annotations import json -from dataclasses import dataclass +from dataclasses import dataclass, field from datetime import timedelta from typing import TYPE_CHECKING, Any from unittest.mock import AsyncMock, Mock, call @@ -43,6 +43,7 @@ class FakeRunApi: steps: list[Step] step_index: int = -1 + items_requests: list[dict[str, str]] = field(default_factory=list) @property def step(self) -> Step: @@ -89,8 +90,14 @@ def handle_items(self, request: Request) -> Response: """Scan exactly `limit` rows from `offset`, like the real endpoint, and derive the count from `itemCount`.""" offset = int(request.args.get('offset', 0)) limit = int(request.args.get('limit', 0)) or 999_999_999_999 + self.items_requests.append(dict(request.args)) rows = range(offset, min(offset + limit, self.step.pushed_rows)) - items = shape_items(rows, clean=request.args.get('clean') == 'true', unwind=bool(request.args.get('unwind'))) + items = shape_items( + rows, + clean=request.args.get('clean') == 'true', + skip_empty=request.args.get('skipEmpty') == 'true', + unwind=bool(request.args.get('unwind')), + ) headers = { 'x-apify-pagination-total': str(self.step.item_count), 'x-apify-pagination-offset': str(offset), @@ -101,12 +108,14 @@ def handle_items(self, request: Request) -> Response: return Response(json.dumps(items), status=200, headers=headers, mimetype='application/json') -def shape_items(rows: range, *, clean: bool = False, unwind: bool = False) -> list[dict[str, Any]]: - """Turn dataset rows into items: `clean` drops every odd row, `unwind` splits a row into `UNWIND_PARTS` items. +def shape_items( + rows: range, *, clean: bool = False, skip_empty: bool = False, unwind: bool = False +) -> list[dict[str, Any]]: + """Turn dataset rows into items: `clean` and `skip_empty` drop every odd row, `unwind` splits a row into items. Under `unwind`, every third row holds an empty array and so unwinds into no items at all. """ - kept_rows = [row for row in rows if not (clean and row % 2)] + kept_rows = [row for row in rows if not ((clean or skip_empty) and row % 2)] if unwind: return [{'row': row, 'part': part} for row in kept_rows if row % 3 != 2 for part in range(UNWIND_PARTS)] return [{'row': row} for row in kept_rows] @@ -156,6 +165,7 @@ async def test_iterate_dataset_items_yields_every_row_once_async( 'shaping', [ pytest.param({'clean': True}, id='clean drops items'), + pytest.param({'skip_empty': True}, id='skip_empty drops items'), pytest.param({'unwind': True}, id='unwind multiplies or drops items'), ], ) @@ -169,6 +179,7 @@ def test_iterate_dataset_items_with_shaped_items_sync( items = list( sync_client.run(RUN_ID).iterate_dataset_items( clean=shaping.get('clean'), + skip_empty=shaping.get('skip_empty'), unwind=['parts'] if shaping.get('unwind') else None, chunk_size=chunk_size, poll_interval=NO_WAIT, @@ -183,6 +194,7 @@ def test_iterate_dataset_items_with_shaped_items_sync( 'shaping', [ pytest.param({'clean': True}, id='clean drops items'), + pytest.param({'skip_empty': True}, id='skip_empty drops items'), pytest.param({'unwind': True}, id='unwind multiplies or drops items'), ], ) @@ -197,6 +209,7 @@ async def test_iterate_dataset_items_with_shaped_items_async( item async for item in async_client.run(RUN_ID).iterate_dataset_items( clean=shaping.get('clean'), + skip_empty=shaping.get('skip_empty'), unwind=['parts'] if shaping.get('unwind') else None, chunk_size=chunk_size, poll_interval=NO_WAIT, @@ -314,6 +327,37 @@ async def test_iterate_dataset_items_limit_ends_reading_past_item_count_async( assert items == shape_items(range(60)) +def test_iterate_dataset_items_forwards_options_without_extra_read_sync( + httpserver: HTTPServer, sync_client: ApifyClient +) -> None: + """Item options reach every page request, and without a dropping filter an empty page ends with no plain read.""" + api = FakeRunApi([Step(pushed_rows=5, item_count=3, status='SUCCEEDED')]) + api.register(httpserver) + + items = list(sync_client.run(RUN_ID).iterate_dataset_items(fields=['row'], chunk_size=10, poll_interval=NO_WAIT)) + + assert items == shape_items(range(5)) + assert [request.get('fields') for request in api.items_requests] == ['row', 'row', 'row'] + + +async def test_iterate_dataset_items_forwards_options_without_extra_read_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: + """Item options reach every page request, and without a dropping filter an empty page ends with no plain read.""" + api = FakeRunApi([Step(pushed_rows=5, item_count=3, status='SUCCEEDED')]) + api.register(httpserver) + + items = [ + item + async for item in async_client.run(RUN_ID).iterate_dataset_items( + fields=['row'], chunk_size=10, poll_interval=NO_WAIT + ) + ] + + assert items == shape_items(range(5)) + assert [request.get('fields') for request in api.items_requests] == ['row', 'row', 'row'] + + def test_iterate_dataset_items_sleeps_between_polls_sync( httpserver: HTTPServer, sync_client: ApifyClient, monkeypatch: pytest.MonkeyPatch ) -> None: From 2addf46ecb6d6500731fa96baf1903027a20f1cf Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 08:33:31 +0200 Subject: [PATCH 5/8] feat: Long-poll the run between iterate_dataset_items polls --- src/apify_client/_resource_clients/run.py | 33 +++++++------- tests/unit/test_run_iterate_dataset_items.py | 47 ++++++++++++++------ 2 files changed, 50 insertions(+), 30 deletions(-) diff --git a/src/apify_client/_resource_clients/run.py b/src/apify_client/_resource_clients/run.py index 1c32b70d..54846f1c 100644 --- a/src/apify_client/_resource_clients/run.py +++ b/src/apify_client/_resource_clients/run.py @@ -1,6 +1,5 @@ from __future__ import annotations -import asyncio import json import random import string @@ -487,11 +486,11 @@ def iterate_dataset_items( ) -> Iterator[dict]: """Iterate over the items of the run's default dataset while the run is still producing them. - While the run has not finished, the dataset is polled every `poll_interval` and the rows below its - `item_count` are yielded. Each page is requested with a `limit` that ends at `item_count`, so it covers - exactly the rows it asks for, whatever the filters or `unwind` do to the items. `item_count` lags a few - seconds behind the pushed items, so once the run reaches a terminal status, the rows past it are read a page - at a time until none are left, and the iterator returns. + While the run has not finished, each poll yields the rows below the dataset's `item_count` and then waits up to + `poll_interval` for the run to finish, so the last rows are read as soon as it does. Each page is requested with + a `limit` that ends at `item_count`, so it covers exactly the rows it asks for, whatever the filters or `unwind` + do to the items. `item_count` lags a few seconds behind the pushed items, so once the run reaches a terminal + status, the rows past it are read a page at a time until none are left, and the iterator returns. https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items @@ -513,7 +512,7 @@ def iterate_dataset_items( skip_hidden: If True, then hidden fields are skipped from the output, i.e. fields starting with the # character. chunk_size: Maximum number of dataset rows requested per API call. - poll_interval: How long to wait between polls while the run has not finished. + poll_interval: How long to wait for the run to finish between polls. timeout: Timeout for each API HTTP request. Yields: @@ -537,8 +536,8 @@ def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: timeout=timeout, ) + run = self.get(timeout=timeout) while True: - run = self.get(timeout=timeout) is_finished = run is None or run.status in _TERMINAL_STATUSES dataset = dataset_client.get(timeout=timeout) item_count = dataset.item_count if dataset else 0 @@ -555,7 +554,7 @@ def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: return if is_finished: break - time.sleep(to_seconds(poll_interval)) + run = self.wait_for_finish(wait_duration=poll_interval, timeout=timeout) while True: page_limit = min(page_size, end - position) if end is not None else page_size @@ -1021,11 +1020,11 @@ async def iterate_dataset_items( ) -> AsyncIterator[dict]: """Iterate over the items of the run's default dataset while the run is still producing them. - While the run has not finished, the dataset is polled every `poll_interval` and the rows below its - `item_count` are yielded. Each page is requested with a `limit` that ends at `item_count`, so it covers - exactly the rows it asks for, whatever the filters or `unwind` do to the items. `item_count` lags a few - seconds behind the pushed items, so once the run reaches a terminal status, the rows past it are read a page - at a time until none are left, and the iterator returns. + While the run has not finished, each poll yields the rows below the dataset's `item_count` and then waits up to + `poll_interval` for the run to finish, so the last rows are read as soon as it does. Each page is requested with + a `limit` that ends at `item_count`, so it covers exactly the rows it asks for, whatever the filters or `unwind` + do to the items. `item_count` lags a few seconds behind the pushed items, so once the run reaches a terminal + status, the rows past it are read a page at a time until none are left, and the iterator returns. https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items @@ -1047,7 +1046,7 @@ async def iterate_dataset_items( skip_hidden: If True, then hidden fields are skipped from the output, i.e. fields starting with the # character. chunk_size: Maximum number of dataset rows requested per API call. - poll_interval: How long to wait between polls while the run has not finished. + poll_interval: How long to wait for the run to finish between polls. timeout: Timeout for each API HTTP request. Yields: @@ -1071,8 +1070,8 @@ async def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: timeout=timeout, ) + run = await self.get(timeout=timeout) while True: - run = await self.get(timeout=timeout) is_finished = run is None or run.status in _TERMINAL_STATUSES dataset = await dataset_client.get(timeout=timeout) item_count = dataset.item_count if dataset else 0 @@ -1090,7 +1089,7 @@ async def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: return if is_finished: break - await asyncio.sleep(to_seconds(poll_interval)) + run = await self.wait_for_finish(wait_duration=poll_interval, timeout=timeout) while True: page_limit = min(page_size, end - position) if end is not None else page_size diff --git a/tests/unit/test_run_iterate_dataset_items.py b/tests/unit/test_run_iterate_dataset_items.py index 558224ca..f52cea8c 100644 --- a/tests/unit/test_run_iterate_dataset_items.py +++ b/tests/unit/test_run_iterate_dataset_items.py @@ -15,6 +15,7 @@ from apify_client import ApifyClient, ApifyClientAsync from apify_client._literals import ActorJobStatus + from apify_client._models import Run pytestmark = pytest.mark.usefixtures('http_client_classes') @@ -43,6 +44,7 @@ class FakeRunApi: steps: list[Step] step_index: int = -1 + run_requests: list[dict[str, str]] = field(default_factory=list) items_requests: list[dict[str, str]] = field(default_factory=list) @property @@ -54,7 +56,8 @@ def register(self, httpserver: HTTPServer) -> None: httpserver.expect_request(f'{RUN_PATH}/dataset', method='GET').respond_with_handler(self.handle_dataset) httpserver.expect_request(f'{RUN_PATH}/dataset/items', method='GET').respond_with_handler(self.handle_items) - def handle_run(self, _request: Request) -> Response: + def handle_run(self, request: Request) -> Response: + self.run_requests.append(dict(request.args)) self.step_index = min(self.step_index + 1, len(self.steps) - 1) run = { 'id': RUN_ID, @@ -358,29 +361,47 @@ async def test_iterate_dataset_items_forwards_options_without_extra_read_async( assert [request.get('fields') for request in api.items_requests] == ['row', 'row', 'row'] -def test_iterate_dataset_items_sleeps_between_polls_sync( +def test_iterate_dataset_items_waits_for_finish_between_polls_sync( httpserver: HTTPServer, sync_client: ApifyClient, monkeypatch: pytest.MonkeyPatch ) -> None: - """The iterator waits `poll_interval` after each poll of an unfinished run and not after the final one.""" + """Each poll of an unfinished run waits up to `poll_interval` for the run to finish, and the final one does not.""" api = FakeRunApi(LAGGING_RUN_STEPS) api.register(httpserver) - sleep = Mock() - monkeypatch.setattr('apify_client._resource_clients.run.time.sleep', sleep) + run_client = sync_client.run(RUN_ID) + wait_with_holding = run_client.wait_for_finish - list(sync_client.run(RUN_ID).iterate_dataset_items(poll_interval=timedelta(seconds=2))) + # The fake API answers at once, so a real wait would re-read the run until the deadline and skip steps. + def wait_without_holding(**kwargs: Any) -> Run | None: + return wait_with_holding(**{**kwargs, 'wait_duration': NO_WAIT}) - assert sleep.call_args_list == [call(2.0)] * 3 + wait_for_finish = Mock(side_effect=wait_without_holding) + monkeypatch.setattr(run_client, 'wait_for_finish', wait_for_finish) + items = list(run_client.iterate_dataset_items(poll_interval=timedelta(seconds=2))) -async def test_iterate_dataset_items_sleeps_between_polls_async( + assert items == shape_items(range(75)) + assert wait_for_finish.call_args_list == [call(wait_duration=timedelta(seconds=2), timeout='long')] * 3 + assert api.run_requests == [{}] + [{'waitForFinish': '0'}] * 3 + + +async def test_iterate_dataset_items_waits_for_finish_between_polls_async( httpserver: HTTPServer, async_client: ApifyClientAsync, monkeypatch: pytest.MonkeyPatch ) -> None: - """The iterator waits `poll_interval` after each poll of an unfinished run and not after the final one.""" + """Each poll of an unfinished run waits up to `poll_interval` for the run to finish, and the final one does not.""" api = FakeRunApi(LAGGING_RUN_STEPS) api.register(httpserver) - sleep = AsyncMock() - monkeypatch.setattr('apify_client._resource_clients.run.asyncio.sleep', sleep) + run_client = async_client.run(RUN_ID) + wait_with_holding = run_client.wait_for_finish - [item async for item in async_client.run(RUN_ID).iterate_dataset_items(poll_interval=timedelta(seconds=2))] + # The fake API answers at once, so a real wait would re-read the run until the deadline and skip steps. + async def wait_without_holding(**kwargs: Any) -> Run | None: + return await wait_with_holding(**{**kwargs, 'wait_duration': NO_WAIT}) - assert sleep.call_args_list == [call(2.0)] * 3 + wait_for_finish = AsyncMock(side_effect=wait_without_holding) + monkeypatch.setattr(run_client, 'wait_for_finish', wait_for_finish) + + items = [item async for item in run_client.iterate_dataset_items(poll_interval=timedelta(seconds=2))] + + assert items == shape_items(range(75)) + assert wait_for_finish.call_args_list == [call(wait_duration=timedelta(seconds=2), timeout='long')] * 3 + assert api.run_requests == [{}] + [{'waitForFinish': '0'}] * 3 From d044643c30745071d05aacb204c5cbb7b3025e7d Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 08:33:32 +0200 Subject: [PATCH 6/8] docs: Cover live iteration over a run's dataset items --- docs/02_concepts/08_pagination.mdx | 2 ++ docs/03_guides/03_retrieve_actor_data.mdx | 30 ++++++++++++++++++++++ docs/03_guides/code/03_live_items_async.py | 22 ++++++++++++++++ docs/03_guides/code/03_live_items_sync.py | 20 +++++++++++++++ 4 files changed, 74 insertions(+) create mode 100644 docs/03_guides/code/03_live_items_async.py create mode 100644 docs/03_guides/code/03_live_items_sync.py diff --git a/docs/02_concepts/08_pagination.mdx b/docs/02_concepts/08_pagination.mdx index 757f51c2..86951f66 100644 --- a/docs/02_concepts/08_pagination.mdx +++ b/docs/02_concepts/08_pagination.mdx @@ -83,3 +83,5 @@ The next example uses `iterate_items` on a dataset client to stream items past a + +To read the items of a run that's still going, use `RunClient.iterate_dataset_items`. Its `offset` and `limit` also count dataset rows. For details, see [Retrieve Actor data](/api/client/python/docs/guides/retrieve-actor-data#read-items-while-the-run-is-going). diff --git a/docs/03_guides/03_retrieve_actor_data.mdx b/docs/03_guides/03_retrieve_actor_data.mdx index 7490f114..e430ef03 100644 --- a/docs/03_guides/03_retrieve_actor_data.mdx +++ b/docs/03_guides/03_retrieve_actor_data.mdx @@ -8,6 +8,10 @@ import Tabs from '@theme/Tabs'; import TabItem from '@theme/TabItem'; import CodeBlock from '@theme/CodeBlock'; +import ApiLink from '@theme/ApiLink'; + +import LiveItemsAsyncExample from '!!raw-loader!./code/03_live_items_async.py'; +import LiveItemsSyncExample from '!!raw-loader!./code/03_live_items_sync.py'; import RetrieveAsyncExample from '!!raw-loader!./code/03_retrieve_async.py'; import RetrieveSyncExample from '!!raw-loader!./code/03_retrieve_sync.py'; @@ -27,3 +31,29 @@ The following example shows how to fetch datasets from an Actor's runs, paginate + +## Read items while the run is going + +To process a run's output before the run finishes, use `RunClient.iterate_dataset_items`. It yields the items of the run's default dataset shortly after the run pushes them, and it returns once the run has finished and every item is read. + +Note that: + +- Between polls, the iterator waits up to `poll_interval` (5 seconds by default) for the run to finish. Once it finishes, the iterator reads the remaining items right away. An explicit `timeout` has to leave room for that wait. +- The item options are the same as in `DatasetClient.iterate_items`, except for `desc` and `signature`. +- `offset` and `limit` count dataset rows, not the items you get back. With `clean`, `skip_empty` or `unwind`, the iterator can yield fewer or more items than `limit`. +- A run that's `ABORTING` or `TIMING-OUT` can still push items, so the iterator keeps polling until the run reaches a terminal status. + +The following example starts an Actor and prints its items as the run produces them: + + + + + {LiveItemsAsyncExample} + + + + + {LiveItemsSyncExample} + + + diff --git a/docs/03_guides/code/03_live_items_async.py b/docs/03_guides/code/03_live_items_async.py new file mode 100644 index 00000000..11030d77 --- /dev/null +++ b/docs/03_guides/code/03_live_items_async.py @@ -0,0 +1,22 @@ +import asyncio + +from apify_client import ApifyClientAsync + +TOKEN = 'MY-APIFY-TOKEN' + + +async def main() -> None: + apify_client = ApifyClientAsync(TOKEN) + + # Start the Actor without waiting for it to finish + actor_client = apify_client.actor('username/actor-name') + run = await actor_client.start(run_input={'query': 'web scraping'}) + + # Each item arrives shortly after the run pushes it. The loop ends once the run + # has finished and every item is read. + async for item in apify_client.run(run.id).iterate_dataset_items(skip_empty=True): + print(item) + + +if __name__ == '__main__': + asyncio.run(main()) diff --git a/docs/03_guides/code/03_live_items_sync.py b/docs/03_guides/code/03_live_items_sync.py new file mode 100644 index 00000000..117c71ee --- /dev/null +++ b/docs/03_guides/code/03_live_items_sync.py @@ -0,0 +1,20 @@ +from apify_client import ApifyClient + +TOKEN = 'MY-APIFY-TOKEN' + + +def main() -> None: + apify_client = ApifyClient(TOKEN) + + # Start the Actor without waiting for it to finish + actor_client = apify_client.actor('username/actor-name') + run = actor_client.start(run_input={'query': 'web scraping'}) + + # Each item arrives shortly after the run pushes it. The loop ends once the run + # has finished and every item is read. + for item in apify_client.run(run.id).iterate_dataset_items(skip_empty=True): + print(item) + + +if __name__ == '__main__': + main() From 3f4961e5a1d3f05c49062d6935f84abdd05646a6 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Thu, 1 Oct 2026 08:40:24 +0200 Subject: [PATCH 7/8] fix: Pin the run last_run() resolves to in iterate_dataset_items --- .../_resource_clients/_resource_client.py | 12 ++++- src/apify_client/_resource_clients/run.py | 44 +++++++++++++++---- tests/unit/test_run_iterate_dataset_items.py | 37 +++++++++++++++- 3 files changed, 83 insertions(+), 10 deletions(-) diff --git a/src/apify_client/_resource_clients/_resource_client.py b/src/apify_client/_resource_clients/_resource_client.py index ecced233..bb5d85fd 100644 --- a/src/apify_client/_resource_clients/_resource_client.py +++ b/src/apify_client/_resource_clients/_resource_client.py @@ -42,6 +42,7 @@ def __init__( client_registry: Any, resource_id: str | None = None, params: dict | None = None, + api_base_url: str | None = None, ) -> None: """Initialize the resource client. @@ -53,11 +54,13 @@ def __init__( client_registry: Bundle of client classes for dependency injection. resource_id: Optional resource ID for single-resource clients. params: Optional default parameters for all requests. + api_base_url: Base URL of the API itself, for clients of top-level resources. Defaults to `base_url`. """ if resource_path.endswith('/'): raise ValueError('resource_path must not end with "/"') self._base_url = base_url + self._api_base_url = api_base_url or base_url self._public_base_url = public_base_url self._http_client = http_client self._default_params = params or {} @@ -82,11 +85,12 @@ def _resource_url(self) -> str: def _base_client_kwargs(self) -> dict[str, Any]: """Base kwargs for creating nested/child clients. - Returns dict with base_url, public_base_url, http_client, and client_registry. Caller adds + Returns dict with base_url, api_base_url, public_base_url, http_client, and client_registry. Caller adds resource_path, resource_id, and params as needed. """ return { 'base_url': self._resource_url, + 'api_base_url': self._api_base_url, 'public_base_url': self._public_base_url, 'http_client': self._http_client, 'client_registry': self._client_registry, @@ -197,6 +201,7 @@ def __init__( client_registry: ClientRegistry, resource_id: str | None = None, params: dict | None = None, + api_base_url: str | None = None, ) -> None: """Initialize the resource client. @@ -208,6 +213,7 @@ def __init__( client_registry: Bundle of client classes for dependency injection. resource_id: Optional resource ID for single-resource clients. params: Optional default parameters for all requests. + api_base_url: Base URL of the API itself, for clients of top-level resources. Defaults to `base_url`. """ super().__init__( base_url=base_url, @@ -217,6 +223,7 @@ def __init__( client_registry=client_registry, resource_id=resource_id, params=params, + api_base_url=api_base_url, ) def _get(self, *, timeout: Timeout) -> dict | None: @@ -389,6 +396,7 @@ def __init__( client_registry: ClientRegistryAsync, resource_id: str | None = None, params: dict | None = None, + api_base_url: str | None = None, ) -> None: """Initialize the resource client. @@ -400,6 +408,7 @@ def __init__( client_registry: Bundle of client classes for dependency injection. resource_id: Optional resource ID for single-resource clients. params: Optional default parameters for all requests. + api_base_url: Base URL of the API itself, for clients of top-level resources. Defaults to `base_url`. """ super().__init__( base_url=base_url, @@ -409,6 +418,7 @@ def __init__( client_registry=client_registry, resource_id=resource_id, params=params, + api_base_url=api_base_url, ) async def _get(self, *, timeout: Timeout) -> dict | None: diff --git a/src/apify_client/_resource_clients/run.py b/src/apify_client/_resource_clients/run.py index 54846f1c..22b34be6 100644 --- a/src/apify_client/_resource_clients/run.py +++ b/src/apify_client/_resource_clients/run.py @@ -490,7 +490,8 @@ def iterate_dataset_items( `poll_interval` for the run to finish, so the last rows are read as soon as it does. Each page is requested with a `limit` that ends at `item_count`, so it covers exactly the rows it asks for, whatever the filters or `unwind` do to the items. `item_count` lags a few seconds behind the pushed items, so once the run reaches a terminal - status, the rows past it are read a page at a time until none are left, and the iterator returns. + status, the rows past it are read a page at a time until none are left, and the iterator returns. On a + `last_run()` client, the iterator sticks to the run that its first request resolves to. https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items @@ -518,11 +519,25 @@ def iterate_dataset_items( Yields: An item from the dataset. """ - dataset_client = self.dataset() page_size = chunk_size or DEFAULT_CHUNK_SIZE position = offset or 0 end = position + limit if limit else None + run = self.get(timeout=timeout) + # A `last_run()` client resolves `runs/last` per request, so a newer run would swap the dataset mid-iteration. + run_client = ( + self._client_registry.run_client( + resource_id=run.id, + base_url=self._api_base_url, + public_base_url=self._public_base_url, + http_client=self._http_client, + client_registry=self._client_registry, + ) + if run is not None and run.id != self._resource_id + else self + ) + dataset_client = run_client.dataset() + def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: return dataset_client.list_items( offset=page_offset, @@ -536,7 +551,6 @@ def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: timeout=timeout, ) - run = self.get(timeout=timeout) while True: is_finished = run is None or run.status in _TERMINAL_STATUSES dataset = dataset_client.get(timeout=timeout) @@ -554,7 +568,7 @@ def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: return if is_finished: break - run = self.wait_for_finish(wait_duration=poll_interval, timeout=timeout) + run = run_client.wait_for_finish(wait_duration=poll_interval, timeout=timeout) while True: page_limit = min(page_size, end - position) if end is not None else page_size @@ -1024,7 +1038,8 @@ async def iterate_dataset_items( `poll_interval` for the run to finish, so the last rows are read as soon as it does. Each page is requested with a `limit` that ends at `item_count`, so it covers exactly the rows it asks for, whatever the filters or `unwind` do to the items. `item_count` lags a few seconds behind the pushed items, so once the run reaches a terminal - status, the rows past it are read a page at a time until none are left, and the iterator returns. + status, the rows past it are read a page at a time until none are left, and the iterator returns. On a + `last_run()` client, the iterator sticks to the run that its first request resolves to. https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items @@ -1052,11 +1067,25 @@ async def iterate_dataset_items( Yields: An item from the dataset. """ - dataset_client = self.dataset() page_size = chunk_size or DEFAULT_CHUNK_SIZE position = offset or 0 end = position + limit if limit else None + run = await self.get(timeout=timeout) + # A `last_run()` client resolves `runs/last` per request, so a newer run would swap the dataset mid-iteration. + run_client = ( + self._client_registry.run_client( + resource_id=run.id, + base_url=self._api_base_url, + public_base_url=self._public_base_url, + http_client=self._http_client, + client_registry=self._client_registry, + ) + if run is not None and run.id != self._resource_id + else self + ) + dataset_client = run_client.dataset() + async def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: return await dataset_client.list_items( offset=page_offset, @@ -1070,7 +1099,6 @@ async def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: timeout=timeout, ) - run = await self.get(timeout=timeout) while True: is_finished = run is None or run.status in _TERMINAL_STATUSES dataset = await dataset_client.get(timeout=timeout) @@ -1089,7 +1117,7 @@ async def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: return if is_finished: break - run = await self.wait_for_finish(wait_duration=poll_interval, timeout=timeout) + run = await run_client.wait_for_finish(wait_duration=poll_interval, timeout=timeout) while True: page_limit = min(page_size, end - position) if end is not None else page_size diff --git a/tests/unit/test_run_iterate_dataset_items.py b/tests/unit/test_run_iterate_dataset_items.py index f52cea8c..3e2cd900 100644 --- a/tests/unit/test_run_iterate_dataset_items.py +++ b/tests/unit/test_run_iterate_dataset_items.py @@ -21,6 +21,7 @@ RUN_ID = 'test-run-id' RUN_PATH = f'/v2/actor-runs/{RUN_ID}' +ACTOR_ID = 'test-actor-id' UNWIND_PARTS = 3 NO_WAIT = timedelta(0) @@ -45,6 +46,7 @@ class FakeRunApi: steps: list[Step] step_index: int = -1 run_requests: list[dict[str, str]] = field(default_factory=list) + last_run_requests: int = 0 items_requests: list[dict[str, str]] = field(default_factory=list) @property @@ -55,13 +57,16 @@ def register(self, httpserver: HTTPServer) -> None: httpserver.expect_request(RUN_PATH, method='GET').respond_with_handler(self.handle_run) httpserver.expect_request(f'{RUN_PATH}/dataset', method='GET').respond_with_handler(self.handle_dataset) httpserver.expect_request(f'{RUN_PATH}/dataset/items', method='GET').respond_with_handler(self.handle_items) + httpserver.expect_request(f'/v2/actors/{ACTOR_ID}/runs/last', method='GET').respond_with_handler( + self.handle_last_run + ) def handle_run(self, request: Request) -> Response: self.run_requests.append(dict(request.args)) self.step_index = min(self.step_index + 1, len(self.steps) - 1) run = { 'id': RUN_ID, - 'actId': 'test-actor-id', + 'actId': ACTOR_ID, 'userId': 'test-user-id', 'startedAt': '2019-11-30T07:34:24.202Z', 'status': self.step.status, @@ -76,6 +81,10 @@ def handle_run(self, request: Request) -> Response: } return Response(json.dumps({'data': run}), status=200, mimetype='application/json') + def handle_last_run(self, request: Request) -> Response: + self.last_run_requests += 1 + return self.handle_run(request) + def handle_dataset(self, _request: Request) -> Response: dataset = { 'id': 'test-dataset-id', @@ -405,3 +414,29 @@ async def wait_without_holding(**kwargs: Any) -> Run | None: assert items == shape_items(range(75)) assert wait_for_finish.call_args_list == [call(wait_duration=timedelta(seconds=2), timeout='long')] * 3 assert api.run_requests == [{}] + [{'waitForFinish': '0'}] * 3 + + +def test_iterate_dataset_items_pins_the_last_run_sync(httpserver: HTTPServer, sync_client: ApifyClient) -> None: + """A `last_run()` client resolves `runs/last` once and reads that run's dataset to the end.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + run_client = sync_client.actor(ACTOR_ID).last_run() + items = list(run_client.iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)) + + assert items == shape_items(range(75)) + assert api.last_run_requests == 1 + + +async def test_iterate_dataset_items_pins_the_last_run_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: + """A `last_run()` client resolves `runs/last` once and reads that run's dataset to the end.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + run_client = async_client.actor(ACTOR_ID).last_run() + items = [item async for item in run_client.iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)] + + assert items == shape_items(range(75)) + assert api.last_run_requests == 1 From 2254321489ce258ae62cbd0411897781137bc9e7 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Thu, 1 Oct 2026 08:49:33 +0200 Subject: [PATCH 8/8] test: Cover iterate_dataset_items on last_run() against the live API --- tests/integration/test_run.py | 92 ++++++++++++++++++++++++++++++++++- 1 file changed, 90 insertions(+), 2 deletions(-) diff --git a/tests/integration/test_run.py b/tests/integration/test_run.py index 811c5794..1cea548d 100644 --- a/tests/integration/test_run.py +++ b/tests/integration/test_run.py @@ -6,8 +6,8 @@ from datetime import UTC, datetime, timedelta from typing import TYPE_CHECKING -from .._utils import maybe_await, poll_until_condition -from apify_client._models import Dataset, KeyValueStore, ListOfRuns, RequestQueue, Run, RunShort +from .._utils import get_random_resource_name, maybe_await, poll_until_condition +from apify_client._models import Actor, Build, Dataset, KeyValueStore, ListOfRuns, RequestQueue, Run, RunShort from apify_client.errors import ApifyApiError if TYPE_CHECKING: @@ -15,6 +15,34 @@ HELLO_WORLD_ACTOR = 'apify/hello-world' +LIVE_ITEM_COUNT = 10 + +LIVE_ITEMS_SOURCE_FILES = [ + { + 'name': 'Dockerfile', + 'format': 'TEXT', + 'content': 'FROM apify/actor-node:22\nCOPY . ./\nCMD ["node", "main.mjs"]\n', + }, + { + 'name': 'main.mjs', + 'format': 'TEXT', + 'content': f""" +const apiUrl = (process.env.APIFY_API_BASE_URL || 'https://api.apify.com').replace(/\\/$/, ''); +const itemsUrl = `${{apiUrl}}/v2/datasets/${{process.env.ACTOR_DEFAULT_DATASET_ID}}/items`; +for (let index = 0; index < {LIVE_ITEM_COUNT}; index++) {{ + const response = await fetch(itemsUrl, {{ + method: 'POST', + headers: {{ 'Content-Type': 'application/json', Authorization: `Bearer ${{process.env.APIFY_TOKEN}}` }}, + body: JSON.stringify({{ runId: process.env.ACTOR_RUN_ID, index }}), + }}); + if (!response.ok) throw new Error(`Pushing item ${{index}} failed with ${{response.status}}`); + await new Promise((resolve) => setTimeout(resolve, 1000)); +}} +""", + }, +] +"""Source of an Actor that pushes `LIVE_ITEM_COUNT` items tagged with its run ID to its dataset, one per second.""" + async def test_run_collection_list_multiple_statuses(client: ApifyClient | ApifyClientAsync) -> None: """Test listing runs with multiple statuses.""" @@ -452,3 +480,63 @@ async def test_run_collection_iterate_actor_runs(client: ApifyClient | ApifyClie assert all(r.act_id == run.act_id for r in collected) finally: await maybe_await(client.run(run.id).delete()) + + +async def test_run_iterate_dataset_items_of_last_run(client: ApifyClient | ApifyClientAsync, *, is_async: bool) -> None: + """`iterate_dataset_items` on `last_run()` reads the run it resolved to, even after a newer run starts.""" + created_actor = await maybe_await( + client.actors().create( + name=get_random_resource_name('actor'), + versions=[ + { + 'versionNumber': '0.0', + 'sourceType': 'SOURCE_FILES', + 'buildTag': 'latest', + 'sourceFiles': LIVE_ITEMS_SOURCE_FILES, + } + ], + ) + ) + assert isinstance(created_actor, Actor) + actor_client = client.actor(created_actor.id) + run_ids: list[str] = [] + + try: + started_build = await maybe_await(actor_client.build(version_number='0.0')) + assert isinstance(started_build, Build) + build = await maybe_await(client.build(started_build.id).wait_for_finish()) + assert isinstance(build, Build) + assert build.status == 'SUCCEEDED' + + async def start_run() -> str: + run = await maybe_await(actor_client.start(memory_mbytes=256, run_timeout=timedelta(seconds=120))) + assert isinstance(run, Run) + return run.id + + first_run_id = await start_run() + run_ids.append(first_run_id) + + items: list[dict] = [] + + async def collect(item: dict) -> None: + items.append(item) + if len(run_ids) == 1: + run_ids.append(await start_run()) + + # A page of one row makes every further read a fresh request, which could land on the newer run. + iterator = actor_client.last_run().iterate_dataset_items(chunk_size=1, poll_interval=timedelta(seconds=1)) + if is_async: + assert isinstance(iterator, AsyncIterator) + async for item in iterator: + await collect(item) + else: + assert isinstance(iterator, Iterator) + for item in iterator: + await collect(item) + + assert items == [{'runId': first_run_id, 'index': index} for index in range(LIVE_ITEM_COUNT)] + finally: + for run_id in run_ids: + await maybe_await(client.run(run_id).wait_for_finish()) + await maybe_await(client.run(run_id).delete()) + await maybe_await(actor_client.delete())