diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-04 11:33:28 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-04 11:33:28 +0100 |
| commit | 714f26894a1a0b048c9f5a43315ac067deb9f93c (patch) | |
| tree | e912de05a474a578ef164e4cc36f0bade0c2ebd2 | |
| parent | 814e2ca1d24b91364744f2564c37b2847313acbb (diff) | |
| download | compiler-714f26894a1a0b048c9f5a43315ac067deb9f93c.tar.gz compiler-714f26894a1a0b048c9f5a43315ac067deb9f93c.zip | |
refactor to use joblib worker class
| -rw-r--r-- | api.py | 34 | ||||
| m--------- | joblib | 0 |
2 files changed, 6 insertions, 28 deletions
@@ -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 |
