Merge pull request #347 from maerteijn/fix-scheduler-concurrency

Fix scheduler concurrency with multiple clusters
This commit is contained in:
Ilan Steemers
2019-08-13 14:13:14 +02:00
committed by GitHub
2 changed files with 65 additions and 67 deletions
+3 -5
View File
@@ -203,7 +203,6 @@ class Sentinel(object):
self.start_event.set() self.start_event.set()
Stat(self).save() Stat(self).save()
logger.info(_('Q Cluster-{} running.').format(self.parent_pid)) logger.info(_('Q Cluster-{} running.').format(self.parent_pid))
scheduler(broker=self.broker)
counter = 0 counter = 0
cycle = Conf.GUARD_CYCLE # guard loop sleep in seconds cycle = Conf.GUARD_CYCLE # guard loop sleep in seconds
# Guard loop. Runs at least once # Guard loop. Runs at least once
@@ -492,7 +491,8 @@ def scheduler(broker=None):
broker = get_broker() broker = get_broker()
db.close_old_connections() db.close_old_connections()
try: try:
for s in Schedule.objects.exclude(repeats=0).filter(next_run__lt=timezone.now()): with db.transaction.atomic():
for s in Schedule.objects.select_for_update().exclude(repeats=0).filter(next_run__lt=timezone.now()):
args = () args = ()
kwargs = {} kwargs = {}
# get args, kwargs and hook # get args, kwargs and hook
@@ -530,9 +530,7 @@ def scheduler(broker=None):
next_run = next_run.shift(years=+1) next_run = next_run.shift(years=+1)
if Conf.CATCH_UP or next_run > arrow.utcnow(): if Conf.CATCH_UP or next_run > arrow.utcnow():
break break
# arrow always returns a tz aware datetime, and we don't want s.next_run = next_run.datetime
# this when we explicitly configured django with USE_TZ=False
s.next_run = next_run.datetime if settings.USE_TZ else next_run.datetime.replace(tzinfo=None)
s.repeats += -1 s.repeats += -1
# send it to the cluster # send it to the cluster
q_options['broker'] = broker q_options['broker'] = broker
+2 -2
View File
@@ -44,8 +44,8 @@ extensions = [
] ]
intersphinx_mapping = {'python': ('https://docs.python.org/3.5', None), intersphinx_mapping = {'python': ('https://docs.python.org/3.5', None),
'django': ('https://docs.djangoproject.com/en/1.8/', 'django': ('https://docs.djangoproject.com/en/2.2/',
'https://docs.djangoproject.com/en/1.8//_objects/')} 'https://docs.djangoproject.com/en/2.2/_objects/')}
# Add any paths that contain templates here, relative to this directory. # Add any paths that contain templates here, relative to this directory.
templates_path = ['_templates'] templates_path = ['_templates']