diff options
Diffstat (limited to 'api.py')
| -rw-r--r-- | api.py | 22 |
1 files changed, 14 insertions, 8 deletions
@@ -52,7 +52,7 @@ class Worker(object): with open(root + "/insert.sql", "r", encoding="utf-8-sig") as file: self._insertquery = file.read() - async def _execute_job(self, jobid, payload): + async def _execute_job(self, jobid, payload, priority): """Finish a job.""" api = crawler.Crawler(self._apitoken) logging.debug("%s: getting matches from API", jobid) @@ -60,20 +60,26 @@ class Worker(object): async for data in api.matches(region=payload["region"], params=payload["params"]): logging.debug("%s: inserting into database", jobid) - ids = await con.fetch(self._insertquery, json.dumps(data)) - logging.info("%s: inserted %s objects", jobid, len(ids)) - object_ids = [i["id"] for i in ids] - for object_id in object_ids: + objects = await con.fetch(self._insertquery, json.dumps(data)) + logging.info("%s: inserted %s", jobid, + {t: len([s for s in objects if s["type"] == t]) + for t in set([e["type"] for e in objects])} + ) + for obj in objects: await self._queue.request(jobtype="process", - payload={"id": object_id}) + priority=priority, + payload={ + "id": obj["id"], + "type": obj["type"] + }) async def _work(self): """Fetch a job and run it.""" - jobid, payload = await self._queue.acquire(jobtype="grab") + jobid, payload, priority = await self._queue.acquire(jobtype="grab") if jobid is None: raise LookupError("no jobs available") logging.debug("%s: starting job", jobid) - await self._execute_job(jobid, payload) + await self._execute_job(jobid, payload, priority) await self._queue.finish(jobid) logging.debug("%s: finished job", jobid) |
