summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-04 11:33:28 +0100
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-04 11:33:28 +0100
commit714f26894a1a0b048c9f5a43315ac067deb9f93c (patch)
treee912de05a474a578ef164e4cc36f0bade0c2ebd2
parent814e2ca1d24b91364744f2564c37b2847313acbb (diff)
downloadcompiler-714f26894a1a0b048c9f5a43315ac067deb9f93c.tar.gz
compiler-714f26894a1a0b048c9f5a43315ac067deb9f93c.zip
refactor to use joblib worker class
-rw-r--r--api.py34
m---------joblib0
2 files changed, 6 insertions, 28 deletions
diff --git a/api.py b/api.py
index e7dbfdd..4b3c890 100644
--- a/api.py
+++ b/api.py
@@ -7,7 +7,7 @@ import logging
import json
import asyncpg
-import joblib.joblib
+import joblib.worker
queue_db = {
"host": os.environ.get("POSTGRESQL_SOURCE_HOST") or "vaindock_postgres_raw",
@@ -26,18 +26,16 @@ db_config = {
}
-class Worker(object):
+class Compiler(joblib.worker.Worker):
def __init__(self):
- self._queue = None
self._pool = None
self._queries = {}
+ super().__init__(jobtype="compile")
async def connect(self, dbconf, queuedb):
"""Connect to database."""
logging.warning("connecting to database")
- self._queue = joblib.joblib.JobQueue()
- await self._queue.connect(**queuedb)
- await self._queue.setup()
+ await super().connect(**queuedb)
self._pool = await asyncpg.create_pool(**dbconf)
async def setup(self):
@@ -55,7 +53,7 @@ class Worker(object):
self._queries[table] = [file.read()]
logging.info("loaded query '%s'", table)
- async def _execute_job(self, jobid, payload):
+ async def _execute_job(self, jobid, payload, priority):
"""Finish a job."""
object_id = payload["id"]
table = payload["type"]
@@ -68,29 +66,9 @@ class Worker(object):
async with con.transaction():
await con.execute(query, object_id)
- async def _work(self):
- """Fetch a job and run it."""
- jobid, payload, _ = await self._queue.acquire(jobtype="compile")
- if jobid is None:
- raise LookupError("no jobs available")
- await self._execute_job(jobid, payload)
- await self._queue.finish(jobid)
-
- async def run(self):
- """Start jobs forever."""
- while True:
- try:
- await self._work()
- except LookupError:
- await asyncio.sleep(1)
-
- async def start(self, number=1):
- """Start jobs in background."""
- for _ in range(number):
- asyncio.ensure_future(self.run())
async def startup():
- worker = Worker()
+ worker = Compiler()
await worker.connect(db_config, queue_db)
await worker.setup()
await worker.start(1)
diff --git a/joblib b/joblib
-Subproject 31d0588a35c8df17a5d3adcfc25af2b72c896cc
+Subproject 8d58a4f5f754b6541af84478a2df3fba70db94b