Update celery routing for celery 4+ (#25567)

* Update celery routing

- Used routing function instead of class
- Move task queues dictionary to Django settings
- Removed routing_key parameter
- Refactored routing for singleton celery instantiation

Co-authored-by: Awais Qureshi <awais.qureshi@arbisoft.com>
This commit is contained in:
Muhammad Soban Javed
2020-12-16 13:40:47 +05:00
committed by GitHub
parent 958313c6cc
commit bd601cf3a6
21 changed files with 134 additions and 151 deletions

View File

@@ -247,7 +247,7 @@ def should_dump_course(course_key, graph):
return last_this_command_was_run < course_last_published_date
@task(routing_key=settings.COURSEGRAPH_JOB_QUEUE)
@task
@set_code_owner_attribute
def dump_course_to_neo4j(course_key_string, credentials):
"""

View File

@@ -14,11 +14,6 @@ from openedx.core.djangoapps.credentials.utils import get_credentials_api_client
logger = get_task_logger(__name__)
# Under cms the following setting is not defined, leading to errors during tests.
# These tasks aren't strictly credentials generation, but are similar in the sense
# that they generate records on the credentials side. And have a similar SLA.
ROUTING_KEY = getattr(settings, 'CREDENTIALS_GENERATION_ROUTING_KEY', None)
# Maximum number of retries before giving up.
# For reference, 11 retries with exponential backoff yields a maximum waiting
# time of 2047 seconds (about 30 minutes). Setting this to None could yield
@@ -26,7 +21,7 @@ ROUTING_KEY = getattr(settings, 'CREDENTIALS_GENERATION_ROUTING_KEY', None)
MAX_RETRIES = 11
@task(bind=True, ignore_result=True, routing_key=ROUTING_KEY)
@task(bind=True, ignore_result=True)
@set_code_owner_attribute
def send_grade_to_credentials(self, username, course_run_key, verified, letter_grade, percent_grade):
""" Celery task to notify the Credentials IDA of a grade change via POST. """

View File

@@ -4,11 +4,10 @@ A trivial task for health checks
from celery.task import task
from django.conf import settings
from edx_django_utils.monitoring import set_code_owner_attribute
@task(routing_key=settings.HEARTBEAT_CELERY_ROUTING_KEY)
@task
@set_code_owner_attribute
def sample_task():
return True

View File

@@ -23,9 +23,6 @@ from openedx.core.djangoapps.programs.utils import ProgramProgressMeter
from openedx.core.djangoapps.site_configuration import helpers as configuration_helpers
LOGGER = get_task_logger(__name__)
# Under cms the following setting is not defined, leading to errors during tests.
ROUTING_KEY = getattr(settings, 'CREDENTIALS_GENERATION_ROUTING_KEY', None)
PROGRAM_CERTIFICATES_ROUTING_KEY = getattr(settings, 'PROGRAM_CERTIFICATES_ROUTING_KEY', None)
# Maximum number of retries before giving up on awarding credentials.
# For reference, 11 retries with exponential backoff yields a maximum waiting
# time of 2047 seconds (about 30 minutes). Setting this to None could yield
@@ -124,7 +121,7 @@ def award_program_certificate(client, username, program_uuid, visible_date):
})
@task(bind=True, ignore_result=True, routing_key=PROGRAM_CERTIFICATES_ROUTING_KEY)
@task(bind=True, ignore_result=True)
@set_code_owner_attribute
def award_program_certificates(self, username):
"""
@@ -287,7 +284,7 @@ def post_course_certificate(client, username, certificate, visible_date):
})
@task(bind=True, ignore_result=True, routing_key=ROUTING_KEY)
@task(bind=True, ignore_result=True)
@set_code_owner_attribute
def award_course_certificate(self, username, course_run_key):
"""
@@ -402,7 +399,7 @@ def revoke_program_certificate(client, username, program_uuid):
})
@task(bind=True, ignore_result=True, routing_key=PROGRAM_CERTIFICATES_ROUTING_KEY)
@task(bind=True, ignore_result=True)
@set_code_owner_attribute
def revoke_program_certificates(self, username, course_key):
"""
@@ -526,7 +523,7 @@ def revoke_program_certificates(self, username, course_key):
LOGGER.info(u'Successfully completed the task revoke_program_certificates for username %s', username)
@task(bind=True, ignore_result=True, routing_key=PROGRAM_CERTIFICATES_ROUTING_KEY)
@task(bind=True, ignore_result=True)
@set_code_owner_attribute
def update_certificate_visible_date_on_course_update(self, course_key):
"""

View File

@@ -150,7 +150,7 @@ class BinnedScheduleMessageBaseTask(ScheduleMessageBaseTask):
raise NotImplementedError
@task(base=LoggedTask, ignore_result=True, routing_key=ROUTING_KEY)
@task(base=LoggedTask, ignore_result=True)
@set_code_owner_attribute
def _recurring_nudge_schedule_send(site_id, msg_str):
_schedule_send(
@@ -161,7 +161,7 @@ def _recurring_nudge_schedule_send(site_id, msg_str):
)
@task(base=LoggedTask, ignore_result=True, routing_key=ROUTING_KEY)
@task(base=LoggedTask, ignore_result=True)
@set_code_owner_attribute
def _upgrade_reminder_schedule_send(site_id, msg_str):
_schedule_send(
@@ -172,7 +172,7 @@ def _upgrade_reminder_schedule_send(site_id, msg_str):
)
@task(base=LoggedTask, ignore_result=True, routing_key=ROUTING_KEY)
@task(base=LoggedTask, ignore_result=True)
@set_code_owner_attribute
def _course_update_schedule_send(site_id, msg_str):
_schedule_send(

View File

@@ -5,62 +5,39 @@ For more, see https://celery.readthedocs.io/en/latest/userguide/routing.html#rou
"""
import logging
from abc import ABCMeta, abstractproperty
from django.conf import settings
import six
log = logging.getLogger(__name__)
class AlternateEnvironmentRouter(six.with_metaclass(ABCMeta, object)):
def route_task(name, args, kwargs, options, task=None, **kw): # pylint: disable=unused-argument
"""
A custom Router class for use in routing celery tasks to non-default queues.
Celery-defined method allowing for custom routing logic.
If None is returned from this method, default routing logic is used.
"""
if name in settings.EXPLICIT_QUEUES:
return settings.EXPLICIT_QUEUES[name]
@abstractproperty
def alternate_env_tasks(self):
"""
Defines the task -> alternate worker environment to be used when routing.
alternate_env = settings.ALTERNATE_ENV_TASKS.get(name, None)
if alternate_env:
return ensure_queue_env(alternate_env)
Subclasses must override this property with their own specific mappings.
"""
return {}
@property
def explicit_queues(self):
"""
Defines the task -> alternate worker queue to be used when routing.
"""
return {}
def ensure_queue_env(desired_env):
"""
Helper method to get the desired type of queue.
def route_for_task(self, task, args=None, kwargs=None): # pylint: disable=unused-argument
"""
Celery-defined method allowing for custom routing logic.
If None is returned from this method, default routing logic is used.
"""
if task in self.explicit_queues:
return self.explicit_queues[task]
alternate_env = self.alternate_env_tasks.get(task, None)
if alternate_env:
return self.ensure_queue_env(alternate_env)
return None
def ensure_queue_env(self, desired_env):
"""
Helper method to get the desired type of queue.
If no such queue is defined, default routing logic is used.
"""
queues = getattr(settings, 'CELERY_QUEUES', None)
return next(
(
queue
for queue in queues
if '.{}.'.format(desired_env) in queue
),
None
)
If no such queue is defined, default routing logic is used.
"""
queues = getattr(settings, 'CELERY_QUEUES', None)
return next(
(
queue
for queue in queues
if '.{}.'.format(desired_env) in queue
),
None
)