summaryrefslogtreecommitdiff
path: root/joblib.py
diff options
context:
space:
mode:
Diffstat (limited to 'joblib.py')
-rw-r--r--joblib.py14
1 files changed, 7 insertions, 7 deletions
diff --git a/joblib.py b/joblib.py
index 6891e11..a0ce985 100644
--- a/joblib.py
+++ b/joblib.py
@@ -36,29 +36,29 @@ class JobQueue(object):
""", jobtype, json.dumps(payload), priority)
async def acquire(self, jobtype):
- """Mark a job as running, return payload and return id.
- Return (None, None) if no job is available."""
+ """Mark a job as running, return id, payload and priority.
+ Return (None, None, None) if no job is available."""
async with self._pool.acquire() as con:
while True:
try:
# do not allow async access
async with con.transaction(isolation="serializable"):
result = await con.fetchrow("""
- SELECT id, payload
+ SELECT id, payload, priority
FROM jobs WHERE
type=$1 AND status='open'
- ORDER BY priority ASC
+ ORDER BY priority DESC
""", jobtype)
if result is None:
# no jobs available
- return None, None
- jobid, payload = result
+ return None, None, None
+ jobid, payload, priority = result
await con.execute("""
UPDATE jobs
SET status='running'
WHERE id=$1
""", jobid)
- return jobid, json.loads(payload)
+ return jobid, json.loads(payload), priority
except asyncpg.exceptions.SerializationError:
# job is being picked up by another worker, try again
pass