From 57ac9c2b03f9e46de6ea4a16eb3954868d1427c9 Mon Sep 17 00:00:00 2001 From: schneefux Date: Mon, 27 Mar 2017 20:56:53 +0200 Subject: (wip) use 2.0 joblib --- Dockerfile | 6 ++ api.py | 242 ------------------------------------------------------- requirements.txt | 7 +- worker.py | 86 ++++++++++++++++++++ 4 files changed, 96 insertions(+), 245 deletions(-) create mode 100644 Dockerfile delete mode 100644 api.py create mode 100644 worker.py diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..6344e0f --- /dev/null +++ b/Dockerfile @@ -0,0 +1,6 @@ +FROM python:3.6-alpine +RUN apk add --no-cache postgresql-dev gcc python3-dev musl-dev +ADD requirements.txt /code/requirements.txt +WORKDIR /code +RUN pip install -r requirements.txt +CMD ["python", "worker.py"] diff --git a/api.py b/api.py deleted file mode 100644 index ff68327..0000000 --- a/api.py +++ /dev/null @@ -1,242 +0,0 @@ -#!/usr/bin/python3 - -import asyncio -import os -import datetime -import glob -import logging -import asyncpg - -import joblib.worker - - -source_db = { - "host": os.environ.get("POSTGRESQL_SOURCE_HOST") or "vaindock_postgres_raw", - "port": os.environ.get("POSTGRESQL_SOURCE_PORT") or 5532, - "user": os.environ.get("POSTGRESQL_SOURCE_USER") or "vainraw", - "password": os.environ.get("POSTGRESQL_SOURCE_PASSWORD") or "vainraw", - "database": os.environ.get("POSTGRESQL_SOURCE_DB") or "vainsocial-raw" -} - -dest_db = { - "host": os.environ.get("POSTGRESQL_DEST_HOST") or "vaindock_postgres_web", - "port": os.environ.get("POSTGRESQL_DEST_PORT") or 5432, - "user": os.environ.get("POSTGRESQL_DEST_USER") or "vainweb", - "password": os.environ.get("POSTGRESQL_DEST_PASSWORD") or "vainweb", - "database": os.environ.get("POSTGRESQL_DEST_DB") or "vainsocial-web" -} - - -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_preload=False, do_analyze=False): - self._queries = {} - super().__init__(jobtype="process") - self._do_preload = do_preload # request preload jobs - self._do_analyze = do_analyze # request machine learning - logging.debug("preload: %s, analyze: %s", - do_preload, do_analyze) - - async def connect(self, sourcea, desta): - """Connect to database.""" - logging.warning("connecting to database") - await super().connect(**sourcea) - self._srccon = await asyncpg.connect(**sourcea) - self._destcon = await asyncpg.connect(**desta) - - async def setup(self): - """Initialize the database.""" - scriptroot = os.path.realpath( - os.path.join(os.getcwd(), os.path.dirname(__file__))) - for path in glob.glob(scriptroot + "/queries/*.sql"): - # utf-8-sig is used by pgadmin, doesn't hurt to specify - # 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] = 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() - self._desttr = self._destcon.transaction() - await self._srctr.start() - await self._desttr.start() - self._priorities = [] - self._compilejobs = [] - self._analyzejobs = [] - self._preloads = [] - - async def _teardown(self, failed): - if failed: - await self._srctr.rollback() - await self._desttr.rollback() - else: - await self._srctr.commit() - await self._desttr.commit() - - await self.request( - jobtype="compile", - payload=self._compilejobs, - priority=self._priorities) - if self._do_analyze: - await self.request( - jobtype="analyze", - payload=self._analyzejobs) - - if self._do_preload: - preloadjobs = [{ - "region": s[0], - "params": { - "filter[playerNames]": s[1], - "filter[createdAt-start]": date2iso(s[2] + datetime.timedelta(seconds=1)), - "filter[gameMode]": "casual,ranked" - } - } for s in self._preloads] - await self.request( - jobtype="preload", - payload=preloadjobs, - priority=[2]*len(preloadjobs)) - - async def _execute_job(self, jobid, payload, priority): - """Finish a job.""" - object_id = payload["id"] - explicit_player = payload.get("playername") or "" - # 1 object in raw : n objects in web - for table, query in self._queries.items(): - logging.debug("%s: running '%s' query", - jobid, table) - # fetch from raw, converted to format for web table - 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) - else: - obj = await self._into( - self._destcon, data, table) - obj_id = None - if obj is not None: - obj_id = obj["api_id"] - except asyncpg.exceptions.DeadlockDetectedError: - logging.error("%s: deadlocked!", jobid) - raise joblib.worker.JobFailed("deadlock", - True) # critical - except asyncpg.exceptions.IntegrityConstraintViolationError as err: - logging.error("%s: SQL error '%s'!", jobid, err) - raise joblib.worker.JobFailed({"id": object_id, - "error": str(err)}, True) - - logging.debug("record processed") - if obj_id is not None: - # run web->web queries - payload = { - "type": table, - "id": obj_id - } - logging.debug("%s: requesting jobs for %s", - jobid, obj_id) - self._priorities.append(priority) - if table == "participant": - self._compilejobs.append(payload) - if table == "player" and data["name"] == explicit_player: - self._compilejobs.append(payload) - if table == "participant": - self._analyzejobs.append(payload) - - if lmcd is not None: - now = datetime.datetime.now() - interval_mins = (now - lmcd).total_seconds()/60 - logging.debug("%s: minutes since last match: %s", - jobid, interval_mins) - if priority == 1 and interval_mins > 30: - # TODO move to config var ^ - # TODO - # only allowed 1 level by ToS - if data["name"] not in [p[1] for p in self._preloads]: - # prevent duplicates - # TODO also prevent across batches! - self._preloads.append((data["shard_id"], - data["name"], lmcd)) - - await self._deletematch.fetchrow(object_id) - - async def _playerinto(self, conn, data, table, update_date): - """Upsert a player named tuple into a table. - Return the object id.""" - # save lmcd to restore later - lmcd = await conn.fetchval(""" - SELECT last_match_created_date FROM player - WHERE api_id=$1 - """, data["api_id"]) - - obj = await self._into(conn, data, table, conflict=""" - DO UPDATE SET ("{1}") = ({2}) - WHERE COALESCE(player.last_match_created_date, - 'epoch'::TIMESTAMP) - <= COALESCE(EXCLUDED.last_match_created_date, - 'epoch'::TIMESTAMP) - 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"] - - # restore lmcd because - # we want to request a preload 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, objlmcd - - async def _into(self, conn, data, table, - conflict="DO UPDATE SET (\"{1}\") = ({2}) " + - "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 {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(): - worker = Processor( - do_preload=os.environ.get("VAINSOCIAL_SPIDER")=="true", - do_analyze=os.environ.get("VAINSOCIAL_ANALYZE")=="true" - ) - await worker.connect( - source_db, dest_db - ) - await worker.setup() - await worker.run(batchlimit=1000) - -logging.basicConfig(level=logging.DEBUG) - -loop = asyncio.get_event_loop() -loop.run_until_complete(startup()) diff --git a/requirements.txt b/requirements.txt index aad7a96..b549b13 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,5 +1,6 @@ -appdirs==1.4.0 -asyncpg==0.8.4 +appdirs==1.4.3 packaging==16.8 -pyparsing==2.1.10 +pika==0.10.0 +psycopg2==2.7.1 +pyparsing==2.2.0 six==1.10.0 diff --git a/worker.py b/worker.py new file mode 100644 index 0000000..9a00f06 --- /dev/null +++ b/worker.py @@ -0,0 +1,86 @@ +#!/usr/bin/python3 + +import os +import logging +import psycopg2 +import psycopg2.extras +import psycopg2.extensions + +import joblib.joblib + + +RABBIT = { + "host": os.environ.get("RABBITMQ_HOST"), + "port": os.environ.get("RABBITMQ_PORT"), + "credentials": os.environ.get("RABBITMQ_CREDS") +} + +DB = { + "host": os.environ.get("POSTGRESQL_HOST") or "vaindock_postgres_web", + "port": os.environ.get("POSTGRESQL_PORT") or 5432, + "user": os.environ.get("POSTGRESQL_USER") or "vainweb", + "password": os.environ.get("POSTGRESQL_PASSWORD") or "vainweb", + "dbname": os.environ.get("POSTGRESQL_DB") or "vainsocial-web" +} + + +class Processor(joblib.joblib.Worker): + def __init__(self): + super().__init__(jobtype="process") + + def connect(self, rabbit, db): + """Connect to database.""" + logging.warning("connecting to database") + super().connect(**rabbit) + self._con = psycopg2.connect(**db) + + def commit(self, failed): + self._con.commit() + + def work(self, payload): + """Finish a job.""" + obj_id = payload["id"] + obj_type = payload["type"] + o = payload["data"]["attributes"] + r = payload["data"].get("relationships") + + c = self._con.cursor( + cursor_factory=psycopg2.extras.DictCursor) + + d = {} + if obj_type == "match": + d["api_id"] = obj_id + d["created_at"] = o["createdAt"] + d["duration"] = o["duration"] + d["game_mode"] = o["gameMode"] + d["patch_version"] = o.get("patchVersion") or "2.2" + d["shard_id"] = o["shardId"] + d["end_game_reason"] = o["stats"]["endGameReason"] + d["queue"] = o["stats"]["queue"] + d["roster_1"] = r["rosters"]["data"][0]["id"] + d["roster_2"] = r["rosters"]["data"][1]["id"] + + if d == {}: + raise Exception + # TODO + + # TODO player upsert + # TODO request compile & preload + l = [(c, v) for c, v in d.items()] + columns = ",".join([t[0] for t in l]) + values = tuple([t[1] for t in l]) + c.execute("INSERT INTO " + obj_type + "(%s) VALUES %s", + ([psycopg2.extensions.AsIs(columns)] + [values])) + c.close() + + +def startup(): + worker = Processor() + worker.connect(RABBIT, DB) + worker.setup() + worker.run() + +logging.basicConfig(level=logging.DEBUG) + +if __name__ == "__main__": + startup() -- cgit v1.3.1