diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-06 20:27:13 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-06 20:27:13 +0100 |
| commit | 32389f971fd299a9284d40a608bea0bd3204aff4 (patch) | |
| tree | 026fc0917d029c851379680bce729409c9f09f29 | |
| parent | 5aee8f86cab36f0a802f7072f937e13fa5a4fb4d (diff) | |
| download | joblib-32389f971fd299a9284d40a608bea0bd3204aff4.tar.gz joblib-32389f971fd299a9284d40a608bea0bd3204aff4.zip | |
fixed multiple teardowns
| -rw-r--r-- | worker.py | 26 |
1 files changed, 13 insertions, 13 deletions
@@ -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.""" |
