diff options
Diffstat (limited to 'api.py')
| -rw-r--r-- | api.py | 173 |
1 files changed, 80 insertions, 93 deletions
@@ -4,7 +4,6 @@ import asyncio import os import glob import logging -import json import asyncpg import joblib.worker @@ -27,14 +26,22 @@ dest_db = { } +def date2iso(d): + """Convert datetime to iso8601 zulu string.""" + date = d.replace(microsecond=0) + date = date.isoformat() + date += "Z" + return date + + class Processor(joblib.worker.Worker): def __init__(self, do_spider=False, do_analyze=False): self._srcpool = None self._destpool = None self._queries = {} super().__init__(jobtype="process") - self._do_spider = do_spider - self._do_analyze = do_analyze + self._do_spider = do_spider # request spider jobs + self._do_analyze = do_analyze # request machine learning async def connect(self, sourcea, desta): """Connect to database.""" @@ -52,8 +59,11 @@ class Processor(joblib.worker.Worker): # file names: raw target table table = os.path.splitext(os.path.basename(path))[0] with open(path, "r", encoding="utf-8-sig") as file: - self._queries[table] = file.read() + self._queries[table] = await self._srccon.prepare( + file.read()) logging.info("loaded query '%s'", table) + self._deletematch = await self._srccon.prepare( + "DELETE FROM match WHERE id=$1") async def _windup(self): self._srctr = self._srccon.transaction() @@ -71,32 +81,24 @@ class Processor(joblib.worker.Worker): else: await self._srctr.commit() await self._desttr.commit() - self._compilejobs = list({ - # uniquify - j["id"]: j for j in self._compilejobs - }.values()) await self._queue.request( jobtype="compile", payload=self._compilejobs)#, # priority=priority) TODO - self._analyzejobs = list({ - # uniquify - j["id"]: j for j in self._analyzejobs - }.values()) if self._do_analyze: await self._queue.request( jobtype="analyze", payload=self._analyzejobs) - spiderjobs = [{ - "region": s[0], - "params": { - "filter[playerNames]": s[1], - "filter[createdAt-start]": "2017-03-01T00:00:00Z" - } - } for s in self._spiders] if self._do_spider: + spiderjobs = [{ + "region": s[0], + "params": { + "filter[playerNames]": s[1], + "filter[createdAt-start]": date2iso(s[2]) + } + } for s in self._spiders] await self._queue.request( jobtype="grab", payload=spiderjobs, @@ -111,98 +113,83 @@ class Processor(joblib.worker.Worker): logging.debug("%s: running '%s' query", jobid, table) # fetch from raw, converted to format for web table - # TODO refactor - duplicated messy code - if table == "player": - # upsert under special conditions - datas = await self._srccon.fetch( - query, object_id, explicit_player) - for data in datas: - try: - obj_id = await self._playerinto( + datas = await query.fetch(object_id) + for data in datas: + lmcd = None + try: + if table == "player": + obj_id, lmcd = await self._playerinto( self._destcon, data, table, data["name"] == explicit_player) - except asyncpg.exceptions.DeadlockDetectedError: - raise joblib.worker.JobFailed("deadlock") - - logging.debug("record processed") - if obj_id: - # run web->web queries - payload = { - "type": table, - "id": obj_id - } - self._compilejobs.append(payload) - self._analyzejobs.append(payload) - if explicit_player != data["name"]: - self._spiders.append((data["shard_id"], - data["name"])) - else: - datas = await self._srccon.fetch( - query, object_id) - for data in datas: - # insert processed result into web table - try: - obj_id = await self._into( + else: + obj = await self._into( self._destcon, data, table) - except asyncpg.exceptions.DeadlockDetectedError: - raise joblib.worker.JobFailed("deadlock") - logging.debug("record processed") + obj_id = None + if obj: + obj_id = obj["api_id"] + except asyncpg.exceptions.DeadlockDetectedError: + raise joblib.worker.JobFailed("deadlock") + + logging.debug("record processed") + if obj_id: + # run web->web queries + payload = { + "type": table, + "id": obj_id + } + self._compilejobs.append(payload) + self._analyzejobs.append(payload) - if obj_id: - # run web->web queries - payload = { - "type": table, - "id": obj_id - } - self._compilejobs.append(payload) - self._analyzejobs.append(payload) + if lmcd is not None: + self._spiders.append((data["shard_id"], + data["name"], lmcd)) - await self._srccon.execute( - "DELETE FROM match WHERE id=$1", object_id) + await self._deletematch.fetchrow(object_id) - async def _playerinto(self, conn, data, table, do_upsert_date): + async def _playerinto(self, conn, data, table, update_date): """Upsert a player named tuple into a table. Return the object id.""" - if do_upsert_date: - # explicit update for this player -> store - lmcd = data["last_match_created_date"] - data = dict(data) - del data["last_match_created_date"] + # save lmcd to restore later + lmcd = data["last_match_created_date"] - items = list(data.items()) - keys, values = [x[0] for x in items], [x[1] for x in items] - placeholders = ["${}".format(i) for i, _ in enumerate(values, 1)] - # upsert all values except lmcd if they are more recent - query = """ - INSERT INTO player ("{0}", "last_match_created_date") - VALUES ({1}, 'epoch'::timestamp) - ON CONFLICT("api_id") DO UPDATE SET ("{0}") = ({1}) - WHERE player.played < EXCLUDED.played - RETURNING "api_id" - """.format( - "\", \"".join(keys), ", ".join(placeholders)) - objid = await conn.fetchval(query, *data.values()) + obj = await self._into(conn, data, table, conflict=""" + DO UPDATE SET ("{1}") = ({2}) + WHERE player.last_match_created_date + < EXCLUDED.last_match_created_date + RETURNING api_id, last_match_created_date + """) + if obj is None: + logging.debug("player was not updated") + return None, None + + objid = obj["api_id"] + objlmcd = obj["last_match_created_date"] - if do_upsert_date: - # upsert lmcd because it was an explicit request - await conn.execute(""" - UPDATE player SET "last_match_created_date"=$2 - WHERE player."api_id"=$1 AND - player."last_match_created_date" < $2 + # restore lmcd because + # we want to request a spider job + if not update_date: + await conn.fetchval(""" + UPDATE player SET last_match_created_date=$2 + WHERE player.api_id=$1 """, objid, lmcd) - return objid + return objid, objlmcd - async def _into(self, conn, data, table): + async def _into(self, conn, data, table, + conflict="DO NOTHING RETURNING api_id"): """Insert a named tuple into a table. Return the object id.""" items = list(data.items()) keys, values = [x[0] for x in items], [x[1] for x in items] placeholders = ["${}".format(i) for i, _ in enumerate(values, 1)] - query = "INSERT INTO {} (\"{}\") VALUES ({}) ON CONFLICT DO NOTHING RETURNING api_id".format( - table, "\", \"".join(keys), ", ".join(placeholders)) - logging.debug(query) - return await conn.fetchval(query, (*data)) + query = ("INSERT INTO {0} (\"{1}\") VALUES ({2}) ON CONFLICT(api_id) " + + conflict).format( + table, + "\", \"".join(keys), + ", ".join(placeholders)) + logging.debug("query: %s", query) + logging.debug("data: %s", data) + return await conn.fetchrow(query, *data) async def startup(): |
