summaryrefslogtreecommitdiff
path: root/api.py
diff options
context:
space:
mode:
Diffstat (limited to 'api.py')
-rw-r--r--api.py173
1 files changed, 80 insertions, 93 deletions
diff --git a/api.py b/api.py
index b37d669..b316a6b 100644
--- a/api.py
+++ b/api.py
@@ -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():