diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-04-04 18:15:41 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-04-04 18:15:41 +0200 |
| commit | 25b13cefa1ee6e39a2861caab66c4ed8bd749d67 (patch) | |
| tree | 6c0b89ebc988ffc4dc29d2b02eb42f612761127f | |
| parent | 3d52234edae83971e5df58f65f91c95e0352ed3f (diff) | |
| download | analyzer-25b13cefa1ee6e39a2861caab66c4ed8bd749d67.tar.gz analyzer-25b13cefa1ee6e39a2861caab66c4ed8bd749d67.zip | |
rewrite
| -rw-r--r-- | .gitmodules | 3 | ||||
| -rw-r--r-- | Dockerfile | 9 | ||||
| -rw-r--r-- | api.py | 295 | ||||
| -rw-r--r-- | cli.py | 66 | ||||
| m--------- | joblib | 0 | ||||
| -rw-r--r-- | requirements.txt | 11 | ||||
| -rw-r--r-- | worker.py | 164 |
7 files changed, 178 insertions, 370 deletions
diff --git a/.gitmodules b/.gitmodules index 47695ea..e69de29 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,3 +0,0 @@ -[submodule "joblib"] - path = joblib - url = https://gitlab.com/vainglorygame/joblib.git diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..d34166f --- /dev/null +++ b/Dockerfile @@ -0,0 +1,9 @@ +FROM tensorflow/tensorflow:latest-py3 + +RUN mkdir -p /usr/src/app +WORKDIR /usr/src/app + +COPY . /usr/src/app +RUN pip install virtualenv && virtualenv venv && pip install -r requirements.txt + +CMD ["python", "worker.py"] @@ -1,295 +0,0 @@ -#!/usr/bin/python - -import os -import itertools -import json -import asyncio -import asyncpg -import logging -import numpy as np -import tensorflow as tf -import psycopg2 - -import joblib.worker - -#tf.logging.set_verbosity(tf.logging.WARNING) - -queue_db = { - "host": os.environ.get("POSTGRESQL_SOURCE_HOST") or "localhost", - "port": os.environ.get("POSTGRESQL_SOURCE_PORT") or 5433, - "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" -} - -db_config = { - "host": os.environ.get("POSTGRESQL_DEST_HOST") or "localhost", - "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" -} - - -# TODO create abstract class, move to own file -class KDAClassifier(object): - """A DNN that classifies loss/win based on KDA.""" - def __init__(self): - self.model = None - self.modeldir = os.path.realpath( - os.path.join(os.getcwd(), os.path.dirname(__file__))) + "/models/kda-win" - self._pool = None - - # DNN configuration - self._categories = [] - self._continuous = ["kills", "deaths", "assists"] - self._classes = ["loss", "win"] - self._label = "win" - - # learn configuration - self._steps = 1000 # per batch - self._batchsize = 500 - self._max_batches = 20 # number of batches to train - - # database mappings - # TODO cache, rm duplicated code - self._trainquery = """ - SELECT - kills/duration::float, - deaths/duration::float, - assists/duration::float, - participant.winner, - participant.api_id - FROM participant - TABLESAMPLE BERNOULLI(5) - JOIN roster on participant.roster_api_id=roster.api_id - JOIN match on roster.match_api_id=match.api_id - LIMIT %s - """ - self._testquery = """ - SELECT - kills/duration::float, - deaths/duration::float, - assists/duration::float, - participant.winner, - participant.api_id - FROM participant - JOIN roster on participant.roster_api_id=roster.api_id - JOIN match on roster.match_api_id=match.api_id - WHERE participant.api_id=ANY(%s) - """ - - def connect(self, **args): - self._conn = psycopg2.connect(**args) - - def _get_sample(self, size=None, filter=None): - """Return a data set from the database. - :param size: (optional) The number of items to get. Defaults to `All`. - :type size: int or str - :param offset: (optional) The SQL OFFSET parameter. - :type offset: int or str - """ - cur = self._conn.cursor() - assert (filter is None) ^ (size is None) - if filter is None: - cur.execute(self._trainquery, (size,)) - else: - cur.execute(self._testquery, (filter,)) - ids = [] - data = {"kills": [], "deaths": [], "assists": [], "win": []} - for rec in cur: - data["kills"].append(rec[0]) - data["deaths"].append(rec[1]) - data["assists"].append(rec[2]) - data["win"].append(rec[3]) - ids.append(rec[4]) - cur.close() - logging.debug(data) - return data, ids - - def _model_setup(self): - """Sets up a model that takes the `num_features` features as input.""" - feature_columns = [ - tf.contrib.layers.sparse_column_with_hash_bucket(feat, 100) - for feat in self._categories - ] - emb_columns = [ - tf.contrib.layers.embedding_column( - sparse_id_column=col, - dimension=2 # log_2(number of unique features) TODO - ) - for col in feature_columns - ] + [ - tf.contrib.layers.real_valued_column(feat) - for feat in self._continuous - ] - - self.model = tf.contrib.learn.DNNClassifier( - feature_columns=emb_columns, - hidden_units=[2], # http://stats.stackexchange.com/a/1097 TODO inp+outp / 2 - n_classes=len(self._classes), - model_dir=self.modeldir, - config=tf.contrib.learn.RunConfig( - save_checkpoints_secs=2 - ) - ) - - # TODO tf warning - dimensions are wrong - def _toinput(self, sample, train=True): - """Convert sample to Tensors. - :param sample: Data dictionary. - :type sample: dict - :param train: (optional) Whether to return labels too. - :type train: bool - :return: features, label - :rtype: tuple - """ - continuous = {k: tf.constant(sample[k]) - for k in self._continuous} - categories = {k: tf.SparseTensor( - indices=[[i, 0] for i in range(len(sample[k]))], - values=sample[k], - shape=[len(sample[k]), 1]) - for k in self._categories} - - features = {**continuous, **categories} - - if train: - label = tf.constant(sample[self._label]) - return features, label - else: - return features - - # TODO use asyncpg cursor https://magicstack.github.io/asyncpg/current/api/index.html#cursors - def _more(self): - """Fetches up to `limit` number of training items - from the database in batches.""" - # TODO maybe you can use an iterator? - sample = self._get_sample(size=self._batchsize)[0] - return self._toinput(sample, train=True) - - def train(self): - """Train a DNN on a static, algorithmic guess.""" - self._model_setup() - - if os.path.isdir(self.modeldir): - logging.warning("already trained, not training again") - return - - # get one batch of testing data - # TODO either terminate with `eval_steps` or with OutOfRangeError - validation_monitor = tf.contrib.learn.monitors.ValidationMonitor( - input_fn=lambda: self._toinput(self._get_sample(size=2000)[0], train=True), - eval_steps=1, - every_n_steps=10 - ) - - for _ in range(self._max_batches): - self.model.fit( - input_fn=lambda: self._more(), - steps=self._steps, - monitors=[validation_monitor] - ) - - def classify(self, sample, only_best=False): - """Classify a data set. - - :param only_best: (optional) Return the predicted result - instead of a dict of propabilities. - :type only_best: bool - :return: Prediction results. - :rtype: list of dict or list - """ - if only_best: - return self.model.predict(input_fn=lambda: self._toinput(sample, train=False)) - else: - return self.model.predict_proba(input_fn=lambda: self._toinput(sample, train=False)) - - def windup(self): - pass - - def teardown(self, failed=False): - if failed: - pass - else: - self._conn.commit() - - def classify_db(self, objids): - """Classify all data in the data base and insert.""" - # split sample (with participant ids) into data and ids - # TODO use parameters - logging.error("sample size: %s", len(objids)) - sample, objids = self._get_sample(filter=objids) - d = self.classify(sample) - cnt = 0 - cur = self._conn.cursor() - for l in itertools.islice(d, len(objids)): - cur.execute(""" - INSERT INTO participant_stats - (patch_version, participant_api_id, score) - VALUES(2.2, %(objid)s, %(score)s) - ON CONFLICT(participant_api_id) DO - UPDATE SET score=%(score)s - """, {"objid": objids[cnt], "score": float(l[1])}) - cnt += 1 - cur.close() - - -class Analyzer(joblib.worker.Worker): - def __init__(self): - self._pool = None - self._queries = {} - super().__init__(jobtype="analyze") - self.classifier = None - - async def connect(self, dbconf, queuedb): - """Connect to database.""" - logging.warning("connecting to database") - await super().connect(**queuedb) - self._pool = await asyncpg.create_pool(**dbconf) - self.classifier = KDAClassifier() - self.classifier.connect(**dbconf) - - async def setup(self): - """Setup the model.""" - self.classifier.train() - - async def _windup(self): - self._con = await self._pool.acquire() - self._tr = self._con.transaction() - await self._tr.start() - self.classifier.windup() - self._participants = [] - - async def _teardown(self, failed): - if len(self._participants) > 0: - # TODO if this fails, job is still marked as finished - self.classifier.classify_db(self._participants) - - if failed: - await self._tr.rollback() - else: - await self._tr.commit() - await self._pool.release(self._con) - self.classifier.teardown() - - async def _execute_job(self, jobid, payload, priority): - object_id = payload["id"] - object_type = payload["type"] - if object_type != "participant": - return - self._participants.append(object_id) - logging.info("%s: classifying '%s', %s", jobid, - object_type, object_id) - -async def startup(): - worker = Analyzer() - await worker.connect(db_config, queue_db) - await worker.setup() - await worker.run(batchlimit=1000) - - -logging.basicConfig(level=logging.DEBUG) - -loop = asyncio.get_event_loop() -loop.run_until_complete(startup()) @@ -1,66 +0,0 @@ -#!/usr/bin/python3 - -import os -import argparse -import asyncio -import asyncpg - -import joblib.joblib - -queue_db = { - "host": os.environ.get("POSTGRESQL_SOURCE_HOST") or "localhost", - "port": os.environ.get("POSTGRESQL_SOURCE_PORT") or 5433, - "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" -} - -db_config = { - "host": os.environ.get("POSTGRESQL_DEST_HOST") or "localhost", - "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" -} - - -async def main(qdb, sdb, name): - queue = joblib.joblib.JobQueue() - await queue.connect(**qdb) - await queue.setup() - pool = await asyncpg.create_pool(**sdb) - - async with pool.acquire() as con: - async with con.transaction(): - participants = await con.fetch(""" -select -unnest(array[ -roster.participant_1, -roster.participant_2, -roster.participant_3 -]) AS api_id -from roster where roster.match_api_id in ( -select -match.api_id -from player -join participant on participant.player_api_id=player.api_id -join roster on participant.roster_api_id=roster.api_id -join match on roster.match_api_id=match.api_id -where player.name=$1 -) - """, name) - payload = [{ - "id": part["api_id"], - "type": "participant" - } for part in participants] - await queue.request(jobtype="analyze", - payload=payload) - -parser = argparse.ArgumentParser(description="Request a Vainsocial analyze.") -parser.add_argument("-n", "--player", - help="Player name", - type=str) -args = parser.parse_args() - -loop = asyncio.get_event_loop() -loop.run_until_complete(main(queue_db, db_config, args.player)) diff --git a/joblib b/joblib deleted file mode 160000 -Subproject 7ced841aa44a7dd7bd2d3a68ceb989ef488ce44 diff --git a/requirements.txt b/requirements.txt index 77e8224..45e8807 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,9 +1,8 @@ -appdirs==1.4.2 -asyncpg==0.9.0 -numpy==1.12.0 +appdirs==1.4.3 +cymysql==0.8.9 +numpy==1.12.1 packaging==16.8 protobuf==3.2.0 -psycopg2==2.7 -pyparsing==2.1.10 +pyparsing==2.2.0 six==1.10.0 -tensorflow==1.0.1 +SQLAlchemy==1.1.8 diff --git a/worker.py b/worker.py new file mode 100644 index 0000000..976be5f --- /dev/null +++ b/worker.py @@ -0,0 +1,164 @@ +#!/usr/bin/python3 +import os +import random +import itertools +import logging + +from sqlalchemy.ext.automap import automap_base +from sqlalchemy.orm import Session, relationship +from sqlalchemy import create_engine + +import tensorflow as tf +import numpy as np + + +DATABASE_URI = os.environ["DATABASE_URI"] +MODEL_ROOT = os.path.join(os.getcwd(), os.path.dirname(__file__))\ + + "/models/" + +# ORM definitions +Match = Roster = Participant = ParticipantExt = Player = None + + +def connect(): + global Match, Roster, Participant, ParticipantExt, Player + # generate schema from db + Base = automap_base() + engine = create_engine(DATABASE_URI) + Base.prepare(engine, reflect=True) + + # definitions + # TODO check whether the primaryjoin clause is the best method to do this + Match = Base.classes.match + Roster = Base.classes.roster + Roster.match = relationship( + "match", foreign_keys="match.api_id", + primaryjoin="and_(match.api_id == roster.match_api_id)") + Participant = Base.classes.participant + Participant.roster = relationship( + "roster", foreign_keys="roster.api_id", + primaryjoin="and_(roster.api_id == participant.roster_api_id)") + Participant.player = relationship( + "player", foreign_keys="player.api_id", + primaryjoin="and_(player.api_id == participant.player_api_id)") + ParticipantExt = Base.classes.participant_ext + Participant.participant_ext = relationship( + "participant_ext", foreign_keys="participant_ext.participant_api_id", + primaryjoin="and_(participant_ext.participant_api_id == participant.api_id)") + Player = Base.classes.player + + return Session(engine) + + +class Model(object): + # override this configuration + features = [] + label = "" + type = "" + batches = 1 + batchsize = 1 + steps = 0 + id = "unlabeled" + + def __init__(self, db): + self._db = db + for feat in self.features: + self._feature_cols = [tf.contrib.layers.real_valued_column( + feat, dimension=1)] + if self.type == "linear": + self._model = tf.contrib.learn.LinearClassifier( + feature_columns=self._feature_cols, + model_dir=MODEL_ROOT + self.id, + config=tf.contrib.learn.RunConfig( + save_checkpoints_secs=1)) + + def _from_record(self, path, record): + table, column = path.split(".") + if table == "participant": + return vars(record)[column] + if table == "participant_ext": + return vars(record.participant_ext[0])[column] + if table == "player": + return vars(record.player[0])[column] + if table == "roster": + return vars(record.roster[0])[column] + if table == "match": + return vars(record.roster[0].match[0])[column] + raise KeyError("Invalid path " + path) + + def _batch(self, ids=[], size=None): + if len(ids) == 0: # training, get random sample + size = size or self.batchsize + # train a model one batch + offset = random.random() * self._db.query(Participant).count() + # have a bit of randomness in the sample + records = self._db.query( + Participant).offset(offset).limit(size).all() + else: + records = self._db.query( + Participant).filter(Participant.api_id.in_(ids)).all() + + data = {} + labels = [] + # populate from db records + for record in records: + labels.append(self.estimate(record)) + for path in self.features: + if path not in data: + data[path] = [] + data[path].append(self._from_record(path, record)) + + # convert to numpy arrs + for key in data: + data[key] = np.array(data[key]) + labels = np.array(labels) + + return tf.contrib.learn.io.numpy_input_fn( + data, labels, batch_size=self.batchsize, + num_epochs=self.steps) + + # TODO at the moment, it's tied to Participant + def train(self, force=False): + if force or not os.path.isdir(MODEL_ROOT + self.id): + monitor = tf.contrib.learn.monitors.ValidationMonitor( + input_fn=self._batch(), + eval_steps=1, every_n_steps=20) + for _ in range(self.batches): + self._model.fit(input_fn=self._batch(), + steps=self.steps, + monitors=[monitor]) + + def predict(self, ids): + return itertools.islice( + self._model.predict_proba(input_fn=self._batch(ids)), + len(ids)) + + def estimate(self, record): + # override: calculate or return the label's value + pass + + +class MVPScoreModel(Model): + def __init__(self, db): + self.features = ["participant.kills", "participant.deaths", + "participant.assists"] + self.label = "participant_ext.rating" + self.type = "linear" + self.batches = 1 + self.batchsize = 500 + self.steps = 500 + self.id = "kda-win" + super().__init__(db) + + def estimate(self, record): + # for training, rating = participant.winner + return record.winner + + +logging.basicConfig(level=logging.INFO) +if __name__ == "__main__": + db = connect() + model = MVPScoreModel(db) + model.train(force=True) + + print(model._model.evaluate(input_fn=model._batch(), steps=1)) |
