summaryrefslogtreecommitdiff
path: root/api.py
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-02-13 21:40:10 +0100
committerschneefux <schneefux+commit@schneefux.xyz>2017-02-13 21:50:29 +0100
commit4a23340f168ff8342b96be49a02616be54f73fba (patch)
tree25711134d7629dbc55a2cc6b0ceec9fc6e761d71 /api.py
parent4ef6c2e89409efca4dc7bae21d9648605a0f2828 (diff)
downloadprocessor-4a23340f168ff8342b96be49a02616be54f73fba.tar.gz
processor-4a23340f168ff8342b96be49a02616be54f73fba.zip
convert samples for each table
Diffstat (limited to 'api.py')
-rw-r--r--api.py180
1 files changed, 131 insertions, 49 deletions
diff --git a/api.py b/api.py
index 921993f..ce0bf56 100644
--- a/api.py
+++ b/api.py
@@ -22,50 +22,23 @@ dest_db = {
}
-class Database(object):
- async def connect(self, **connect_kwargs):
- """Connect to database by arguments."""
- self._pool = await asyncpg.create_pool(**connect_kwargs)
-
- async def each(self, fetchquery, func, **func_args):
- """Execute a function for each row that is fetched with query."""
- # TODO use iterator instead of callback
- async with self._pool.acquire() as conn:
- async with conn.transaction():
- tasks = []
- async for record in conn.cursor(fetchquery):
- tasks.append(
- asyncio.ensure_future(func(record, **func_args))
- )
- await asyncio.gather(*tasks)
-
- async def into(self, data, table):
- """Insert a named tuple into a database."""
- 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 {} (\"{}\") VALUES ({})".format(
- table, "\", \"".join(keys), ", ".join(placeholders))
+class Processor(object):
+ def __init__(self):
+ self._queries = {}
+ self._srcpool = self._destpool = None
- async with self._pool.acquire() as conn:
- await conn.fetch(query, (*data))
-
-
-async def process():
- sdb = Database()
- await sdb.connect(**source_db)
- ddb = Database()
- await ddb.connect(**dest_db)
-
- selqry = """
+ def setup(self):
+ """Load conversion queries."""
+ self._queries = {
+ "match": """
SELECT
-(data->'attributes'->>'duration')::int AS "duration",
-data->'attributes'->>'gameMode' AS "gameMode",
-COALESCE(NULLIF(data->'attributes'->>'patchVersion', ''), '0')::int AS "patchVersion",
-data->'attributes'->>'shardId' AS "shard",
-data->'attributes'->'stats'->>'endGameReason' AS "result",
-data->'attributes'->'stats'->>'queue' AS "queue",
+(attributes->>'duration')::int AS "duration",
+attributes->>'gameMode' AS "gameMode",
+COALESCE(NULLIF(attributes->>'patchVersion', ''), '0')::int AS "patchVersion",
+attributes->>'shardId' AS "shard",
+attributes->'stats'->>'endGameReason' AS "result",
+attributes->'stats'->>'queue' AS "queue",
false AS "anyAFK",
'' AS "winningTeam",
@@ -74,13 +47,122 @@ false AS "anyAFK",
0 AS "jungleMinionsSlayed",
0 AS "heroDeaths"
-FROM match LIMIT 100
- """
+FROM apidata WHERE type='match'
+ """,
+ "match_results": """
+SELECT
+
+attributes->'stats'->>'side' AS "team",
+FALSE AS "winner",
+FALSE AS "surrender",
+0 AS "teamSize",
+'' AS "hero_1",
+'' AS "hero_2",
+'' AS "hero_3",
+0 AS "laneMinionsSlayed",
+0 AS "jungleMinionsSlayed",
+0 AS "turretsDestroyed",
+(attributes->'stats'->>'heroKills')::int AS "heroKills",
+0 AS "heroAssits",
+0 AS "heroDeaths",
+(attributes->'stats'->>'krakenCaptures')::int AS "krakensCaptured",
+0 AS "goldMineCaptures",
+0 AS "crystalMineCaptures",
+0 AS "afkCount",
+0 AS "afkTime"
+
+FROM apidata WHERE type='roster'
+ """,
+ "match_participation": """
+SELECT
+
+0 AS "rosterId",
+attributes->>'actor' AS "hero",
+(attributes->'stats'->>'kills')::int AS "kills",
+(attributes->'stats'->>'deaths')::int AS "deaths",
+(attributes->'stats'->>'assists')::int AS "assists",
+0.0 AS "kd",
+0.0 AS "kda",
+(attributes->'stats'->>'nonJungleMinionKills')::int AS "laneMinionsSlayed",
+(attributes->'stats'->>'minionKills')::int AS "jungleMinionsSlayed",
+(attributes->'stats'->>'turretCaptures')::int AS "turretsDestroyed",
+0 AS "heroKills",
+0 AS "heroDeaths",
+0 AS "heroAssits",
+(attributes->'stats'->>'krakenCaptures')::int AS "krakensCaptured",
+(attributes->'stats'->>'goldMineCaptures')::int AS "goldMineCaptures",
+(attributes->'stats'->>'crystalMineCaptures')::int AS "crystalMineCaptures",
+(attributes->'stats'->>'wentAfk')::bool::int AS "afkCount",
+(attributes->'stats'->>'firstAfkTime')::float::int AS "afkTime",
+FALSE AS "perfectGame"
+
+FROM apidata WHERE type='participant'
+ """,
+ "player": """
+SELECT
+
+id AS "apiId",
+attributes->>'name' AS "name",
+(attributes->'stats'->>'level')::int AS "level",
+(attributes->'stats'->>'xp')::int AS "xp",
+(attributes->'stats'->>'played')::int AS "played",
+(attributes->'stats'->>'played_ranked')::int AS "playedRanked",
+(attributes->'stats'->>'wins')::int AS "wins",
+(attributes->'stats'->>'winStreak')::int AS "streak",
+'' AS "herosUnlocked",
+'' "skinsUnlocked",
+0 AS "totalGamePlaytime",
+attributes->'stats'->>'lifetimeGold' AS "lifeTimeGold",
+0 AS "lifeTimeKills",
+0 AS "lifeTimeDeaths",
+0 AS "lifeTimeAssists",
+'' AS "bestHeroPlayingWith",
+'' AS "worstHeroPlayingWith",
+'' AS "bestHeroPlayingAgainst",
+'' AS "worstHeroPlayingAgainst",
+0 AS "lifeTimeKD",
+0 AS "lifeTimeKDA",
+'[]'::jsonb AS "heroPerformance",
+'[]'::jsonb AS "rolePerformance"
+
+FROM apidata WHERE type='player'
+ """
+ }
+
+ async def connect(self, source, dest):
+ """Connect to database by arguments."""
+ self._srcpool = await asyncpg.create_pool(**source)
+ self._destpool = await asyncpg.create_pool(**dest)
+
+ async def run(self, stop_after=-1):
+ """Execute a function for each row that is fetched with query."""
+ # TODO unnest
+ async with self._srcpool.acquire() as srccon:
+ async with self._destpool.acquire() as destcon:
+ for table, query in self._queries.items():
+ stop_after -= 1 # TODO for debugging, don't process whole table
+ logging.info("processing table %s", table)
+ async with srccon.transaction():
+ async with destcon.transaction():
+ async for rec in srccon.cursor(query):
+ await self.into(destcon, rec, table)
+
+ async def into(self, conn, data, table):
+ """Insert a named tuple into a table."""
+ 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 {} (\"{}\") VALUES ({})".format(
+ table, "\", \"".join(keys), ", ".join(placeholders))
+ await conn.execute(query, (*data))
+
- await sdb.each(selqry, ddb.into, table="match")
+async def main():
+ pr = Processor()
+ await pr.connect(source_db, dest_db)
+ pr.setup()
+ await pr.run(stop_after=1000)
-if __name__ == "__main__":
- loop = asyncio.get_event_loop()
- loop.run_until_complete(
- process()
- )
+logging.basicConfig(level=logging.DEBUG)
+loop = asyncio.get_event_loop()
+loop.run_until_complete(main())