From 714f26894a1a0b048c9f5a43315ac067deb9f93c Mon Sep 17 00:00:00 2001 From: schneefux Date: Sat, 4 Mar 2017 11:33:28 +0100 Subject: refactor to use joblib worker class --- api.py | 34 ++++++---------------------------- joblib | 2 +- 2 files changed, 7 insertions(+), 29 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 index 31d0588..8d58a4f 160000 --- a/joblib +++ b/joblib @@ -1 +1 @@ -Subproject commit 31d0588a35c8df17a5d3adcfc25af2b72c896cc7 +Subproject commit 8d58a4f5f754b6541af84478a2df3fba70db94be -- cgit v1.3.1