From b8ea5a875043fa3ab13c81187404c6a389213184 Mon Sep 17 00:00:00 2001 From: schneefux Date: Sat, 4 Mar 2017 11:12:41 +0100 Subject: refactor to use joblib worker class --- api.py | 62 +++++++++++++++++++------------------------------------------- joblib | 2 +- 2 files changed, 20 insertions(+), 44 deletions(-) diff --git a/api.py b/api.py index 43343db..f66b880 100644 --- a/api.py +++ b/api.py @@ -7,22 +7,20 @@ import json import asyncpg import crawler -import joblib.joblib +import joblib.worker -class Worker(object): +class Apigrabber(joblib.worker.Worker): def __init__(self, apitoken): self._apitoken = apitoken - self._queue = None self._pool = None self._insertquery = "" + super().__init__(jobtype="grab") async def connect(self, **args): """Connect to database.""" logging.warning("connecting to database") - self._queue = joblib.joblib.JobQueue() - await self._queue.connect(**args) - await self._queue.setup() + await super().connect(**args) self._pool = await asyncpg.create_pool(**args) async def setup(self): @@ -49,47 +47,25 @@ class Worker(object): playername = "" async with self._pool.acquire() as con: - async for data in api.matches(region=payload["region"], - params=payload["params"]): - matchids = await con.fetch(self._insertquery, json.dumps(data)) - logging.debug("%s: inserted %s matches from API into database", - jobid, len(matchids)) - for matchid in matchids: - await self._queue.request(jobtype="process", - priority=priority, - payload={ - "id": matchid["id"], - "playername": playername - }) - - async def _work(self): - """Fetch a job and run it.""" - jobid, payload, priority = await self._queue.acquire(jobtype="grab") - if jobid is None: - raise LookupError("no jobs available") - try: - await self._execute_job(jobid, payload, priority) - await self._queue.finish(jobid) - except crawler.ApiError as error: - logging.warning("%s: failed with %s", jobid, - error.args[0]) - await self._queue.fail(jobid, error.args[0]) - - async def run(self): - """Start jobs forever.""" - while True: try: - await self._work() - except LookupError: - await asyncio.sleep(1) + async for data in api.matches(region=payload["region"], + params=payload["params"]): + matchids = await con.fetch(self._insertquery, json.dumps(data)) + logging.debug("%s: inserted %s matches from API into database", + jobid, len(matchids)) + for matchid in matchids: + await self._queue.request(jobtype="process", + priority=priority, + payload={ + "id": matchid["id"], + "playername": playername + }) + except crawler.ApiError as error: + raise joblib.worker.JobFailed(error.args[0]) - async def start(self, number=1): - """Start jobs in background.""" - for _ in range(number): - asyncio.ensure_future(self.run()) async def startup(): - worker = Worker( + worker = Apigrabber( apitoken=os.environ["VAINSOCIAL_APITOKEN"] ) await worker.connect( diff --git a/joblib b/joblib index 31d0588..8d58a4f 160000 --- a/joblib +++ b/joblib @@ -1 +1 @@ -Subproject commit 31d0588a35c8df17a5d3adcfc25af2b72c896cc7 +Subproject commit 8d58a4f5f754b6541af84478a2df3fba70db94be -- cgit v1.3.1