diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-04 11:12:41 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-04 11:12:41 +0100 |
| commit | b8ea5a875043fa3ab13c81187404c6a389213184 (patch) | |
| tree | 8422cff4f64ee39b6a6aca4fa2af86a8d148228f | |
| parent | 206b66d19045f7157e9e54d84318cbffbd18d7d6 (diff) | |
| download | apigrabber-b8ea5a875043fa3ab13c81187404c6a389213184.tar.gz apigrabber-b8ea5a875043fa3ab13c81187404c6a389213184.zip | |
refactor to use joblib worker class
| -rw-r--r-- | api.py | 62 | ||||
| m--------- | joblib | 0 |
2 files changed, 19 insertions, 43 deletions
@@ -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 -Subproject 31d0588a35c8df17a5d3adcfc25af2b72c896cc +Subproject 8d58a4f5f754b6541af84478a2df3fba70db94b |
