mirror of https://github.com/rq/rq.git
Oops, fix some old references to current_connection.
This commit is contained in:
parent
d721f0708b
commit
227e107a82
|
@ -6,7 +6,7 @@ import procname
|
||||||
from logbook import Logger
|
from logbook import Logger
|
||||||
from pickle import loads, dumps
|
from pickle import loads, dumps
|
||||||
from .queue import Queue
|
from .queue import Queue
|
||||||
from .conn import current_connection
|
from .proxy import conn
|
||||||
|
|
||||||
class NoQueueError(Exception): pass
|
class NoQueueError(Exception): pass
|
||||||
|
|
||||||
|
@ -45,7 +45,7 @@ class Worker(object):
|
||||||
def work(self):
|
def work(self):
|
||||||
while True:
|
while True:
|
||||||
self.procline('Waiting on %s' % (', '.join(self.queue_names()),))
|
self.procline('Waiting on %s' % (', '.join(self.queue_names()),))
|
||||||
queue, msg = current_connection().blpop(self.queue_keys())
|
queue, msg = conn.blpop(self.queue_keys())
|
||||||
self.fork_and_perform_job(queue, msg)
|
self.fork_and_perform_job(queue, msg)
|
||||||
|
|
||||||
def fork_and_perform_job(self, queue, msg):
|
def fork_and_perform_job(self, queue, msg):
|
||||||
|
@ -80,7 +80,7 @@ class Worker(object):
|
||||||
else:
|
else:
|
||||||
self.log.info('Job ended normally without result')
|
self.log.info('Job ended normally without result')
|
||||||
if rv is not None:
|
if rv is not None:
|
||||||
p = current_connection().pipeline()
|
p = conn.pipeline()
|
||||||
p.set(key, dumps(rv))
|
p.set(key, dumps(rv))
|
||||||
p.expire(key, self.rv_ttl)
|
p.expire(key, self.rv_ttl)
|
||||||
p.execute()
|
p.execute()
|
||||||
|
|
Loading…
Reference in New Issue