From 3229571141b638d9d5d6985e98b7e61642229b77 Mon Sep 17 00:00:00 2001 From: schneefux Date: Sun, 5 Mar 2017 10:36:14 +0100 Subject: support job batching (#12) * support job batching * worker: pull more frequently --- worker.py | 60 ++++++++++++++++++++++++++++++++++++++++++++---------------- 1 file changed, 44 insertions(+), 16 deletions(-) (limited to 'worker.py') diff --git a/worker.py b/worker.py index 2737671..163ebfd 100644 --- a/worker.py +++ b/worker.py @@ -24,33 +24,61 @@ class Worker(object): # override pass + async def _windup(self): + # override + pass + async def _execute_job(self, jobid, payload, priority): # override pass + async def _teardown(self, failed): + # override + pass + async def _work(self): - """Fetch a job and run it.""" + """Fetch a job and run it. + Return id.""" jobid, payload, priority = await self._queue.acquire( jobtype=self._jobtype) if jobid is None: - raise LookupError("no jobs available") + raise LookupError try: await self._execute_job(jobid, payload, priority) - await self._queue.finish(jobid) - except JobFailed as error: - logging.warning("%s: failed with %s", jobid, - error.args[0]) - await self._queue.fail(jobid, error.args[0]) + except JobFailed as err: + raise JobFailed(err.args[0], jobid) + return jobid - async def run(self): + async def run(self, batchlimit=1): """Start jobs forever.""" while True: + await self._windup() + jobids = [] + error = None + low_load = False try: - await self._work() - except LookupError: - await asyncio.sleep(1) - - async def start(self, number=1): - """Start jobs in background.""" - for _ in range(number): - asyncio.ensure_future(self.run()) + for _ in range(batchlimit): + try: + jobids.append(await self._work()) + except LookupError: + low_load = True + break + except JobFailed as err: + error = err.args[0] + jobids.append(err.args[1]) + break + finally: + if error is not None: + await self._queue.reset(jobids[:-1]) + logging.debug(jobids) + await self._queue.fail(jobids[-1], error) + logging.warning("batch failed, reset") + else: + await self._queue.finish(jobids) + await self._teardown(failed=error is not None) + if low_load: + await asyncio.sleep(0.1) + + async def start(self, batchlimit=1): + """Start in background.""" + asyncio.ensure_future(self.run(batchlimit)) -- cgit v1.3.1