mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-03 15:08:11 +08:00
+2
-2
@@ -58,8 +58,8 @@ target/
|
|||||||
|
|
||||||
dev-requirements.txt
|
dev-requirements.txt
|
||||||
manage.py
|
manage.py
|
||||||
testq
|
|
||||||
db.sqlite3
|
db.sqlite3
|
||||||
.venv
|
.venv
|
||||||
.idea
|
.idea
|
||||||
djq
|
djq
|
||||||
|
node_modules
|
||||||
+14
-14
@@ -72,7 +72,7 @@ class Cluster(object):
|
|||||||
self.sentinel.start()
|
self.sentinel.start()
|
||||||
logger.info(_('Q Cluster-{} starting.').format(self.pid))
|
logger.info(_('Q Cluster-{} starting.').format(self.pid))
|
||||||
while not self.start_event.is_set():
|
while not self.start_event.is_set():
|
||||||
sleep(0.2)
|
sleep(0.1)
|
||||||
return self.pid
|
return self.pid
|
||||||
|
|
||||||
def stop(self):
|
def stop(self):
|
||||||
@@ -176,7 +176,7 @@ class Sentinel(object):
|
|||||||
return self.spawn_process(pusher, self.task_queue, self.event_out, self.list_key, self.r)
|
return self.spawn_process(pusher, self.task_queue, self.event_out, self.list_key, self.r)
|
||||||
|
|
||||||
def spawn_worker(self):
|
def spawn_worker(self):
|
||||||
self.spawn_process(worker, self.task_queue, self.result_queue, Value('b', -1), self.timeout)
|
self.spawn_process(worker, self.task_queue, self.result_queue, Value('f', -1), self.timeout)
|
||||||
|
|
||||||
def spawn_monitor(self):
|
def spawn_monitor(self):
|
||||||
return self.spawn_process(monitor, self.result_queue)
|
return self.spawn_process(monitor, self.result_queue)
|
||||||
@@ -226,17 +226,18 @@ class Sentinel(object):
|
|||||||
logger.info(_('Q Cluster-{} running.').format(self.parent_pid))
|
logger.info(_('Q Cluster-{} running.').format(self.parent_pid))
|
||||||
scheduler(list_key=self.list_key)
|
scheduler(list_key=self.list_key)
|
||||||
counter = 0
|
counter = 0
|
||||||
|
cycle = 0.5 # guard loop sleep in seconds
|
||||||
# Guard loop. Runs at least once
|
# Guard loop. Runs at least once
|
||||||
while not self.stop_event.is_set() or not counter:
|
while not self.stop_event.is_set() or not counter:
|
||||||
# Check Workers
|
# Check Workers
|
||||||
for p in self.pool:
|
for p in self.pool:
|
||||||
# Are you alive?
|
# Are you alive?
|
||||||
if not p.is_alive() or (self.timeout and int(p.timer.value) == 0):
|
if not p.is_alive() or (self.timeout and p.timer.value == 0):
|
||||||
self.reincarnate(p)
|
self.reincarnate(p)
|
||||||
continue
|
continue
|
||||||
# Decrement timer if work is being done
|
# Decrement timer if work is being done
|
||||||
if p.timer.value > 0:
|
if self.timeout and p.timer.value > 0:
|
||||||
p.timer.value -= 1
|
p.timer.value -= cycle
|
||||||
# Check Monitor
|
# Check Monitor
|
||||||
if not self.monitor.is_alive():
|
if not self.monitor.is_alive():
|
||||||
self.reincarnate(self.monitor)
|
self.reincarnate(self.monitor)
|
||||||
@@ -244,13 +245,13 @@ class Sentinel(object):
|
|||||||
if not self.pusher.is_alive():
|
if not self.pusher.is_alive():
|
||||||
self.reincarnate(self.pusher)
|
self.reincarnate(self.pusher)
|
||||||
# Call scheduler once a minute (or so)
|
# Call scheduler once a minute (or so)
|
||||||
counter += 1
|
counter += cycle
|
||||||
if counter > 120:
|
if counter == 30:
|
||||||
counter = 0
|
counter = 0
|
||||||
scheduler(list_key=self.list_key)
|
scheduler(list_key=self.list_key)
|
||||||
# Save current status
|
# Save current status
|
||||||
Stat(self).save()
|
Stat(self).save()
|
||||||
sleep(0.5)
|
sleep(cycle)
|
||||||
self.stop()
|
self.stop()
|
||||||
|
|
||||||
def stop(self):
|
def stop(self):
|
||||||
@@ -261,7 +262,7 @@ class Sentinel(object):
|
|||||||
self.event_out.set()
|
self.event_out.set()
|
||||||
# Wait for it to stop
|
# Wait for it to stop
|
||||||
while self.pusher.is_alive():
|
while self.pusher.is_alive():
|
||||||
sleep(0.2)
|
sleep(0.1)
|
||||||
Stat(self).save()
|
Stat(self).save()
|
||||||
# Put poison pills in the queue
|
# Put poison pills in the queue
|
||||||
for _ in range(len(self.pool)):
|
for _ in range(len(self.pool)):
|
||||||
@@ -274,7 +275,7 @@ class Sentinel(object):
|
|||||||
for p in self.pool:
|
for p in self.pool:
|
||||||
if not p.is_alive():
|
if not p.is_alive():
|
||||||
self.pool.remove(p)
|
self.pool.remove(p)
|
||||||
sleep(0.2)
|
sleep(0.1)
|
||||||
Stat(self).save()
|
Stat(self).save()
|
||||||
# Finally stop the monitor
|
# Finally stop the monitor
|
||||||
self.result_queue.put('STOP')
|
self.result_queue.put('STOP')
|
||||||
@@ -286,8 +287,8 @@ class Sentinel(object):
|
|||||||
count = 0
|
count = 0
|
||||||
if not self.timeout:
|
if not self.timeout:
|
||||||
self.timeout = 30
|
self.timeout = 30
|
||||||
while self.status() == Conf.STOPPING and count < self.timeout * 5:
|
while self.status() == Conf.STOPPING and count < self.timeout * 10:
|
||||||
sleep(0.2)
|
sleep(0.1)
|
||||||
Stat(self).save()
|
Stat(self).save()
|
||||||
count += 1
|
count += 1
|
||||||
# Final status
|
# Final status
|
||||||
@@ -311,8 +312,7 @@ def pusher(task_queue, e, list_key=Conf.Q_LIST, r=redis_client):
|
|||||||
sleep(10)
|
sleep(10)
|
||||||
break
|
break
|
||||||
if task:
|
if task:
|
||||||
task = task[1]
|
task_queue.put(task[1])
|
||||||
task_queue.put(task)
|
|
||||||
logger.debug(_('queueing from {}').format(list_key))
|
logger.debug(_('queueing from {}').format(list_key))
|
||||||
if e.is_set():
|
if e.is_set():
|
||||||
break
|
break
|
||||||
|
|||||||
Reference in New Issue
Block a user