Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions docs/02_concepts/06_interacting_with_other_actors.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import InteractingCallExample from '!!raw-loader!roa-loader!./code/06_interactin
import InteractingNamedCallExample from '!!raw-loader!roa-loader!./code/06_interacting_named_call.py';
import InteractingChildRunsExample from '!!raw-loader!roa-loader!./code/06_interacting_child_runs.py';
import InteractingAbortWithParentExample from '!!raw-loader!roa-loader!./code/06_interacting_abort_with_parent.py';
import InteractingChildRunLimitExample from '!!raw-loader!roa-loader!./code/06_interacting_child_run_limit.py';
import InteractingCallTaskExample from '!!raw-loader!roa-loader!./code/06_interacting_call_task.py';
import InteractingMetamorphExample from '!!raw-loader!roa-loader!./code/06_interacting_metamorph.py';
import InteractingAbortExample from '!!raw-loader!roa-loader!./code/06_interacting_abort.py';
Expand Down Expand Up @@ -86,6 +87,23 @@ Note that:
- Child runs started after your Actor run received `ABORTING` aren't aborted, so don't start new ones while it's shutting down.
- A child run started with its own `token` is aborted with that token. After a migration or resurrection, the SDK uses your Actor's token for it until the same named call runs again. If that token can't access the child run, the abort fails and the error is logged.

### Limiting concurrent child runs

An Actor that starts many child runs at once can hit the concurrency or memory limit of your account. To cap how many named child runs are active at once, call <ApiLink to="class/Actor#set_child_run_limits">`Actor.set_child_run_limits`</ApiLink>. While the limit is reached, a named `Actor.start` or `Actor.call` that would start or resurrect a run waits until one of the active child runs finishes. The SDK counts the child runs from the registry, so child runs started before a migration or resurrection count too.

<RunnableCodeBlock className="language-python" language="python">
{InteractingChildRunLimitExample}
</RunnableCodeBlock>

Note that:

- A child run counts as active while it's `READY`, `RUNNING`, `ABORTING` or `TIMING-OUT`.
- Only named child runs count, and only named calls wait. A call without `name` starts its run right away.
- Reattaching to an active child run never waits, since the run already holds a slot.
- `Actor.call` frees the slot as soon as its run finishes. The SDK doesn't learn right away about a child run that nothing waits for, so it fetches the run again before counting it, once its status is more than 10 seconds old.
- The limit isn't persisted. After a migration or resurrection, call `Actor.set_child_run_limits` again before you start child runs.
- Once your Actor run receives the `ABORTING` event, a call waiting for a slot raises a `RuntimeError` without starting its run.

## Actor call task

The <ApiLink to="class/Actor#call_task">`Actor.call_task`</ApiLink> method starts an [Actor task](https://docs.apify.com/platform/actors/tasks) on the Apify platform, and waits for the started Actor run to finish.
Expand Down
30 changes: 30 additions & 0 deletions docs/02_concepts/code/06_interacting_child_run_limit.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
import asyncio

from apify import Actor


async def main() -> None:
async with Actor:
# Keep at most 3 named child runs active at once.
Actor.set_child_run_limits(max_concurrent_runs=3)

urls = [f'https://www.apify.com/?page={page}' for page in range(6)]

# Each call waits for a free slot before it starts its child run.
actor_runs = await asyncio.gather(
*(
Actor.call(
actor_id='apify/screenshot-url',
run_input={'urls': [{'url': url}]},
name=f'screenshot-{index}',
)
for index, url in enumerate(urls)
)
)

for actor_run in actor_runs:
Actor.log.info(f'Child run {actor_run.id} finished as {actor_run.status}')


if __name__ == '__main__':
asyncio.run(main())
17 changes: 17 additions & 0 deletions src/apify/_actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -1193,6 +1193,8 @@ async def call(
run = await self._wait_for_child_run(
client.run(started_run.id), started_run, wait=wait, logger=logger, from_start=is_new
)
if run is not None:
await self._child_run_registry.run_finished(name, run)

if run is None:
raise RuntimeError(f'Failed to call Actor with ID "{actor_id}".')
Expand Down Expand Up @@ -1270,6 +1272,21 @@ async def child_runs(self) -> dict[str, ChildRunInfo]:
"""
return await self._child_run_registry.list_runs(self.apify_client)

def set_child_run_limits(self, *, max_concurrent_runs: int | None) -> None:
"""Limit the named child runs of this Actor run.

While `max_concurrent_runs` named child runs are `READY`, `RUNNING`, `ABORTING` or `TIMING-OUT`, a named
`Actor.start` or `Actor.call` that would start or resurrect a run waits until one of them finishes. Reattaching
to a recorded run never waits. Runs started without a `name` are not counted and never wait. A child run not
awaited by `Actor.call` is fetched again before it is counted, if its status is more than 10 seconds old.

The limit is kept in memory, so call this method again after a migration or resurrection of this Actor run.

Args:
max_concurrent_runs: How many named child runs may be active at once, or `None` for no limit.
"""
self._child_run_registry.set_max_concurrent_runs(max_concurrent_runs)

@_ensure_context
async def call_task(
self,
Expand Down
149 changes: 132 additions & 17 deletions src/apify/_child_runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@

import asyncio
from collections import defaultdict
from contextlib import asynccontextmanager, suppress
from dataclasses import dataclass
from datetime import timedelta
from logging import getLogger
from typing import TYPE_CHECKING

Expand All @@ -12,7 +14,7 @@
from apify._utils import docs_group

if TYPE_CHECKING:
from collections.abc import Awaitable, Callable
from collections.abc import AsyncIterator, Awaitable, Callable

from apify_client import ApifyClientAsync
from apify_client._models import Run
Expand All @@ -32,6 +34,11 @@

_ABORTABLE_STATUSES = frozenset({'READY', 'RUNNING'})

_ACTIVE_STATUSES = frozenset({'READY', 'RUNNING', 'ABORTING', 'TIMING-OUT'})

_STATUS_MAX_AGE = timedelta(seconds=10)
"""How long an observed active status counts toward the concurrency limit before the run is fetched again."""


class ChildRunRecord(BaseModel):
"""A child run tracked under a name in the child run registry."""
Expand Down Expand Up @@ -90,6 +97,19 @@ def __init__(self, open_key_value_store: Callable[[], Awaitable[KeyValueStore]])
self._name_locks: defaultdict[str, asyncio.Lock] = defaultdict(asyncio.Lock)
self._clients: dict[str, ApifyClientAsync] = {}
"""Client each name was last started or reattached with in this process, used to abort its run."""
self._max_concurrent_runs: int | None = None
self._slots = asyncio.Condition()
self._starting: set[str] = set()
"""Names holding a slot for a start or resurrection in flight, not counted from the records yet."""
self._observed: dict[str, tuple[str, float]] = {}
"""Last observed status of the run recorded under each name, with the event loop time it was observed at."""
self._parent_aborting = False

def set_max_concurrent_runs(self, max_concurrent_runs: int | None) -> None:
"""Set how many recorded runs may be active at once, or remove the limit with `None`."""
if max_concurrent_runs is not None and max_concurrent_runs < 1:
raise ValueError(f'`max_concurrent_runs` must be at least 1, got {max_concurrent_runs}.')
self._max_concurrent_runs = max_concurrent_runs

async def find_or_start(
self,
Expand All @@ -107,6 +127,8 @@ async def find_or_start(
An `ABORTED` or `TIMED-OUT` run is resurrected, since Actors are expected to resume from their state.
A `FAILED` run, or one the API no longer knows, is replaced by a new run under the same name.

Starting or resurrecting a run waits while the concurrency limit is reached. Reattaching never waits.

Args:
name: Name of the child run, unique within the parent run.
actor_id: The Actor to start. It must match the Actor already recorded under `name`.
Expand All @@ -132,13 +154,14 @@ async def find_or_start(
self._clients[name] = client

if record is None:
run = await self._start(
name,
actor_id=actor_id,
start_run=start_run,
previous_run_ids=[],
abort_with_parent=abort_with_parent,
)
async with self._slot(name, client):
run = await self._start(
name,
actor_id=actor_id,
start_run=start_run,
previous_run_ids=[],
abort_with_parent=abort_with_parent,
)
return run, True

run_client = client.run(record.run_id)
Expand All @@ -148,25 +171,41 @@ async def find_or_start(
run = await run_client.wait_for_finish()

if run is None or run.status == 'FAILED':
run = await self._start(
name,
actor_id=actor_id,
start_run=start_run,
previous_run_ids=[*record.previous_run_ids, record.run_id],
abort_with_parent=abort_with_parent,
)
async with self._slot(name, client):
run = await self._start(
name,
actor_id=actor_id,
start_run=start_run,
previous_run_ids=[*record.previous_run_ids, record.run_id],
abort_with_parent=abort_with_parent,
)
return run, True

self._observe(name, run)

if record.abort_with_parent != abort_with_parent:
await self._save(name, record.model_copy(update={'abort_with_parent': abort_with_parent}))

if run.status in _RESURRECTABLE_STATUSES:
logger.info(f'Resurrecting child run "{name}"', extra={'run_id': run.id, 'status': run.status})
return await resurrect_run(run_client), False
async with self._slot(name, client):
logger.info(f'Resurrecting child run "{name}"', extra={'run_id': run.id, 'status': run.status})
run = await resurrect_run(run_client)
self._observe(name, run)
return run, False

logger.info(f'Reattaching to child run "{name}"', extra={'run_id': run.id, 'status': run.status})
return run, False

async def run_finished(self, name: str, run: Run) -> None:
"""Record the status of a run under `name` that was awaited, releasing its slot when it is no longer active."""
records = await self._load()
record = records.get(name)
if record is None or record.run_id != run.id:
return
self._observe(name, run)
async with self._slots:
self._slots.notify_all()

async def list_runs(self, client: ApifyClientAsync) -> dict[str, ChildRunInfo]:
"""Return every recorded child run by name, with its current state fetched from the API.

Expand All @@ -176,6 +215,11 @@ async def list_runs(self, client: ApifyClientAsync) -> dict[str, ChildRunInfo]:
# Copy the records, since a named start can add one while the runs are fetched.
records = dict(await self._load())
runs = await asyncio.gather(*(client.run(record.run_id).get() for record in records.values()))
# A name whose run was replaced during the fetch keeps the status observed for its new run.
current = await self._load()
for name, run in zip(records, runs, strict=True):
if run is not None and name in current and current[name].run_id == run.id:
self._observe(name, run)
return {
name: ChildRunInfo(
actor_id=record.actor_id,
Expand All @@ -195,6 +239,10 @@ async def abort_runs_with_parent(self, client: ApifyClientAsync) -> None:
Args:
client: Client used for a name not started or reattached in this process, e.g. after a migration.
"""
self._parent_aborting = True
async with self._slots:
self._slots.notify_all()

records = await self._load()
# Names with a start in flight are not recorded yet, so their locks are awaited too.
await asyncio.gather(*(self._abort(name, client) for name in {*records, *self._name_locks}))
Expand Down Expand Up @@ -232,8 +280,75 @@ async def _start(
abort_with_parent=abort_with_parent,
)
await self._save(name, record)
self._observe(name, run)
return run

@asynccontextmanager
async def _slot(self, name: str, client: ApifyClientAsync) -> AsyncIterator[None]:
"""Hold a slot for starting or resurrecting the run under `name`, waiting while the limit is reached."""
if self._max_concurrent_runs is None:
yield
return

async with self._slots:
while (
self._max_concurrent_runs is not None
and await self._count_active(client, exclude=name) >= self._max_concurrent_runs
):
if self._parent_aborting:
raise RuntimeError(
f'Child run "{name}" was not started, since this Actor run is being aborted and the limit '
f'of {self._max_concurrent_runs} concurrent child runs is reached.'
)
logger.debug(f'Child run "{name}" is waiting for a free slot')
with suppress(TimeoutError):
await asyncio.wait_for(self._slots.wait(), timeout=_STATUS_MAX_AGE.total_seconds())
self._starting.add(name)

try:
yield
finally:
self._starting.discard(name)
async with self._slots:
self._slots.notify_all()

async def _count_active(self, client: ApifyClientAsync, *, exclude: str) -> int:
"""Count active recorded runs and slots held by others, fetching runs whose active status is not fresh."""
records = await self._load()
names = [name for name in records if name != exclude and name not in self._starting]
now = asyncio.get_running_loop().time()
stale = [
name
for name in names
if name not in self._observed
or (
self._observed[name][0] in _ACTIVE_STATUSES
and now - self._observed[name][1] >= _STATUS_MAX_AGE.total_seconds()
)
]
runs = await asyncio.gather(
*(self._clients.get(name, client).run(records[name].run_id).get() for name in stale),
return_exceptions=True,
)
for name, run in zip(stale, runs, strict=True):
if isinstance(run, BaseException):
logger.warning(
f'Failed to fetch child run "{name}" to count it toward the concurrency limit',
extra={'run_id': records[name].run_id},
exc_info=run,
)
elif run is None:
self._observed[name] = ('MISSING', now)
else:
self._observe(name, run)

# A run that could not be fetched counts only when an earlier observation saw it active.
active = sum(1 for name in names if name in self._observed and self._observed[name][0] in _ACTIVE_STATUSES)
return active + len(self._starting)

def _observe(self, name: str, run: Run) -> None:
self._observed[name] = (run.status, asyncio.get_running_loop().time())

async def _load(self) -> dict[str, ChildRunRecord]:
async with self._load_lock:
if self._records is None:
Expand Down
33 changes: 33 additions & 0 deletions tests/e2e/test_actor_child_runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -134,3 +134,36 @@ async def main() -> None:
child_run = await apify_client_async.run(child_run_id).wait_for_finish(wait_duration=timedelta(seconds=120))
assert child_run is not None
assert child_run.status == 'ABORTED'


async def test_named_child_runs_respect_the_concurrency_limit(
make_actor: MakeActorFunction,
run_actor: RunActorFunction,
) -> None:
"""Named child runs started concurrently under a limit of one run one after another."""

async def main() -> None:
async with Actor:
actor_input = (await Actor.get_input()) or {}
if actor_input.get('is_child') is True:
await asyncio.sleep(10)
return

actor_id = Actor.configuration.actor_id or ''
Actor.set_child_run_limits(max_concurrent_runs=1)
runs = await asyncio.gather(
*(
Actor.call(actor_id=actor_id, run_input={'is_child': True}, name=f'child-{index}')
for index in range(2)
)
)
first, second = sorted(runs, key=lambda run: run.started_at)
assert first.finished_at is not None, 'first.finished_at is None'
assert second.started_at >= first.finished_at, f'first={first}, second={second}'

actor = await make_actor(label='child-run-limit', main_func=main)
run_result = await run_actor(actor)

assert run_result.status == 'SUCCEEDED'
# The parent run and its two child runs.
assert (await actor.runs().list()).total == 3
Loading
Loading