【问题标题】:Celery: Routing tasks issue - only one worker consume all tasks from all queuesCelery:路由任务问题 - 只有一名工作人员使用所有队列中的所有任务
【发布时间】:2019-11-19 22:14:52
【问题描述】:

我有一些任务具有手动配置的路线和 3 个工作人员,这些工作人员配置为使用特定队列中的任务。但是只有一名工作人员完成了所有任务,我不知道如何解决这个问题。

我的celeryconfig.py

    class CeleryConfig:
    enable_utc = True
    timezone = 'UTC'

    imports = ('events.tasks')
    broker_url = Config.BROKER_URL
    broker_transport_options = {'visibility_timeout': 10800}  # 3H

    worker_hijack_root_logger = False

    task_protocol = 2
    task_ignore_result = True
    task_publish_retry_policy = {'max_retries': 3, 'interval_start': 0, 'interval_step': 0.2, 'interval_max': 0.2}
    task_time_limit = 30  # sec
    task_soft_time_limit = 15  # sec

    task_default_queue = 'low'
    task_default_exchange = 'low'
    task_default_routing_key = 'low'

    task_queues = (
        Queue('daily', Exchange('daily'), routing_key='daily'),
        Queue('high', Exchange('high'), routing_key='high'),
        Queue('normal', Exchange('normal'), routing_key='normal'),
        Queue('low', Exchange('low'), routing_key='low'),
        Queue('service', Exchange('service'), routing_key='service'),
        Queue('award', Exchange('award'), routing_key='award'),
    )

    task_route = {
        # -- SCHEDULE QUEUE --
        base_path.format(task='refresh_rank'): {'queue': 'daily'}
        # -- HIGH QUEUE --
        base_path.format(task='execute_order'): {'queue': 'high'},
        # -- NORMAL QUEUE --
        base_path.format(task='calculate_cost'): {'queue': 'normal'},
        # -- SERVICE QUEUE --
        base_path.format(task='send_pin'): {'queue': 'service'},
        # -- LOW QUEUE
        base_path.format(task='invite_to_tournament'): {'queue': 'low'},
        # -- AWARD QUEUE
        base_path.format(task='get_lesson_award'): {'queue': 'award'},
        # -- TEST TASK

    worker_concurrency = multiprocessing.cpu_count() * 2 + 1
    worker_prefetch_multiplier = 1  #
    worker_max_tasks_per_child = 1
    worker_max_memory_per_child = 90000  # 90MB

    beat_max_loop_interval = 60 * 5  # 5 min

我在 docker 中运行工人,这是我的 stack.yml 的一部分

    version: "3.7"

    services:

      worker_high:
        command: celery worker -l debug -A runcelery.celery -Q high -n worker.high@%h

      worker_normal:
        command: celery worker -l debug -A runcelery.celery -Q normal,award,service,low -n worker.normal@%h

      worker_schedule:
        command: celery worker -l debug -A runcelery.celery -Q daily -n worker.schedule@%h

      beat:
        command: celery beat -l debug -A runcelery.celery

      flower:
        command: flower -l debug -A runcelery.celery --port=5555 

      broker:
        image: redis:5.0-alpine

我认为我的配置是正确的并且运行命令也是正确的,但是 docker logs 和flower 显示只有 worker.normal 消耗所有任务。

更新

这是task.py的一部分:

def refresh_rank_in_tournaments():
    logger.debug(f'Start task refresh_rank_in_tournaments')
    return AnalyticBackgroundManager.refresh_tournaments_rank()

base_path 是完整任务路径的快捷方式: base_path = 'events.tasks.{task}'

execute_order任务代码:

    @celery.task(bind=True, default_retry_delay=5)
    def execute_order(self, private_id, **kwargs):
        try:
            return OrderBackgroundManager.execute_order(private_id, **kwargs)
        except IEXException as exc:
            raise self.retry(exc=exc)

此任务将在视图中调用tasks.execute_order.delay(id)

【问题讨论】:

  • 你能添加'tasks.py'@kotmsk
  • 什么是base_path?并添加tasks.py
  • 我已更新主题消息。

标签: celery celery-task


【解决方案1】:

您的 worker.normal 订阅了 normal,award,service,low 队列。此外,low 队列是默认队列,因此每个没有明确设置队列的任务都将在 worker.normal 上执行。

【讨论】:

  • 我在celeryconfig.py中指定了队列和路由
  • 好的,但是任务execute_order 在worker.normal 上执行,尽管它已经明确设置了队列high
  • 你的 Flower 截图没有说明执行了哪些任务,所以我不可能知道那里发生了什么......另外你没有给我们execute_order 的代码,以及你如何调用它...
  • 任务执行得很好,数字195是执行任务的数量,毫无疑问。我已经更新了我最初的帖子并添加了一些代码。
  • 您已经为您的工作人员明确设置了队列,并且只有 worker.normal 订阅了设置为默认值的 low 队列。因此,当您运行 tasks.send_pin.delay(pin, id) 时,它将进入默认队列(在您的设置中为 low)。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2012-08-11
  • 1970-01-01
  • 2012-04-22
  • 2022-10-18
  • 2023-04-09
  • 1970-01-01
  • 2021-12-12
相关资源
最近更新 更多