summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-04-04 18:15:41 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-04-04 18:15:41 +0200
commit25b13cefa1ee6e39a2861caab66c4ed8bd749d67 (patch)
tree6c0b89ebc988ffc4dc29d2b02eb42f612761127f
parent3d52234edae83971e5df58f65f91c95e0352ed3f (diff)
downloadanalyzer-25b13cefa1ee6e39a2861caab66c4ed8bd749d67.tar.gz
analyzer-25b13cefa1ee6e39a2861caab66c4ed8bd749d67.zip
rewrite
-rw-r--r--.gitmodules3
-rw-r--r--Dockerfile9
-rw-r--r--api.py295
-rw-r--r--cli.py66
m---------joblib0
-rw-r--r--requirements.txt11
-rw-r--r--worker.py164
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"]
diff --git a/api.py b/api.py
deleted file mode 100644
index afc08aa..0000000
--- a/api.py
+++ /dev/null
@@ -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())
diff --git a/cli.py b/cli.py
deleted file mode 100644
index 8765e7d..0000000
--- a/cli.py
+++ /dev/null
@@ -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))