summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-06 20:27:13 +0100
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-06 20:27:13 +0100
commit32389f971fd299a9284d40a608bea0bd3204aff4 (patch)
tree026fc0917d029c851379680bce729409c9f09f29
parent5aee8f86cab36f0a802f7072f937e13fa5a4fb4d (diff)
downloadjoblib-32389f971fd299a9284d40a608bea0bd3204aff4.tar.gz
joblib-32389f971fd299a9284d40a608bea0bd3204aff4.zip
fixed multiple teardowns
-rw-r--r--worker.py26
1 files changed, 13 insertions, 13 deletions
diff --git a/worker.py b/worker.py
index 88ca110..27fd7d8 100644
--- a/worker.py
+++ b/worker.py
@@ -49,20 +49,20 @@ class Worker(object):
continue
error = None
- for jobid, payload, priority in jobs:
- try:
+ try:
+ for jobid, payload, priority in jobs:
await self._execute_job(jobid, payload, priority)
- except JobFailed as err:
- error = err.args[0]
- finally:
- if error is not None:
- await self._queue.reset([j[0] for j in jobs])
- await self._queue.fail(jobid, error)
- logging.warning("batch failed, reset")
- await self._teardown(failed=True)
- else:
- await self._queue.finish([j[0] for j in jobs])
- await self._teardown(failed=False)
+ except JobFailed as err:
+ error = err.args[0]
+ finally:
+ if error is not None:
+ await self._queue.reset([j[0] for j in jobs])
+ await self._queue.fail(jobid, error)
+ logging.warning("batch failed, reset")
+ await self._teardown(failed=True)
+ else:
+ await self._queue.finish([j[0] for j in jobs])
+ await self._teardown(failed=False)
async def start(self, batchlimit=1):
"""Start in background."""