summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+github@schneefux.xyz>2017-03-05 10:36:28 +0100
committerGitHub <noreply@github.com>2017-03-05 10:36:28 +0100
commit259dc118f12083fb3040bf8407e903e02579fe57 (patch)
treeca3b422d7a2e46c7b7f3e2cc806a0f3f305c14ac
parentd6d94e86c7b44c99dd6974453ba27637e8fb7196 (diff)
downloadcompiler-259dc118f12083fb3040bf8407e903e02579fe57.tar.gz
compiler-259dc118f12083fb3040bf8407e903e02579fe57.zip
run jobs in batches (#7)
* run jobs in batches * upsert player_stats
-rw-r--r--api.py35
m---------joblib0
-rw-r--r--queries/player/count_player_matches.sql2
3 files changed, 27 insertions, 10 deletions
diff --git a/api.py b/api.py
index 4b3c890..c88a39a 100644
--- a/api.py
+++ b/api.py
@@ -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)