diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-02-26 12:20:05 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-02-26 12:21:12 +0100 |
| commit | ce1f18227b3b8e3d087e2b920151d63d0061fe33 (patch) | |
| tree | 6f70d7bf416c8e2735cbadd3934d81f85677c09b /api.py | |
| parent | 8db59a181e50dac958f20fd6ca851e2c54d4c4da (diff) | |
| download | apigrabber-ce1f18227b3b8e3d087e2b920151d63d0061fe33.tar.gz apigrabber-ce1f18227b3b8e3d087e2b920151d63d0061fe33.zip | |
pass object type to process job payload
Diffstat (limited to 'api.py')
| -rw-r--r-- | api.py | 17 |
1 files changed, 11 insertions, 6 deletions
@@ -60,13 +60,18 @@ 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", - priority=priority, - payload={"id": object_id}) + priority=priority, + payload={ + "id": obj["id"], + "type": obj["type"] + }) async def _work(self): """Fetch a job and run it.""" |
