diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-04 11:17:32 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-04 11:17:32 +0100 |
| commit | 1eb1bf90e796d0814eb8e501f856cf21c1e15f60 (patch) | |
| tree | b789062cdaf79678f92f97d0fab35fe9b4b1e736 /api.py | |
| parent | 6c8bbce9c961b9876a93a8aac7e4b4276978ca59 (diff) | |
| download | shrinker-1eb1bf90e796d0814eb8e501f856cf21c1e15f60.tar.gz shrinker-1eb1bf90e796d0814eb8e501f856cf21c1e15f60.zip | |
refactor to use joblib worker class
Diffstat (limited to 'api.py')
| -rw-r--r-- | api.py | 33 |
1 files changed, 5 insertions, 28 deletions
@@ -7,7 +7,7 @@ import logging import json import asyncpg -import joblib.joblib +import joblib.worker source_db = { @@ -27,19 +27,17 @@ dest_db = { } -class Worker(object): +class Processor(joblib.worker.Worker): def __init__(self): - self._queue = None self._srcpool = None self._destpool = None self._queries = {} + super().__init__(jobtype="process") async def connect(self, sourcea, desta): """Connect to database.""" logging.warning("connecting to database") - self._queue = joblib.joblib.JobQueue() - await self._queue.connect(**sourcea) - await self._queue.setup() + await super().connect(**sourcea) self._srcpool = await asyncpg.create_pool(**sourcea) self._destpool = await asyncpg.create_pool(**desta) @@ -155,30 +153,9 @@ class Worker(object): logging.debug(query) return await conn.fetchval(query, (*data)) - async def _work(self): - """Fetch a job and run it.""" - jobid, payload, priority = await self._queue.acquire( - jobtype="process") - if jobid is None: - raise LookupError("no jobs available") - await self._execute_job(jobid, payload, priority) - 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 = Processor() await worker.connect( source_db, dest_db ) |
