Python批量异步上传音频文件并获取特征结果的问题
批量异步处理音频文件上传与特征查询解决方案
问题背景
本地存储数千个音频文件,需通过FastAPI搭建的API提取特征,流程为:
- 上传音频文件获取token
- 通过token轮询查询处理状态,获取最终特征结果
已实现单文件上传和查询,但批量异步上传时返回422 Unprocessable Entity错误,需解决批量异步/并行处理的完整流程。
错误原因分析
原异步上传代码存在两个核心问题:
- 请求参数格式错误:使用
data={'data': data}嵌套结构,不符合API预期的files={'file': 文件对象}格式 - 未正确读取异步文件流:定义了
file_reader但未实际使用,同步打开文件也不符合异步场景要求
修正后的完整实现代码
1. 异步文件上传模块
import os import asyncio import aiohttp import aiofiles from tqdm import tqdm class FileManager(): def __init__(self, file_name: str): self.name = file_name self.size = os.path.getsize(self.name) self.pbar = None self.token = None # 存储上传后返回的token def __init_pbar(self): self.pbar = tqdm( total=self.size, desc=self.name, unit='B', unit_scale=True, unit_divisor=1024, leave=True) async def file_reader(self): self.__init_pbar() chunk_size = 64 * 1024 async with aiofiles.open(self.name, 'rb') as f: chunk = await f.read(chunk_size) while chunk: self.pbar.update(len(chunk)) # 用实际读取长度更新进度,避免文件不足chunk_size时的误差 yield chunk chunk = await f.read(chunk_size) self.pbar.close() async def upload_file(file: FileManager, upload_url: str, session: aiohttp.ClientSession): try: # 构造符合API要求的multipart/form-data请求 data = aiohttp.FormData() data.add_field('file', file.file_reader(), filename=os.path.basename(file.name), content_type='audio/wav') async with session.post(upload_url, data=data) as resp: if resp.status == 200: result = await resp.json() file.token = result.get('token') return file else: print(f"上传失败 {file.name}: {resp.status} - {await resp.text()}") return file except Exception as e: print(f"上传异常 {file.name}: {str(e)}") return file
2. 异步轮询查询特征模块
async def query_features(file: FileManager, query_url_template: str, session: aiohttp.ClientSession, poll_interval=5): if not file.token: print(f"{file.name} 无有效token,跳过查询") return (file.name, None) query_url = query_url_template.format(token=file.token) while True: try: async with session.get(query_url) as resp: if resp.status == 200: features = await resp.json() # 假设API返回中包含状态字段,比如'status',处理完成时返回特征 if features.get('status', 'completed') == 'completed': return (file.name, features) else: await asyncio.sleep(poll_interval) elif resp.status == 404: print(f"{file.name} token无效: {file.token}") return (file.name, None) else: print(f"查询失败 {file.name}: {resp.status} - {await resp.text()}") await asyncio.sleep(poll_interval) except Exception as e: print(f"查询异常 {file.name}: {str(e)}") await asyncio.sleep(poll_interval)
3. 主流程入口
async def main(file_paths, upload_url, query_url_template): # 初始化文件管理器列表 file_managers = [FileManager(path) for path in file_paths] # 异步批量上传文件 async with aiohttp.ClientSession() as session: print("开始批量上传文件...") uploaded_files = await asyncio.gather(*[upload_file(f, upload_url, session) for f in file_managers]) # 过滤上传成功获取到token的文件 valid_files = [f for f in uploaded_files if f.token] print(f"上传完成,成功获取{len(valid_files)}个文件的token") # 异步批量轮询查询特征 print("开始批量查询特征...") feature_results = await asyncio.gather(*[query_features(f, query_url_template, session) for f in valid_files]) # 整理结果为字典:文件名 -> 特征数据 result_dict = {name: features for name, features in feature_results if features} print(f"特征查询完成,共获取{len(result_dict)}个文件的特征") return result_dict # 示例调用 if __name__ == "__main__": # 配置API地址 UPLOAD_URL = "https://something/submit/" QUERY_URL_TEMPLATE = "https://something/features/{token}" # 示例文件路径列表(实际可替换为批量获取的路径) file_paths = [ '/media/SPEAKER_01/interview-10001.wav', '/media/SPEAKER_00/interview-10001.wav' # 更多文件路径... ] # 运行异步主流程 final_results = asyncio.run(main(file_paths, UPLOAD_URL, QUERY_URL_TEMPLATE)) # 可将结果保存到文件 # import json # with open('audio_features.json', 'w', encoding='utf-8') as f: # json.dump(final_results, f, indent=2)
关键优化点
- 请求格式修正:使用
aiohttp.FormData构造正确的multipart/form-data请求,匹配API预期的file参数 - 异步文件流读取:实际使用
file_reader异步生成文件块,配合进度条显示上传进度 - 错误处理增强:针对上传和查询阶段的不同状态码及异常捕获,提升鲁棒性
- 轮询逻辑优化:基于API状态字段(需根据实际API调整)进行智能轮询,避免无效请求
- 结果整理:自动过滤失败任务,返回清晰的文件名-特征映射字典
内容的提问来源于stack exchange,提问作者Tütü
相关产品推荐
相关产品推荐

