You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Flask服务中管理请求队列?现有子进程方案存问题

嘿,你的这个Flask请求处理方式确实踩了几个坑——无限制创建子进程很容易把服务器资源耗干,而且完全没做队列控制,高并发下肯定出问题。我给你梳理几个靠谱的解决方案,从简单到生产级都有:

先说说你当前代码的核心问题

  • 每次请求都新建一个Process,请求量上来后系统进程数会爆炸,直接把CPU、内存占满
  • 子进程结束后如果主进程没调用join(),会变成僵尸进程,持续占用系统资源
  • 跨进程传递request_data依赖pickle序列化,如果数据结构复杂(比如包含无法序列化的对象)会直接报错

方案1:用进程池做轻量队列控制(适合小规模场景)

用multiprocessing.Pool或者concurrent.futures.ProcessPoolExecutor,它们会自动维护一个进程队列,控制并发数,还能复用进程减少开销:

from multiprocessing import Pool
from flask import Flask, request

app = Flask(__name__)
# 根据服务器CPU核心数设置进程数,比如4核就设4,避免资源浪费
process_pool = Pool(processes=4)

def do_the_processing(request_data):
    # 这里写你的实际处理逻辑
    print(f"Processing data: {request_data}")
    # 处理完直接返回即可,不需要手动调用sys.exit()

@app.route('/', methods=['POST'])
def index():
    request_data = request.get_json()
    # 把任务提交到进程池,自动排队等待空闲进程处理
    process_pool.apply_async(do_the_processing, args=(request_data,))
    return '200 OK'

if __name__ == '__main__':
    app.run()

注意:如果你的处理任务是IO密集型(比如调用外部API、读写文件),用ThreadPoolExecutor线程池会更高效,代码结构和上面几乎一样,只是把Pool换成ThreadPoolExecutor就行。


方案2:用专业任务队列(适合生产环境大规模请求)

如果你的服务需要处理大量请求,或者需要任务持久化、重试、监控等功能,直接上Celery+消息中间件(比如Redis/RabbitMQ),这是Python生态里最成熟的异步任务方案:

首先安装依赖:

pip install celery redis

然后编写代码:

from flask import Flask, request
from celery import Celery

app = Flask(__name__)
# 配置Redis作为消息中间件和结果存储(如果不需要结果可以只配broker)
app.config['CELERY_BROKER_URL'] = 'redis://localhost:6379/0'
app.config['CELERY_RESULT_BACKEND'] = 'redis://localhost:6379/0'

# 初始化Celery实例
celery = Celery(app.name, broker=app.config['CELERY_BROKER_URL'])
celery.conf.update(app.config)

# 定义异步任务
@celery.task
def do_the_processing(request_data):
    # 你的处理逻辑
    print(f"Processing data with Celery: {request_data}")

@app.route('/', methods=['POST'])
def index():
    request_data = request.get_json()
    # 把任务发送到队列,Celery会自动处理排队、并发、重试
    do_the_processing.delay(request_data)
    return '200 OK'

if __name__ == '__main__':
    app.run()

启动Celery worker:

celery -A your_app_name worker --loglevel=info

这个方案的优势是:

  • 任务会被持久化到消息中间件,即使Flask重启也不会丢失
  • 可以设置任务优先级、超时时间、重试次数
  • 有完善的监控和管理工具,适合生产环境

额外提醒

  • 不管用哪种方案,都要确保request_data能被序列化(pickle或者JSON),如果有特殊对象,最好先转换成字典再传递
  • 如果是CPU密集型任务,不要用线程池(GIL限制),必须用进程池或者Celery的进程worker
  • 生产环境部署Flask时,不要用自带的app.run(),用Gunicorn+Nginx这类生产级服务器

内容的提问来源于stack exchange,提问作者Savage Panda

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 03:35:34