diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-27 20:56:53 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-27 20:56:53 +0200 |
| commit | 57ac9c2b03f9e46de6ea4a16eb3954868d1427c9 (patch) | |
| tree | 4e844799b542b31b62ce7a8ffcf779234ad46d91 /worker.py | |
| parent | c393152d1f84fc19a0d04c91d7d545c89a7d3693 (diff) | |
| download | shrinker-57ac9c2b03f9e46de6ea4a16eb3954868d1427c9.tar.gz shrinker-57ac9c2b03f9e46de6ea4a16eb3954868d1427c9.zip | |
(wip) use 2.0 joblib
Diffstat (limited to 'worker.py')
| -rw-r--r-- | worker.py | 86 |
1 files changed, 86 insertions, 0 deletions
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() |
