【发布时间】: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