diff options
| -rw-r--r-- | api.py | 95 | ||||
| -rw-r--r-- | queries/match.sql | 2 | ||||
| -rw-r--r-- | queries/participant.sql | 2 | ||||
| -rw-r--r-- | queries/player.sql | 2 | ||||
| -rw-r--r-- | queries/roster.sql | 2 |
5 files changed, 86 insertions, 17 deletions
@@ -28,7 +28,7 @@ class Processor(object): self._queries = {} self._srcpool = self._destpool = None - def setup(self): + async def setup(self): """Load .sql files from queries/ directory. File name is the table to insert into.""" self._queries = {} @@ -42,23 +42,91 @@ class Processor(object): logging.info("loaded query for '%s'", table) self._queries[table] = file.read() + # prepare worker queue + async with self._srcpool.acquire() as con: + await con.execute(""" + CREATE TABLE IF NOT EXISTS processjobs + (id SERIAL, tablename TEXT, finished BOOL) + """) + # purge failed jobs from last time and retry + await con.execute("DELETE FROM processjobs WHERE finished=false") + 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 def acquire_job(self): + """Reserve a job and return (id, table).""" + async with self._srcpool.acquire() as con: + while True: + try: + async with con.transaction(isolation="serializable"): + # select us our job + + table_name = await con.fetchval(""" + SELECT table_name + FROM ( + SELECT table_name FROM information_schema.tables WHERE table_schema='public' AND + (table_name ~ '^(match|roster|participant)_\w{2,3}_\d{4}_\d{2}_\d{2}$' + OR table_name ~ '^(player|team)_\w{2,3}$') + ) AS table_names + WHERE table_name NOT IN + (SELECT tablename FROM processjobs) + """) + if table_name == None: + logging.warn("no jobs available") + return (None, None) + + # store our job as pending + jobid = await con.fetchval(""" + INSERT INTO processjobs(tablename, finished) + VALUES ($1, FALSE) + RETURNING id + """, table_name) + # exit loop + return (jobid, table_name) + except asyncpg.exceptions.SerializationError: + # job is being picked up by another worker, try again + pass + + async def run_job(self): + """Processes all data in table `table_name`.""" + # reserve one job + jobid, table_name = await self.acquire_job() + assert jobid is not None + logging.info("%s: processing '%s'", jobid, table_name) + + # process. + query_name = table_name.split("_")[0] # `match` async with self._srcpool.acquire() as srccon: async with self._destpool.acquire() as destcon: - for table, query in self._queries.items(): - logging.info("running query for '%s'", table) - stop_after -= 1 # TODO for debugging, don't process whole table - async with srccon.transaction(): - async with destcon.transaction(): - async for rec in srccon.cursor(query): - await self.into(destcon, rec, table) + # load SQL + # replace generic `FROM` by `FROM partition` + query = self._queries[query_name].replace( + "FROM \"" + query_name + "\"", + "FROM \"" + table_name + "\"" + ) + logging.debug("%s: using query '%s'", jobid, query_name) + async with srccon.transaction(): + async with destcon.transaction(): + async for rec in srccon.cursor(query): + # loop over src, insert into dest + await self.into(destcon, rec, query_name) + + async def run(self): + """Execute a function for each row that is fetched with query.""" + async def worker(): + """Spawn tasks forever.""" + try: + await self.run_job() + except AssertionError: + logging.info("worker idling") + await asyncio.sleep(30) + asyncio.ensure_future(worker()) + + for _ in range(5): + asyncio.ensure_future(worker()) async def into(self, conn, data, table): """Insert a named tuple into a table.""" @@ -73,9 +141,10 @@ class Processor(object): async def main(): pr = Processor() await pr.connect(source_db, dest_db) - pr.setup() - await pr.run(stop_after=1000) + await pr.setup() + await pr.run() logging.basicConfig(level=logging.DEBUG) loop = asyncio.get_event_loop() loop.run_until_complete(main()) +loop.run_forever() diff --git a/queries/match.sql b/queries/match.sql index 4f82933..2963172 100644 --- a/queries/match.sql +++ b/queries/match.sql @@ -12,4 +12,4 @@ attributes->'stats'->>'queue' AS "queue", relationships->'rosters'->'data'->0->>'id' as roster_1, relationships->'rosters'->'data'->1->>'id' as roster_2 -FROM match +FROM "match" diff --git a/queries/participant.sql b/queries/participant.sql index b7685a9..e9f8264 100644 --- a/queries/participant.sql +++ b/queries/participant.sql @@ -25,4 +25,4 @@ attributes->'stats'->'itemSells' AS "itemSells", attributes->'stats'->'itemUses' AS "itemUses", attributes->'stats'->'items' AS "items" -FROM participant +FROM "participant" diff --git a/queries/player.sql b/queries/player.sql index add926a..992cb44 100644 --- a/queries/player.sql +++ b/queries/player.sql @@ -12,4 +12,4 @@ attributes->>'name' AS "name", 0 AS "streak" -FROM player +FROM "player" diff --git a/queries/roster.sql b/queries/roster.sql index b5c4dff..eaf0ed6 100644 --- a/queries/roster.sql +++ b/queries/roster.sql @@ -14,4 +14,4 @@ relationships->'participants'->'data'->0->>'id' AS "participant_1", relationships->'participants'->'data'->1->>'id' AS "participant_2", relationships->'participants'->'data'->2->>'id' AS "participant_3" -FROM roster +FROM "roster" |
