diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-09 18:45:28 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-09 18:45:28 +0100 |
| commit | 535a14f4ea1973eaabcaf4624604a42335b15981 (patch) | |
| tree | 82bcf619fafeafc8cb87cb824e9973bb4b0233e2 | |
| parent | 7fa00a9e4f2723f7d93fc1ecc8369ccbbc9e4f76 (diff) | |
| download | joblib-535a14f4ea1973eaabcaf4624604a42335b15981.tar.gz joblib-535a14f4ea1973eaabcaf4624604a42335b15981.zip | |
allow non critical job fails
| -rw-r--r-- | worker.py | 35 |
1 files changed, 21 insertions, 14 deletions
@@ -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.""" |
