diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-10 22:14:00 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-10 22:14:00 +0100 |
| commit | 542df395ce186ccc28d2c4122b7b12024ade8f47 (patch) | |
| tree | 9671c225bd3f2970e7a78f56a14b7662f2056dac | |
| parent | e4896dc8e460ec5c1a32348e2137e19ac2c90d3a (diff) | |
| download | joblib-1.0.0.tar.gz joblib-1.0.0.zip | |
bugfix: do not commit failed jobs as finished1.0.0
| -rw-r--r-- | worker.py | 6 |
1 files changed, 4 insertions, 2 deletions
@@ -50,9 +50,11 @@ class Worker(object): await self._windup() critical_error = False failed = [] + succeeded = [] for jobid, payload, priority in jobs: try: await self._execute_job(jobid, payload, priority) + succeeded.append(jobid) except JobFailed as err: failed.append((jobid, err.args[0])) if len(err.args) > 1: @@ -63,10 +65,10 @@ class Worker(object): else: critical_error = True # default if critical_error: - await self._queue.reset([j[0] for j in jobs]) + await self._queue.reset(succeeded) logging.warning("batch failed, reset") else: - await self._queue.finish([j[0] for j in jobs]) + await self._queue.finish(succeeded) for jobid, err in failed: await self._queue.fail(jobid, err) |
