From e2f5d6aa939bdc08f2a94fdb62e256b2cc4669d2 Mon Sep 17 00:00:00 2001 From: schneefux Date: Mon, 20 Feb 2017 17:34:09 +0100 Subject: partition data into tables per day and region --- api.py | 48 ++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 38 insertions(+), 10 deletions(-) (limited to 'api.py') diff --git a/api.py b/api.py index d178828..1dcad0f 100644 --- a/api.py +++ b/api.py @@ -42,6 +42,22 @@ class Apigrabber(object): (id SERIAL, start_date TIMESTAMP, end_date TIMESTAMP, finished BOOL, region TEXT) """) + # create master tables + for objecttype in ["match", "roster", "participant", + "team", "player"]: + await con.execute(""" + CREATE TABLE IF NOT EXISTS """ + objecttype + """ ( + id TEXT PRIMARY KEY NOT NULL, + type TEXT NOT NULL, + attributes JSONB, + relationships JSONB) + """) + # create region partitions + # TODO use CHECK for regions once every object has a shardId + for region in self.regions: + await con.execute("CREATE TABLE IF NOT EXISTS " + + objecttype + "_" + region + + " (id TEXT PRIMARY KEY) INHERITS (" + objecttype + ")") # create past zombie job that marks the last data to fetch async with con.transaction(): @@ -55,21 +71,31 @@ class Apigrabber(object): ON CONFLICT DO NOTHING; """, json.dumps([{"region": r} for r in self.regions])) - async def _db_insert(self, con, objects): + async def _db_insert(self, con, objects, ddate, region): """Insert a list of API response objects into respective tables.""" + day = ddate.strftime("%Y_%m_%d") + + def table(objtype): + """Return the partition the object belongs in.""" + return (objtype + "_" + region + + ("_" + day + if objtype != "player" + and objtype != "team" + else "")) + objectmap = {} for o in objects: try: objectmap[o["type"]].append(o) except KeyError: objectmap[o["type"]] = [o] - await con.execute(""" - CREATE TABLE IF NOT EXISTS """ + o["type"] + """ ( - id TEXT PRIMARY KEY NOT NULL, - type TEXT NOT NULL, - attributes JSONB, - relationships JSONB) - """) + # create a partition for the day + # players are not partitioned by day + await con.execute("CREATE TABLE IF NOT EXISTS " + + table(o["type"]) + + " (id TEXT PRIMARY KEY) INHERITS (" + + o["type"] + "_" + region + + ")") for otype, objs in objectmap.items(): async with con.transaction(): # create savepoint @@ -83,7 +109,7 @@ class Apigrabber(object): # recent. TODO always keep the update condition in line # with the data that is returned by the API. await con.execute(""" - INSERT INTO """ + otype + """ AS j + INSERT INTO """ + table(otype) + """ AS j SELECT DISTINCT ON(id) * FROM JSONB_TO_RECORDSET($1::JSONB) AS jsn(id TEXT, type TEXT, attributes JSONB, relationships JSONB) @@ -117,7 +143,9 @@ class Apigrabber(object): logging.info("%s: (%s) received %s data objects", region, jobid, len(matches)) - await self._db_insert(con, matches) + # TODO to be more precise, jobs that return date over midnight + # should be split so they insert in the according table + await self._db_insert(con, matches, jobstart, region) logging.info("%s: (%s) inserted", region, jobid) -- cgit v1.3.1 From 7425736ee1e9db515a8a772ddd515690ade67285 Mon Sep 17 00:00:00 2001 From: schneefux Date: Mon, 20 Feb 2017 19:48:55 +0100 Subject: work around table creation race condition, delete unfinished crawljobs on start, fix async error --- api.py | 36 +++++++++++++++++++++--------------- 1 file changed, 21 insertions(+), 15 deletions(-) (limited to 'api.py') diff --git a/api.py b/api.py index 1dcad0f..2ecc2d2 100644 --- a/api.py +++ b/api.py @@ -84,18 +84,23 @@ class Apigrabber(object): else "")) objectmap = {} - for o in objects: - try: - objectmap[o["type"]].append(o) - except KeyError: - objectmap[o["type"]] = [o] - # create a partition for the day - # players are not partitioned by day - await con.execute("CREATE TABLE IF NOT EXISTS " + - table(o["type"]) + - " (id TEXT PRIMARY KEY) INHERITS (" + - o["type"] + "_" + region + - ")") + async with con.transaction(): # create savepoint + for o in objects: + try: + objectmap[o["type"]].append(o) + except KeyError: + objectmap[o["type"]] = [o] + # create a partition for the day + # players are not partitioned by day + try: + await con.execute("CREATE TABLE IF NOT EXISTS " + + table(o["type"]) + + " (id TEXT PRIMARY KEY) INHERITS (" + + o["type"] + "_" + region + + ")") + except asyncpg.DuplicateTableError: + # ninja'd by another worker + pass for otype, objs in objectmap.items(): async with con.transaction(): # create savepoint @@ -177,10 +182,11 @@ class Apigrabber(object): WHERE start_date-previous_end>INTERVAL '0' -- get gaps ORDER BY start_date DESC LIMIT 1 """, region, default_diff) - if row_res == None: + if row_res is None: logging.warn("%s: no jobs available. idling.", region) await asyncio.sleep(60) # a minute TODO make this smarter asyncio.ensure_future(self.crawl_region(region)) + return jobdate, delta_minutes = row_res delta = datetime.timedelta(minutes=delta_minutes) @@ -242,8 +248,8 @@ class Apigrabber(object): async def setup(self): async with self._pool.acquire() as con: await self._db_setup(con) - ## clean up after force quit (TODO - disabled for dev) - #await con.execute("DELETE FROM crawljobs WHERE finished=false") + # clean up after force quit (TODO - disabled for dev) + await con.execute("DELETE FROM crawljobs WHERE finished=false") async def startup(): -- cgit v1.3.1 From baaed2458bbd3bb6b3d3ba597b7558e8bca66b7d Mon Sep 17 00:00:00 2001 From: Kapil Viren Ahuja Date: Tue, 21 Feb 2017 10:06:19 +0530 Subject: ghost commit to re-initiate a PR --- api.py | 1 + 1 file changed, 1 insertion(+) (limited to 'api.py') diff --git a/api.py b/api.py index 2ecc2d2..aafe6a5 100644 --- a/api.py +++ b/api.py @@ -268,3 +268,4 @@ logging.basicConfig(level=logging.DEBUG) loop = asyncio.get_event_loop() loop.run_until_complete(startup()) loop.run_forever() + \ No newline at end of file -- cgit v1.3.1