mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-07 21:28:11 +08:00
#76 adds connection checked with timeout
Adds a check for old connections in both the workers and the monitor , every DB_TIMEOUT seconds
This commit is contained in:
@@ -319,10 +319,17 @@ def monitor(result_queue, broker=None):
|
|||||||
name = current_process().name
|
name = current_process().name
|
||||||
logger.info(_("{} monitoring at {}").format(name, current_process().pid))
|
logger.info(_("{} monitoring at {}").format(name, current_process().pid))
|
||||||
db.close_old_connections()
|
db.close_old_connections()
|
||||||
|
connection_timer = timezone.now()
|
||||||
for task in iter(result_queue.get, 'STOP'):
|
for task in iter(result_queue.get, 'STOP'):
|
||||||
|
# check db connection timeout
|
||||||
|
if (timezone.now() - connection_timer).total_seconds() >= Conf.DB_TIMEOUT:
|
||||||
|
db.close_old_connections()
|
||||||
|
connection_timer = timezone.now()
|
||||||
|
# acknowledge
|
||||||
ack_id = task.pop('ack_id', False)
|
ack_id = task.pop('ack_id', False)
|
||||||
if ack_id:
|
if ack_id:
|
||||||
broker.acknowledge(ack_id)
|
broker.acknowledge(ack_id)
|
||||||
|
# save the result
|
||||||
save_task(task)
|
save_task(task)
|
||||||
if task['success']:
|
if task['success']:
|
||||||
logger.info(_("Processed [{}]").format(task['name']))
|
logger.info(_("Processed [{}]").format(task['name']))
|
||||||
@@ -341,6 +348,7 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
|
|||||||
name = current_process().name
|
name = current_process().name
|
||||||
logger.info(_('{} ready for work at {}').format(name, current_process().pid))
|
logger.info(_('{} ready for work at {}').format(name, current_process().pid))
|
||||||
db.close_old_connections()
|
db.close_old_connections()
|
||||||
|
connection_timer = timezone.now()
|
||||||
task_count = 0
|
task_count = 0
|
||||||
# Start reading the task queue
|
# Start reading the task queue
|
||||||
for task in iter(task_queue.get, 'STOP'):
|
for task in iter(task_queue.get, 'STOP'):
|
||||||
@@ -360,6 +368,10 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
|
|||||||
result = (e, False)
|
result = (e, False)
|
||||||
# We're still going
|
# We're still going
|
||||||
if not result:
|
if not result:
|
||||||
|
# check db connection timeout
|
||||||
|
if (timezone.now() - connection_timer).total_seconds() >= Conf.DB_TIMEOUT:
|
||||||
|
db.close_old_connections()
|
||||||
|
connection_timer = timezone.now()
|
||||||
# execute the payload
|
# execute the payload
|
||||||
timer.value = task['kwargs'].pop('timeout', timeout or 0) # Busy
|
timer.value = task['kwargs'].pop('timeout', timeout or 0) # Busy
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -64,6 +64,9 @@ class Conf(object):
|
|||||||
# Failures are always saved
|
# Failures are always saved
|
||||||
SAVE_LIMIT = conf.get('save_limit', 250)
|
SAVE_LIMIT = conf.get('save_limit', 250)
|
||||||
|
|
||||||
|
# Sets the time in seconds to check for stale database connections
|
||||||
|
DB_TIMEOUT = conf.get('db_timeout', 60)
|
||||||
|
|
||||||
# Disable the scheduler
|
# Disable the scheduler
|
||||||
SCHEDULER = conf.get('scheduler', True)
|
SCHEDULER = conf.get('scheduler', True)
|
||||||
|
|
||||||
|
|||||||
@@ -113,6 +113,7 @@ def test_cluster(broker):
|
|||||||
@pytest.mark.django_db
|
@pytest.mark.django_db
|
||||||
def test_async(broker, admin_user):
|
def test_async(broker, admin_user):
|
||||||
broker.list_key = 'cluster_test:q'
|
broker.list_key = 'cluster_test:q'
|
||||||
|
Conf.DB_TIMEOUT = 0
|
||||||
broker.delete_queue()
|
broker.delete_queue()
|
||||||
a = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_cluster.assert_result',
|
a = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_cluster.assert_result',
|
||||||
broker=broker)
|
broker=broker)
|
||||||
@@ -237,6 +238,7 @@ def test_async(broker, admin_user):
|
|||||||
assert fetch(k, 100) is None
|
assert fetch(k, 100) is None
|
||||||
assert result(k, 100) is None
|
assert result(k, 100) is None
|
||||||
broker.delete_queue()
|
broker.delete_queue()
|
||||||
|
Conf.DB_TIMEOUT = 60
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.django_db
|
@pytest.mark.django_db
|
||||||
|
|||||||
Reference in New Issue
Block a user