-
Notifications
You must be signed in to change notification settings - Fork 2
Prevent worker crashes from being treated as successful completion #97
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -16,6 +16,8 @@ | |
| from multiprocessing.queues import Queue as MPQueue | ||
| from typing import Any, ParamSpec, TypeVar | ||
|
|
||
| from .exceptions import UnexpectedWorkerExitError | ||
|
|
||
| logger = logging.getLogger(__name__) | ||
|
|
||
| T = TypeVar('T') | ||
|
|
@@ -55,9 +57,15 @@ async def _iterator_cpu_bound_inner( | |
| try: | ||
| item = await asyncio.to_thread(state_queue.get, True, 0.5) | ||
| except queue.Empty: | ||
| if not process.is_alive(): | ||
| break | ||
| continue | ||
| if process.is_alive(): | ||
| continue | ||
| # Completion may have arrived between the timeout and the exit check. | ||
| try: | ||
| item = await asyncio.to_thread(state_queue.get_nowait) | ||
| except queue.Empty as e: | ||
| raise UnexpectedWorkerExitError( | ||
| f'{process.name} exited with code {process.exitcode} without sending IteratorDone' | ||
| ) from e | ||
|
Comment on lines
+66
to
+68
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
|
||
| match item: | ||
| case IteratorDone(): | ||
| break | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
A
RuntimeErroris what_perform_statetreats as retryable: it resets the state toTrainModelDownloadedand starts_trainagain. The non-finite loss from the description fails the same way every time, so the node would keep retraining and never reachReadyForCleanup. A worker killed by the OOM killer should be retried though. MaybeCriticalErrorwhen the worker exited on its own, and a retry only when a signal killed it?