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

Flask应用如何通过任务队列限制异步处理的最大并发线程数

可行实现方案

方案1:使用Python内置ThreadPoolExecutor(无额外依赖,适配当前场景)

这个方案改动最小,仅需新增3行左右代码即可实现固定线程数限制,完全匹配你当前的需求:

  • 核心逻辑:全局初始化固定大小的线程池,所有音频处理任务提交到线程池排队执行,同时运行的线程数不会超过你设置的最大值,高并发场景下任务自动排队不会爆线程。

修改后的代码

首先导入依赖,全局初始化线程池:

from concurrent.futures import ThreadPoolExecutor
import boto3
import os
import json
from flask import Flask, request
# 其他你原有依赖保持不变

# 配置最大同时运行线程数,可根据服务器性能调整为5/8等值
MAX_WORKER_NUM = 5
# 全局唯一线程池实例,不要在请求处理函数内初始化
executor = ThreadPoolExecutor(max_workers=MAX_WORKER_NUM)

# 你原有的process_audio函数完全不用改
def process_audio(bucket_name, key, _id, extension):
    S3_CLIENT = boto3.client('s3', region_name=S3_REGION)
    print('Running audio proccessing')
    INPUT_FILE = os.path.join(TEMP_PATH, f'{_id}.{extension}')
    print(f'Saving downloaded file to {INPUT_FILE}')
    S3_CLIENT.download_file(bucket_name, key, INPUT_FILE)
    print('File downloaded')
    process = stt.process_audio(INPUT_FILE)
    print(f'Audio processed by AI returned: "{process}"')
    stt.reset()
    ai = get_sentimentAI_results(process)
    if ai:
        print(f'Text processed by AI returned class {ai[0]} with a certainty of {ai[1]}%')
        return True

    print('Request to sentiment AI endpoint failed for an unkown reason. Check CloudWhatch for more information!')
    return False

# 路由函数仅需修改线程启动部分
@app.route('/process/audio', methods=['POST'])
def process_new_audio():
    print('Recieving new request')
    data = request.data

    if not data:
        return '', 404

    data = json.loads(data)
    bucket_name = data.get('bucket_name')
    key = data.get('key')
    _id = data.get('id')
    extension = data.get('file_extension')

    if not key or not bucket_name or not _id or not extension:
        return '', 404

    # 替换原有线程启动逻辑
    executor.submit(process_audio, bucket_name=bucket_name, key=key, _id=_id, extension=extension)
    
    return '', 200

可选优化

如果需要避免任务队列积压过多导致内存溢出,可以在提交任务前判断队列长度,超出阈值时直接返回限流响应:

# 提交任务前添加判断
if executor._work_queue.qsize() > 100: # 最大允许100个任务排队
    return 'Too many requests, please retry later', 429
executor.submit(...)

方案2:使用Celery任务队列(适用于后续分布式/高级功能扩展)

如果后续需要任务持久化、失败自动重试、多实例分布式处理等功能,可以改用Celery配合Redis/RabbitMQ做任务队列:

  • 安装依赖:pip install celery redis
  • 初始化Celery实例,将process_audio注册为Celery任务
  • 路由中直接提交任务到Celery队列即可,并发数通过Celery worker的concurrency参数控制

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 08:24:01