summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-10 22:14:00 +0100
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-10 22:14:00 +0100
commit542df395ce186ccc28d2c4122b7b12024ade8f47 (patch)
tree9671c225bd3f2970e7a78f56a14b7662f2056dac
parente4896dc8e460ec5c1a32348e2137e19ac2c90d3a (diff)
downloadjoblib-1.0.0.tar.gz
joblib-1.0.0.zip
bugfix: do not commit failed jobs as finished1.0.0
-rw-r--r--worker.py6
1 files changed, 4 insertions, 2 deletions
diff --git a/worker.py b/worker.py
index 1e16237..e5f0b25 100644
--- a/worker.py
+++ b/worker.py
@@ -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)