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

Python批量异步上传音频文件并获取特征结果的问题

批量异步处理音频文件上传与特征查询解决方案

问题背景

本地存储数千个音频文件,需通过FastAPI搭建的API提取特征,流程为:

  1. 上传音频文件获取token
  2. 通过token轮询查询处理状态,获取最终特征结果

已实现单文件上传和查询,但批量异步上传时返回422 Unprocessable Entity错误,需解决批量异步/并行处理的完整流程。

错误原因分析

原异步上传代码存在两个核心问题:

  1. 请求参数格式错误:使用data={'data': data}嵌套结构,不符合API预期的files={'file': 文件对象}格式
  2. 未正确读取异步文件流:定义了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ü

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 21:32:36