【问题标题】:Celery: prevent adding more tasks when too many queued芹菜:当排队太多时防止添加更多任务
【发布时间】:2019-08-19 09:05:09
【问题描述】:

我有一个 Flask REST API,它利用 Celery 运行异步请求。

这个想法是async=1 查询参数指示应异步处理请求(立即返回客户端稍后将使用的任务 ID)。

同时我想在等待处理的任务过多时阻止接受新任务

下面的代码有效,但accepting_new_tasks() 需要大约 2 秒,这太慢了。

Celery 中是否有允许限制等待任务数量的配置(或其他东西);还是更快的方法来获取等待任务的数量?

import math

from celery import Celery
from flask import abort, Flask, jsonify, request


flask_app = Flask(__name__)
celery_app = Celery("tasks", broker="rabbit...")


@flask_app.route("/")
def home():
    async_ = request.args.get("async")
    settings = request.args.get("settings")

    if async_:
        if not accepting_new_tasks(celery_app):
            return abort(503)

        task = celery_app.send_task(name="my-task", kwargs={"settings": settings})
        return jsonify({"taskId": task.id})

    return jsonify({})


def accepting_new_tasks(celery_app):
    inspector = celery_app.control.inspect()
    nodes_stats = inspector.stats()
    nodes_reserved = inspector.reserved()

    workers = 0
    for stats in nodes_stats.values():
        workers += stats["pool"]["max-concurrency"]

    waiting_tasks = 0
    for reserved in nodes_reserved.values():
        waiting_tasks += len(reserved)

    return waiting_tasks < math.ceil(workers / 3)

【问题讨论】:

  • 你能补充一些细节吗?这是为了什么?我的意思是,你想解决什么问题?
  • @DanilaGanchar 我想处理大量数据,根据传递的settings 可能需要数小时。
  • 好的。现在让我们想象一下所有工人都在工作(现在12:12:12.52....)。您返回了abort(503),但在12:12:12.53.... 中,一些工作人员变得可用。这是正确的行为吗?我的意思是你有等待的时间范围吗?
  • @DanilaGanchar 不,没有等待的时间范围。这个想法是在每个请求之前检查。无论如何,我已经通过使用 RabbitMQ 管理 API 解决了它:stackoverflow.com/a/27074594/4183498

标签: python asynchronous flask celery amqp


【解决方案1】:

最终我通过查询 RabbitMQ 管理 API 解决了这个问题,正如https://stackoverflow.com/a/27074594/4183498 指出的那样。

import math

from celery import Celery
from flask import abort, Flask, jsonify, request
from requests import get
from requests.auth import HTTPBasicAuth


flask_app = Flask(__name__)
celery_app = Celery("tasks", broker="rabbit...")


def get_workers_count():
    inspector = celery_app.control.inspect()
    nodes_stats = inspector.stats()
    nodes_reserved = inspector.reserved()

    workers = 0
    for stats in nodes_stats.values():
        workers += stats["pool"]["max-concurrency"]

    return workers


WORKERS_COUNT = get_workers_count()


@flask_app.route("/")
def home():
    async_ = request.args.get("async")
    settings = request.args.get("settings")

    if async_:
        if not accepting_new_tasks(celery_app):
            return abort(503)

        task = celery_app.send_task(name="my-task", kwargs={"settings": settings})
        return jsonify({"taskId": task.id})

    return jsonify({})


def accepting_new_tasks(celery_app):WORKERS_COUNT
    auth = HTTPBasicAuth("guest", "guest")
    response = get(
        "http://localhost:15672/api/queues/my_vhost/celery",
         auth=auth
    )
    waiting_tasks = response.json()["messages"]
    return waiting_tasks < math.ceil(WORKERS_COUNT / 3)

【讨论】:

    猜你喜欢
    • 2021-05-30
    • 1970-01-01
    • 2020-01-04
    • 1970-01-01
    • 2016-01-28
    • 2018-04-22
    • 2019-09-13
    • 2020-02-01
    • 1970-01-01
    相关资源
    最近更新 更多