Repository navigation
Keep a single EventManager per process (event loop) #2242
Description
Activity
- addedt-toolingIssues with this label are in the ownership of the tooling team.Issues with this label are in the ownership of the tooling team.
on Sep 19, 2026 Hey @Mantisus, thanks for filing this. We'll probably want to do this right in crawlee-js 4.0. So are you saying that we should simply not allow passing the
event_manageroption to theBasicCrawlerconstructor?So are you saying that we should simply not allow passing the event_manager option to the BasicCrawler constructor?
Yes, and the same goes for the other places that accept it today,
SessionPool, for example. The other half is decoupling theEventManagerlifecycle from the crawler. Today, it only runs while its context is entered, and in practice only the crawler does that.So that users don't have to start the event manager themselves, I'd propose lazy initialization: the manager starts on first use and keeps running until the asyncio event loop shuts down or the user closes it explicitly.
Lazy initialization and async don't play well together, do you have something concrete in mind?
Starting it isn't a problem, since
EventManager.__aenter__never actually awaits anything. It only starts the recurring tasks, which a synconoremitcan do just as well.The part that needs thought is closing it. I'm considering two options.
- A task created with
loop.create_taskthat waits until it's cancelled. On exit,asyncio.runcancels every pending task and then awaits them, so the task'sfinallyblock still runs with a working loop.
class SingleEventManager(EventManager): def __init__(self, **kwargs: Any) -> None: super().__init__(**kwargs) self._shutdown_task: asyncio.Task[None] | None = None def on(self, *, event: Event, listener: EventListener[Any]) -> None: self._ensure_started() super().on(event=event, listener=listener) def emit(self, *, event: Event, event_data: EventData) -> None: self._ensure_started() super().emit(event=event, event_data=event_data) async def close(self) -> None: if (shutdown_task := self._shutdown_task) is None: return shutdown_task.cancel() await asyncio.wait([shutdown_task]) if not shutdown_task.cancelled() and (exc := shutdown_task.exception()) is not None: raise exc await self._close() def _ensure_started(self) -> None: if self.active: return try: loop = asyncio.get_running_loop() except RuntimeError: raise RuntimeError( 'The event manager can only be used from within a running asyncio event loop.' ) from None self._emit_persist_state_event_rec_task.start() self._shutdown_task = loop.create_task(self._wait_for_shutdown(), name='single-event-manager-shutdown') async def _wait_for_shutdown(self) -> None: try: await asyncio.Event().wait() finally: # Stops the recurring task, emits a last `PersistState` and waits for the listeners. await self._close()
- An async generator that's closed by
loop.shutdown_asyncgens(), whichasyncio.runcalls on exit, after all tasks are done.
class SingleEventManager(EventManager): def __init__(self, **kwargs: Any) -> None: super().__init__(**kwargs) self._shutdown_hook: AsyncGenerator[None, None] | None = None def on(self, *, event: Event, listener: EventListener[Any]) -> None: self._ensure_started() super().on(event=event, listener=listener) def emit(self, *, event: Event, event_data: EventData) -> None: self._ensure_started() super().emit(event=event, event_data=event_data) async def close(self) -> None: if self._shutdown_hook is not None: await self._shutdown_hook.aclose() def _ensure_started(self) -> None: if self.active: return try: asyncio.get_running_loop() except RuntimeError: raise RuntimeError( 'The event manager can only be used from within a running asyncio event loop.' ) from None self._emit_persist_state_event_rec_task.start() self._hook_into_loop() def _hook_into_loop(self) -> None: hook = self._wait_for_shutdown() # Start the async generator from a sync method, so that the loop registers it. send = hook.asend(None) try: send.send(None) except StopIteration: self._shutdown_hook = hook return send.close() raise RuntimeError('The shutdown hook must not await anything before its `yield`.') async def _wait_for_shutdown(self) -> AsyncGenerator[None, None]: try: yield finally: # Stops the recurring task, emits a last `PersistState` and waits for the listeners. await self._close()
The main difference is ordering. With the task, the final flush runs concurrently with the teardown of the other tasks. The generator runs after all of them have finished.
- A task created with
Tying the event manager lifecycle to the lifetime of the event loop is a viable option, sure.
Then again, in crawlee-js, we went with the owned-or-injected pattern for most crawler dependencies (either crawler constructs and tears down a default implementation, or it receives a readymade instance via constructor and the caller is responsible for the lifecycle management). This is nice for services that the user may actually want to override, because it makes the lifecycle clear.
I guess that pretty much nobody wants to override the event manager though, so we shouldn't force people to own the lifecycle management in this case.
Regarding the async startup, this is true for the current implementations, but what if some alternative event manager actually needs async initialization?
The context manager stays for the cases that need explicit lifecycle management. In the SDK the Actor owns the
ApifyEventManager: it initializes it and closes it when the Actor exits. I believe, that's correct behaviour and it shouldn't change.The rule is then the same in both directions: whoever enters the context also closes it, and a manager nobody owns closes itself with the event loop.
So a custom implementation that needs to
awaitduring initialization uses the context manager. And users who don't want to think about the event manager at all get the lazy initialization, with the manager working whether or not a crawler is running at that moment.For implementations that need the async initialization, we can add a class attribute like
supports_lazy_start, so that the lazy path raises instead of leaving the manager half-initialized.@janbuchar I guess we forgot about this one, if you think we should make this in the JS part too, let's create an issue on that side
Created the JS counterpart: apify/crawlee#4224
Same two problems there. Passing any custom service to a crawler (even just
configuration) can create a second event manager, and standalone storages never get periodicpersistState. Node has no event loop shutdown hook, so the close side would go through anunref()'d interval plus a final flush inprocess.once('beforeExit').One thing to watch here: the lazy start should hold its own ref in
_active_ref_count. Otherwise, when the Actor enters and exits a manager that already started lazily, the last exit clears the listeners that standalone components registered.Reacted by Max Bohomolov
Make
EventManagera single service per process (per event loop, in practice), owned by the globalservice_locator.Motivation
Today it's possible to create several
EventManagerinstances that emit the same events.BasicCrawlerbuilds its ownServiceLocatorand may put a crawler-specific event manager into it, whileSnapshotterandRecoverableStatealways resolve the globalservice_locator. So a crawler-scoped manager runs a secondSystemInfosampling loop and emitsPersistStatethat nobody listens to, while the components it was supposed to serve keep listening on the global one. Working around that is whyBasicCrawler._run_crawlercurrently enters both managers.The second problem is lifetime. An event manager only runs while at least one crawler is active, but several systems that depend on it are usable without a crawler.
FileSystemRequestQueueClientkeeps its ordering and in-progress state inRecoverableState, and the same holds for autosavedKeyValueStorevalues,RequestListandSitemapRequestLoader. Used standalone, they register aPersistStatelistener on a manager that is never started, so their state is only written on explicit teardown.Implementation
Keep exactly one
EventManager, reachable through a single access point inservice_locator, resolved per running event loop so that a manager never leaks from oneasyncio.runinto the next.service_locator.set_event_managerstays as the way to install a custom manager, for tests or for the Apify SDK's platform event manager.The manager starts itself on the first
onoremitcall and closes itself when its event loop shuts down, so it works regardless of whether a crawler is running. There is no context to enter and noasync within the calling code. All currentoncall sites are already inside async code, and so is everyRecurringTask.startcall site, so the lazy start always has a running loop to attach to.