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
相关产品推荐
相关产品推荐

