diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-04-13 20:05:01 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-04-13 20:05:01 +0200 |
| commit | c21015de0966635d6fc651d7906ff258e6036b36 (patch) | |
| tree | f192746ccb4fb1c40812334d91ba40893cc2e75d | |
| parent | b8fcdaa9b852cce63badb2665151615e1865a02d (diff) | |
| download | analyzer-c21015de0966635d6fc651d7906ff258e6036b36.tar.gz analyzer-c21015de0966635d6fc651d7906ff258e6036b36.zip | |
make analyzer a tool not a service
| -rw-r--r-- | rolerate.py | 165 | ||||
| -rw-r--r-- | worker.py | 87 |
2 files changed, 178 insertions, 74 deletions
diff --git a/rolerate.py b/rolerate.py new file mode 100644 index 0000000..86ca2d8 --- /dev/null +++ b/rolerate.py @@ -0,0 +1,165 @@ +#!/usr/bin/python3 +import os +import time +import random +import logging +import itertools + +from sqlalchemy.orm import Session, relationship +from sqlalchemy.ext.automap import automap_base +from sqlalchemy.exc import OperationalError +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 = ParticipantStats = Player = None +db = None + +# models +mvpmodel = None + + +def connect(): + global db + global Match, Roster, Participant, Hero, ParticipantStats, Player + + # generate schema from db + Base = automap_base() + engine = create_engine(DATABASE_URI) + while True: + try: + Base.prepare(engine, reflect=True) + break + except OperationalError as err: + logging.error(err) + time.sleep(5) + + # 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)") + Participant.hero = relationship( + "hero", foreign_keys="hero.id", + primaryjoin="and_(hero.id == participant.hero_id)") + Hero = Base.classes.hero + ParticipantStats = Base.classes.participant_stats + Participant.participant_stats = relationship( + "participant_stats", foreign_keys="participant_stats.participant_api_id", + primaryjoin="and_(participant_stats.participant_api_id == participant.api_id)") + Player = Base.classes.player + + db = Session(engine) + + +class RoleModel(object): + def __init__(self, db): + self.batches = 15 + self.batchsize = 300 + self.steps = 500 + self.id = "role-kda" + + self._db = db + self._feature_cols = [ + tf.contrib.layers.real_valued_column("kills_per_min", dimension=1), + tf.contrib.layers.real_valued_column("deaths_per_min", dimension=1), + tf.contrib.layers.real_valued_column("assists_per_min", dimension=1) + ] + 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 _batch(self, ids=None, size=None): + if ids is None: # training, get random sample + size = size or self.batchsize + # train a model one batch + offset = int(random.random() * self._db.query(Participant)\ + .filter(Participant.hero.any(is_captain=True))\ + .count()) + # have a bit of randomness in the sample + records = self._db.query(Participant)\ + .filter(Participant.hero.any(is_captain=True))\ + .offset(offset).limit(size)\ + .all() + else: + records = self._db.query(Participant)\ + .filter(Hero.is_carry)\ + .filter(Participant.api_id.in_(ids))\ + .order_by(Participant.api_id.desc())\ + .all() + + data = {} + labels = [] + # populate from db records + data["kills_per_min"] = [] + data["deaths_per_min"] = [] + data["assists_per_min"] = [] + for record in records: + labels.append(record.winner) + data["kills_per_min"].append(record.participant_stats[0].kills + / record.roster[0].match[0].duration) + data["deaths_per_min"].append(record.participant_stats[0].deaths + / record.roster[0].match[0].duration) + data["assists_per_min"].append(record.participant_stats[0].assists + / record.roster[0].match[0].duration) + + logging.info("---------- %s data points for training ---------", len(labels)) + # convert to numpy arrs + for key in data: + data[key] = np.array(data[key]) + labels = np.array(labels) + + if ids is not None: + assert len(ids) == len(labels), "got nonexisting participant" + + 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 + + +logging.basicConfig(level=logging.INFO) +if __name__ == "__main__": + connect() + rolemodel = RoleModel(db) + rolemodel.train() + for name in rolemodel._model.get_variable_names(): + logging.info("%s: %s", name, + rolemodel._model.get_variable_value(name)) @@ -13,23 +13,14 @@ from sqlalchemy import create_engine import tensorflow as tf import numpy as np -import pika - -RABBITMQ_URI = os.environ.get("RABBITMQ_URI") or "amqp://localhost" DATABASE_URI = os.environ["DATABASE_URI"] -BATCHSIZE = os.environ.get("BATCHSIZE") or 1000 # objects -IDLE_TIMEOUT = os.environ.get("IDLE_TIMEOUT") or 1 # s MODEL_ROOT = os.path.join(os.getcwd(), os.path.dirname(__file__))\ + "/models/" # ORM definitions Match = Roster = Participant = ParticipantStats = Player = None -db = rabbit = channel = None - -# batch storage -queue = [] -timer = None +db = None # models mvpmodel = None @@ -37,7 +28,7 @@ mvpmodel = None def connect(): global Match, Roster, Participant, ParticipantStats, Player - global db, rabbit, channel + global db # generate schema from db Base = automap_base() @@ -64,6 +55,9 @@ def connect(): Participant.player = relationship( "player", foreign_keys="player.api_id", primaryjoin="and_(player.api_id == participant.player_api_id)") + Participant.hero = relationship( + "hero", foreign_keys="hero.id", + primaryjoin="and_(hero.id == participant.hero_id)") ParticipantStats = Base.classes.participant_stats Participant.participant_stats = relationship( "participant_stats", foreign_keys="participant_stats.participant_api_id", @@ -72,18 +66,6 @@ def connect(): db = Session(engine) - while True: - try: - rabbit = pika.BlockingConnection(pika.URLParameters(RABBITMQ_URI)) - break - except pika.exceptions.ConnectionClosed as err: - logging.error(err) - time.sleep(5) - channel = rabbit.channel() - channel.queue_declare(queue="analyze", durable=True) - channel.basic_qos(prefetch_count=BATCHSIZE) - channel.basic_consume(newjob, queue="analyze") - class Model(object): # override this configuration @@ -97,9 +79,8 @@ class Model(object): def __init__(self, db): self._db = db - for feat in self.features: - self._feature_cols = [tf.contrib.layers.real_valued_column( - feat, dimension=1)] + self._feature_cols = [tf.contrib.layers.real_valued_column( + feat, dimension=1) for feat in self.features] if self.type == "linear": self._model = tf.contrib.learn.LinearClassifier( feature_columns=self._feature_cols, @@ -183,7 +164,8 @@ class MVPScoreModel(Model): def __init__(self, db): self.features = ["participant_stats.kills", "participant_stats.deaths", - "participant_stats.assists"] + "participant_stats.assists", + "roster.hero_kills"] self.label = "participant_stats.impact_score" self.type = "linear" self.batches = 1 @@ -200,54 +182,11 @@ class MVPScoreModel(Model): return [w for l, w in self.predict(ids)] -def newjob(_, method, properties, body): - global timer, queue - queue.append((method, properties, body)) - if timer is None: - timer = rabbit.add_timeout(IDLE_TIMEOUT, process) - if len(queue) == BATCHSIZE: - process() - - -def process(): - global timer, queue, db, mvpmodel - if timer is not None: - rabbit.remove_timeout(timer) - timer = None - jobs = queue[:] - queue = [] - - logging.info("analyzing batch %s", str(len(jobs))) - ids = list(set([str(id, "utf-8") for _, _, id in jobs])) - ratings = mvpmodel.rate(ids) - - with db.no_autoflush: - for c in range(len(ratings)): - rating = int(100 * float(ratings[c])) - pstat = db.query(ParticipantStats).filter( - ParticipantStats.participant_api_id == ids[c]).first() - # manual upsert - if pstat is None: - db.add(ParticipantStats(participant_api_id=ids[c], - impact_score=rating)) - else: - pstat.impact_score = rating - - db.commit() - # ack all until this one - logging.info("acking batch") - channel.basic_ack(jobs[-1][0].delivery_tag, multiple=True) - # notify web - for api_id in ids: - channel.basic_publish("amq.topic", "participant." + api_id, - "stats_update") - - logging.basicConfig(level=logging.INFO) if __name__ == "__main__": connect() mvpmodel = MVPScoreModel(db) - mvpmodel.train() - logging.info(mvpmodel._model.evaluate(input_fn=mvpmodel._batch(), steps=1)) - - channel.start_consuming() + mvpmodel.train(force=True) + for name in mvpmodel._model.get_variable_names(): + logging.info("%s: %s", name, + mvpmodel._model.get_variable_value(name)) |
