diff options
| author | schneefux <schneefux+github@schneefux.xyz> | 2017-03-05 10:36:28 +0100 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2017-03-05 10:36:28 +0100 |
| commit | 259dc118f12083fb3040bf8407e903e02579fe57 (patch) | |
| tree | ca3b422d7a2e46c7b7f3e2cc806a0f3f305c14ac | |
| parent | d6d94e86c7b44c99dd6974453ba27637e8fb7196 (diff) | |
| download | compiler-259dc118f12083fb3040bf8407e903e02579fe57.tar.gz compiler-259dc118f12083fb3040bf8407e903e02579fe57.zip | |
run jobs in batches (#7)
* run jobs in batches
* upsert player_stats
| -rw-r--r-- | api.py | 35 | ||||
| m--------- | joblib | 0 | ||||
| -rw-r--r-- | queries/player/count_player_matches.sql | 2 |
3 files changed, 27 insertions, 10 deletions
@@ -53,25 +53,40 @@ class Compiler(joblib.worker.Worker): self._queries[table] = [file.read()] logging.info("loaded query '%s'", table) + + async def _windup(self): + self._con = await self._pool.acquire() + self._tr = self._con.transaction() + await self._tr.start() + + async def _teardown(self, failed): + if failed: + await self._tr.rollback() + else: + await self._tr.commit() + await self._pool.release(self._con) + async def _execute_job(self, jobid, payload, priority): """Finish a job.""" object_id = payload["id"] table = payload["type"] if table not in self._queries: return - async with self._pool.acquire() as con: - logging.debug("%s: compiling '%s' from '%s'", - jobid, object_id, table) - for query in self._queries[table]: - async with con.transaction(): - await con.execute(query, object_id) + logging.debug("%s: compiling '%s' from '%s'", + jobid, object_id, table) + tasks = [] + for query in self._queries[table]: + tasks.append(asyncio.ensure_future( + self._con.execute(query, object_id))) + await asyncio.gather(*tasks) async def startup(): - worker = Compiler() - await worker.connect(db_config, queue_db) - await worker.setup() - await worker.start(1) + for _ in range(1): + worker = Compiler() + await worker.connect(db_config, queue_db) + await worker.setup() + await worker.start(batchlimit=20) logging.basicConfig( diff --git a/joblib b/joblib -Subproject 8d58a4f5f754b6541af84478a2df3fba70db94b +Subproject 7799048a5927670b9492f80cdbbf2fd8d319c50 diff --git a/queries/player/count_player_matches.sql b/queries/player/count_player_matches.sql index fbdd283..226cacc 100644 --- a/queries/player/count_player_matches.sql +++ b/queries/player/count_player_matches.sql @@ -25,3 +25,5 @@ select ) INSERT INTO player_stats(player_api_id, patch_version, played, wins, casual_played, casual_wins, ranked_played, ranked_wins, streak) SELECT * FROM stats +ON CONFLICT(player_api_id) DO + UPDATE SET (player_api_id, patch_version, played, wins, casual_played, casual_wins, ranked_played, ranked_wins, streak) = (SELECT * FROM stats) |
