mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-07 07:18:12 +08:00
Replaces format with f-strings
This commit is contained in:
+23
-21
@@ -6,43 +6,45 @@ from django_q.conf import Conf
|
|||||||
|
|
||||||
class Disque(Broker):
|
class Disque(Broker):
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
retry = Conf.RETRY if Conf.RETRY > 0 else '{} REPLICATE 1'.format(Conf.RETRY)
|
retry = Conf.RETRY if Conf.RETRY > 0 else f"{Conf.RETRY} REPLICATE 1"
|
||||||
return self.connection.execute_command(
|
return self.connection.execute_command(
|
||||||
'ADDJOB {} {} 500 RETRY {}'.format(self.list_key, task, retry)).decode()
|
f"ADDJOB {self.list_key} {task} 500 RETRY {retry}"
|
||||||
|
).decode()
|
||||||
|
|
||||||
def dequeue(self):
|
def dequeue(self):
|
||||||
tasks = self.connection.execute_command(
|
tasks = self.connection.execute_command(
|
||||||
'GETJOB COUNT {} TIMEOUT 1000 FROM {}'.format(Conf.BULK, self.list_key))
|
f"GETJOB COUNT {Conf.BULK} TIMEOUT 1000 FROM {self.list_key}"
|
||||||
if tasks:
|
)
|
||||||
return [(t[1].decode(), t[2].decode()) for t in tasks]
|
if tasks:
|
||||||
|
return [(t[1].decode(), t[2].decode()) for t in tasks]
|
||||||
|
|
||||||
def queue_size(self):
|
def queue_size(self):
|
||||||
return self.connection.execute_command('QLEN {}'.format(self.list_key))
|
return self.connection.execute_command(f"QLEN {self.list_key}")
|
||||||
|
|
||||||
def acknowledge(self, task_id):
|
def acknowledge(self, task_id):
|
||||||
command = 'FASTACK' if Conf.DISQUE_FASTACK else 'ACKJOB'
|
command = "FASTACK" if Conf.DISQUE_FASTACK else "ACKJOB"
|
||||||
return self.connection.execute_command('{} {}'.format(command,task_id))
|
return self.connection.execute_command(f"{command} {task_id}")
|
||||||
|
|
||||||
def ping(self):
|
def ping(self):
|
||||||
return self.connection.execute_command('HELLO')[0] > 0
|
return self.connection.execute_command("HELLO")[0] > 0
|
||||||
|
|
||||||
def delete(self, task_id):
|
def delete(self, task_id):
|
||||||
return self.connection.execute_command('DELJOB {}'.format(task_id))
|
return self.connection.execute_command(f"DELJOB {task_id}")
|
||||||
|
|
||||||
def fail(self, task_id):
|
def fail(self, task_id):
|
||||||
return self.delete(task_id)
|
return self.delete(task_id)
|
||||||
|
|
||||||
def delete_queue(self):
|
def delete_queue(self):
|
||||||
jobs = self.connection.execute_command('JSCAN QUEUE {}'.format(self.list_key))[1]
|
jobs = self.connection.execute_command(f"JSCAN QUEUE {self.list_key}")[1]
|
||||||
if jobs:
|
if jobs:
|
||||||
job_ids = ' '.join(jid.decode() for jid in jobs)
|
job_ids = " ".join(jid.decode() for jid in jobs)
|
||||||
self.connection.execute_command('DELJOB {}'.format(job_ids))
|
self.connection.execute_command(f"DELJOB {job_ids}")
|
||||||
return len(jobs)
|
return len(jobs)
|
||||||
|
|
||||||
def info(self):
|
def info(self):
|
||||||
if not self._info:
|
if not self._info:
|
||||||
info = self.connection.info('server')
|
info = self.connection.info("server")
|
||||||
self._info= 'Disque {}'.format(info['disque_version'])
|
self._info = f'Disque {info["disque_version"]}'
|
||||||
return self._info
|
return self._info
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
@@ -51,15 +53,15 @@ class Disque(Broker):
|
|||||||
random.shuffle(Conf.DISQUE_NODES)
|
random.shuffle(Conf.DISQUE_NODES)
|
||||||
# find one that works
|
# find one that works
|
||||||
for node in Conf.DISQUE_NODES:
|
for node in Conf.DISQUE_NODES:
|
||||||
host, port = node.split(':')
|
host, port = node.split(":")
|
||||||
kwargs = {'host': host, 'port': port}
|
kwargs = {"host": host, "port": port}
|
||||||
if Conf.DISQUE_AUTH:
|
if Conf.DISQUE_AUTH:
|
||||||
kwargs['password'] = Conf.DISQUE_AUTH
|
kwargs["password"] = Conf.DISQUE_AUTH
|
||||||
redis_client = redis.Redis(**kwargs)
|
redis_client = redis.Redis(**kwargs)
|
||||||
redis_client.decode_responses = True
|
redis_client.decode_responses = True
|
||||||
try:
|
try:
|
||||||
redis_client.execute_command('HELLO')
|
redis_client.execute_command("HELLO")
|
||||||
return redis_client
|
return redis_client
|
||||||
except redis.exceptions.ConnectionError:
|
except redis.exceptions.ConnectionError:
|
||||||
continue
|
continue
|
||||||
raise redis.exceptions.ConnectionError('Could not connect to any Disque nodes')
|
raise redis.exceptions.ConnectionError("Could not connect to any Disque nodes")
|
||||||
|
|||||||
@@ -29,14 +29,14 @@ class Mongo(Broker):
|
|||||||
try:
|
try:
|
||||||
Conf.MONGO_DB = self.connection.get_default_database().name
|
Conf.MONGO_DB = self.connection.get_default_database().name
|
||||||
except ConfigurationError:
|
except ConfigurationError:
|
||||||
Conf.MONGO_DB = 'django-q'
|
Conf.MONGO_DB = "django-q"
|
||||||
return self.connection[Conf.MONGO_DB][self.list_key]
|
return self.connection[Conf.MONGO_DB][self.list_key]
|
||||||
|
|
||||||
def queue_size(self):
|
def queue_size(self):
|
||||||
return self.collection.count({'lock': {'$lte': _timeout()}})
|
return self.collection.count({"lock": {"$lte": _timeout()}})
|
||||||
|
|
||||||
def lock_size(self):
|
def lock_size(self):
|
||||||
return self.collection.count({'lock': {'$gt': _timeout()}})
|
return self.collection.count({"lock": {"$gt": _timeout()}})
|
||||||
|
|
||||||
def purge_queue(self):
|
def purge_queue(self):
|
||||||
return self.delete_queue()
|
return self.delete_queue()
|
||||||
@@ -46,20 +46,24 @@ class Mongo(Broker):
|
|||||||
|
|
||||||
def info(self):
|
def info(self):
|
||||||
if not self._info:
|
if not self._info:
|
||||||
self._info = 'MongoDB {}'.format(self.connection.server_info()['version'])
|
self._info = f"MongoDB {self.connection.server_info()['version']}"
|
||||||
return self._info
|
return self._info
|
||||||
|
|
||||||
def fail(self, task_id):
|
def fail(self, task_id):
|
||||||
self.delete(task_id)
|
self.delete(task_id)
|
||||||
|
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
inserted_id = self.collection.insert_one({'payload': task, 'lock': _timeout()}).inserted_id
|
inserted_id = self.collection.insert_one(
|
||||||
|
{"payload": task, "lock": _timeout()}
|
||||||
|
).inserted_id
|
||||||
return str(inserted_id)
|
return str(inserted_id)
|
||||||
|
|
||||||
def dequeue(self):
|
def dequeue(self):
|
||||||
task = self.collection.find_one_and_update({'lock': {'$lte': _timeout()}}, {'$set': {'lock': timezone.now()}})
|
task = self.collection.find_one_and_update(
|
||||||
|
{"lock": {"$lte": _timeout()}}, {"$set": {"lock": timezone.now()}}
|
||||||
|
)
|
||||||
if task:
|
if task:
|
||||||
return [(str(task['_id']), task['payload'])]
|
return [(str(task["_id"]), task["payload"])]
|
||||||
# empty queue, spare the cpu
|
# empty queue, spare the cpu
|
||||||
sleep(Conf.POLL)
|
sleep(Conf.POLL)
|
||||||
|
|
||||||
@@ -67,7 +71,7 @@ class Mongo(Broker):
|
|||||||
return self.collection.drop()
|
return self.collection.drop()
|
||||||
|
|
||||||
def delete(self, task_id):
|
def delete(self, task_id):
|
||||||
self.collection.delete_one({'_id': ObjectId(task_id)})
|
self.collection.delete_one({"_id": ObjectId(task_id)})
|
||||||
|
|
||||||
def acknowledge(self, task_id):
|
def acknowledge(self, task_id):
|
||||||
return self.delete(task_id)
|
return self.delete(task_id)
|
||||||
|
|||||||
+20
-7
@@ -27,10 +27,16 @@ class ORM(Broker):
|
|||||||
return OrmQ.objects.using(Conf.ORM)
|
return OrmQ.objects.using(Conf.ORM)
|
||||||
|
|
||||||
def queue_size(self):
|
def queue_size(self):
|
||||||
return self.get_connection().filter(key=self.list_key, lock__lte=_timeout()).count()
|
return (
|
||||||
|
self.get_connection()
|
||||||
|
.filter(key=self.list_key, lock__lte=_timeout())
|
||||||
|
.count()
|
||||||
|
)
|
||||||
|
|
||||||
def lock_size(self):
|
def lock_size(self):
|
||||||
return self.get_connection().filter(key=self.list_key, lock__gt=_timeout()).count()
|
return (
|
||||||
|
self.get_connection().filter(key=self.list_key, lock__gt=_timeout()).count()
|
||||||
|
)
|
||||||
|
|
||||||
def purge_queue(self):
|
def purge_queue(self):
|
||||||
return self.get_connection().filter(key=self.list_key).delete()
|
return self.get_connection().filter(key=self.list_key).delete()
|
||||||
@@ -40,22 +46,30 @@ class ORM(Broker):
|
|||||||
|
|
||||||
def info(self):
|
def info(self):
|
||||||
if not self._info:
|
if not self._info:
|
||||||
self._info = 'ORM {}'.format(Conf.ORM)
|
self._info = f"ORM {Conf.ORM}"
|
||||||
return self._info
|
return self._info
|
||||||
|
|
||||||
def fail(self, task_id):
|
def fail(self, task_id):
|
||||||
self.delete(task_id)
|
self.delete(task_id)
|
||||||
|
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
package = self.get_connection().create(key=self.list_key, payload=task, lock=_timeout())
|
package = self.get_connection().create(
|
||||||
|
key=self.list_key, payload=task, lock=_timeout()
|
||||||
|
)
|
||||||
return package.pk
|
return package.pk
|
||||||
|
|
||||||
def dequeue(self):
|
def dequeue(self):
|
||||||
tasks = self.get_connection().filter(key=self.list_key, lock__lt=_timeout())[0:Conf.BULK]
|
tasks = self.get_connection().filter(key=self.list_key, lock__lt=_timeout())[
|
||||||
|
0 : Conf.BULK
|
||||||
|
]
|
||||||
if tasks:
|
if tasks:
|
||||||
task_list = []
|
task_list = []
|
||||||
for task in tasks:
|
for task in tasks:
|
||||||
if self.get_connection().filter(id=task.id, lock=task.lock).update(lock=timezone.now()):
|
if (
|
||||||
|
self.get_connection()
|
||||||
|
.filter(id=task.id, lock=task.lock)
|
||||||
|
.update(lock=timezone.now())
|
||||||
|
):
|
||||||
task_list.append((task.pk, task.payload))
|
task_list.append((task.pk, task.payload))
|
||||||
# else don't process, as another cluster has been faster than us on that task
|
# else don't process, as another cluster has been faster than us on that task
|
||||||
return task_list
|
return task_list
|
||||||
@@ -70,4 +84,3 @@ class ORM(Broker):
|
|||||||
|
|
||||||
def acknowledge(self, task_id):
|
def acknowledge(self, task_id):
|
||||||
return self.delete(task_id)
|
return self.delete(task_id)
|
||||||
|
|
||||||
|
|||||||
@@ -11,7 +11,7 @@ except ImportError:
|
|||||||
|
|
||||||
class Redis(Broker):
|
class Redis(Broker):
|
||||||
def __init__(self, list_key=Conf.PREFIX):
|
def __init__(self, list_key=Conf.PREFIX):
|
||||||
super(Redis, self).__init__(list_key='django_q:{}:q'.format(list_key))
|
super(Redis, self).__init__(list_key=f"django_q:{list_key}:q")
|
||||||
|
|
||||||
def enqueue(self, task):
|
def enqueue(self, task):
|
||||||
return self.connection.rpush(self.list_key, task)
|
return self.connection.rpush(self.list_key, task)
|
||||||
@@ -34,13 +34,13 @@ class Redis(Broker):
|
|||||||
try:
|
try:
|
||||||
return self.connection.ping()
|
return self.connection.ping()
|
||||||
except redis.ConnectionError as e:
|
except redis.ConnectionError as e:
|
||||||
logger.error('Can not connect to Redis server.')
|
logger.error("Can not connect to Redis server.")
|
||||||
raise e
|
raise e
|
||||||
|
|
||||||
def info(self):
|
def info(self):
|
||||||
if not self._info:
|
if not self._info:
|
||||||
info = self.connection.info('server')
|
info = self.connection.info("server")
|
||||||
self._info = 'Redis {}'.format(info['redis_version'])
|
self._info = f"Redis {info['redis_version']}"
|
||||||
return self._info
|
return self._info
|
||||||
|
|
||||||
def set_stat(self, key, value, timeout):
|
def set_stat(self, key, value, timeout):
|
||||||
|
|||||||
@@ -8,34 +8,44 @@ from django_q.monitor import info, get_ids
|
|||||||
|
|
||||||
class Command(BaseCommand):
|
class Command(BaseCommand):
|
||||||
# Translators: help text for qinfo management command
|
# Translators: help text for qinfo management command
|
||||||
help = _('General information over all clusters.')
|
help = _("General information over all clusters.")
|
||||||
|
|
||||||
def add_arguments(self, parser):
|
def add_arguments(self, parser):
|
||||||
parser.add_argument(
|
parser.add_argument(
|
||||||
'--config',
|
"--config",
|
||||||
action='store_true',
|
action="store_true",
|
||||||
dest='config',
|
dest="config",
|
||||||
default=False,
|
default=False,
|
||||||
help='Print current configuration.',
|
help="Print current configuration.",
|
||||||
)
|
)
|
||||||
parser.add_argument(
|
parser.add_argument(
|
||||||
'--ids',
|
"--ids",
|
||||||
action='store_true',
|
action="store_true",
|
||||||
dest='ids',
|
dest="ids",
|
||||||
default=False,
|
default=False,
|
||||||
help='Print cluster task ID(s) (PIDs).',
|
help="Print cluster task ID(s) (PIDs).",
|
||||||
)
|
)
|
||||||
|
|
||||||
def handle(self, *args, **options):
|
def handle(self, *args, **options):
|
||||||
if options.get('ids', True):
|
if options.get("ids", True):
|
||||||
get_ids()
|
get_ids()
|
||||||
elif options.get('config', False):
|
elif options.get("config", False):
|
||||||
hide = ['conf', 'IDLE', 'STOPPING', 'STARTING', 'WORKING', 'SIGNAL_NAMES', 'STOPPED']
|
hide = [
|
||||||
settings = [a for a in dir(Conf) if not a.startswith('__') and a not in hide]
|
"conf",
|
||||||
self.stdout.write('VERSION: {}'.format('.'.join(str(v) for v in VERSION)))
|
"IDLE",
|
||||||
|
"STOPPING",
|
||||||
|
"STARTING",
|
||||||
|
"WORKING",
|
||||||
|
"SIGNAL_NAMES",
|
||||||
|
"STOPPED",
|
||||||
|
]
|
||||||
|
settings = [
|
||||||
|
a for a in dir(Conf) if not a.startswith("__") and a not in hide
|
||||||
|
]
|
||||||
|
self.stdout.write(f"VERSION: {'.'.join(str(v) for v in VERSION)}")
|
||||||
for setting in settings:
|
for setting in settings:
|
||||||
value = getattr(Conf, setting)
|
value = getattr(Conf, setting)
|
||||||
if value is not None:
|
if value is not None:
|
||||||
self.stdout.write('{}: {}'.format(setting, value))
|
self.stdout.write(f"{setting}: {value}")
|
||||||
else:
|
else:
|
||||||
info()
|
info()
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ def count_letters2(obj):
|
|||||||
return count_letters(obj.get_words())
|
return count_letters(obj.get_words())
|
||||||
|
|
||||||
|
|
||||||
def word_multiply(x, word=''):
|
def word_multiply(x, word=""):
|
||||||
return len(word) * x
|
return len(word) * x
|
||||||
|
|
||||||
|
|
||||||
@@ -39,8 +39,8 @@ def get_user_id(user):
|
|||||||
|
|
||||||
|
|
||||||
def hello():
|
def hello():
|
||||||
return 'hello'
|
return "hello"
|
||||||
|
|
||||||
|
|
||||||
def result(obj):
|
def result(obj):
|
||||||
print('RESULT HOOK {} : {}'.format(obj.name, obj.result))
|
print(f"RESULT HOOK {obj.name} : {obj.result()}")
|
||||||
|
|||||||
Reference in New Issue
Block a user