summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-09 18:45:28 +0100
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-09 18:45:28 +0100
commit535a14f4ea1973eaabcaf4624604a42335b15981 (patch)
tree82bcf619fafeafc8cb87cb824e9973bb4b0233e2
parent7fa00a9e4f2723f7d93fc1ecc8369ccbbc9e4f76 (diff)
downloadjoblib-535a14f4ea1973eaabcaf4624604a42335b15981.tar.gz
joblib-535a14f4ea1973eaabcaf4624604a42335b15981.zip
allow non critical job fails
-rw-r--r--worker.py35
1 files changed, 21 insertions, 14 deletions
diff --git a/worker.py b/worker.py
index 2115847..2d33537 100644
--- a/worker.py
+++ b/worker.py
@@ -48,21 +48,28 @@ class Worker(object):
continue
await self._windup()
- error = None
- try:
- for jobid, payload, priority in jobs:
+ critical_error = True
+ failed = []
+ for jobid, payload, priority in jobs:
+ try:
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:
+ failed.append((jobid, err.args[0]))
+ if len(err.args) > 1:
+ # rollback
+ critical_error = err.args[1]
+ if critical_error:
+ break
+ if critical_error:
+ await self._queue.reset([j[0] for j in jobs])
+ logging.warning("batch failed, reset")
+ else:
+ await self._queue.finish([j[0] for j in jobs])
+
+ for jobid, err in failed:
+ await self._queue.fail(jobid, err)
+
+ await self._teardown(failed=critical_error)
async def start(self, batchlimit=1):
"""Start in background."""