使用asyncio异步处理音频转写时遭遇未关闭文件的ResourceWarning问题求助
asyncio异步处理音频转写时遭遇未关闭文件的ResourceWarning问题求助
我正在开发一个Python服务,用来处理MP3/MP4音频文件——把它们分割成片段,再通过REST API异步转写。代码本身能正常运行,但日志里一直弹出一堆ResourceWarnings,提示有未关闭的文件:
ResourceWarning: unclosed file <_io.BufferedRandom name='/var/folders/vw/k_0mpg690rg44m0mw3vj9l2c0000gp/T/tmpoico59vl.mp3'>
handle = None
ResourceWarning: Enable tracemalloc to get the object allocation traceback
我的代码是并行处理多个音频片段的,虽然已经显式关闭并删除了临时文件,但这些警告还是会反复出现。下面是我的相关代码和整体流程,实在找不到问题出在哪,希望能得到大家的帮助。
相关代码
import logging import aiohttp import asyncio import json import os import tempfile import sys import subprocess from pydub import AudioSegment logger = logging.getLogger(__name__) MODEL_ENDPOINT = ( "http://example" ) MODEL_NAME = "example" def load_audio_file(file_path): """Load audio file from disk.""" logger.info(f"Loading audio file from {file_path}") file_extension = os.path.splitext(file_path)[1][1:].lower() if file_extension == 'mp3': audio = AudioSegment.from_mp3(file_path) elif file_extension == 'mp4': temp_mp3 = tempfile.NamedTemporaryFile(delete=False, suffix='.mp3') temp_mp3.close() try: cmd = [ 'ffmpeg', '-i', file_path, '-vn', '-ar', '44100', '-ac', '2', '-b:a', '192k', '-f', 'mp3', '-y', temp_mp3.name ] logger.info(f"MP4 file detected, converting to mp3") subprocess.run( cmd, check=True, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) audio = AudioSegment.from_mp3(temp_mp3.name) finally: if os.path.exists(temp_mp3.name): os.unlink(temp_mp3.name) else: raise Exception("Only mp3 and mp4 files supported") logger.info(f"Loaded audio file: {len(audio)}ms duration") return audio def split_audio(audio, chunk_duration_ms=29000): chunks = [] total_duration_ms = len(audio) overlap_ms = 500 for start_ms in range(0, total_duration_ms, chunk_duration_ms - overlap_ms): end_ms = min(start_ms + chunk_duration_ms, total_duration_ms) chunk = audio[start_ms:end_ms] chunks.append(chunk) logger.info(f"Split audio into {len(chunks)} chunks") return chunks async def export_audio_chunk(chunk, chunk_index): temp_file = tempfile.NamedTemporaryFile(delete=False, suffix='.mp3') temp_file.close() # run the export in a thread to avoid blocking the event loop await asyncio.to_thread(chunk.export, temp_file.name, format="mp3") def read_file(filename): with open(filename, 'rb') as f: return f.read() audio_data = await asyncio.to_thread(read_file, temp_file.name) await asyncio.to_thread(os.unlink, temp_file.name) return audio_data async def transcribe_chunk_async(session, chunk, chunk_index, api_url): try: logger.info(f"Preparing chunk {chunk_index+1} for API...") audio_data = await export_audio_chunk(chunk, chunk_index) form_data = aiohttp.FormData() form_data.add_field('file', audio_data, filename=f'chunk_{chunk_index}.mp3', content_type='audio/mpeg') form_data.add_field('model', MODEL_NAME) form_data.add_field('response_format', 'json') form_data.add_field('temperature', '0.0') logger.info( f"Sending chunk {chunk_index+1} ({len(audio_data)/1024:.1f} KB) to API...") async with session.post(f"{api_url}/audio/transcriptions", data=form_data) as response: if response.status != 200: logger.error( f"Error for chunk {chunk_index+1}: HTTP {response.status}") text = await response.text() logger.error(f"Response: {text[:500]}") return chunk_index, None try: result = await response.json() if 'text' in result: logger.info( f"Chunk {chunk_index+1} transcribed: {result['text'][:50]}...") return chunk_index, result except json.JSONDecodeError: logger.error( f"Failed to decode JSON response for chunk {chunk_index+1}") text = await response.text() logger.error(f"Response content: {text[:500]}...") return chunk_index, None except Exception as e: logger.error( f"Exception during transcription of chunk {chunk_index+1}: {e}") return chunk_index, None async def transcribe_audio_async(file_path): audio = await asyncio.to_thread(load_audio_file, file_path) chunks = await asyncio.to_thread(split_audio, audio) connector = aiohttp.TCPConnector(limit=10) timeout = aiohttp.ClientTimeout(total=600) async with aiohttp.ClientSession(connector=connector, timeout=timeout) as session: tasks = [ transcribe_chunk_async(session, chunk, i, MODEL_ENDPOINT) for i, chunk in enumerate(chunks) ] results = await asyncio.gather(*tasks) ordered_results = sorted(results, key=lambda x: x[0]) transcriptions = [r[1]['text'] if r[1] and 'text' in r[1] else "" for r in ordered_results] full_transcription = " ".join(transcriptions) logger.info( f"Transcription complete: {len(full_transcription)} characters") return full_transcription async def transcribe_audio_async_wrapper(file): try: logger.info( f"Creating temporary file for uploaded file {file.filename}") temp_file = tempfile.NamedTemporaryFile( delete=False, suffix=os.path.splitext(file.filename)[1]) file.save(temp_file.name) temp_file.close() transcript = await transcribe_audio_async(temp_file.name) await asyncio.to_thread(os.unlink, temp_file.name) return transcript except Exception as e: logger.error(f"Transcription failed: {e}") try: if 'temp_file' in locals() and os.path.exists(temp_file.name): await asyncio.to_thread(os.unlink, temp_file.name) except: pass raise def transcribe_audios(files): async def transcribe_all(): return await asyncio.gather(*[transcribe_audio_async_wrapper(file) for file in files]) return asyncio.run(transcribe_all())
整体处理流程
- 使用pydub加载音频文件(MP3直接加载,MP4先通过ffmpeg转成MP3再加载)
- 将音频分割为带500ms重叠的片段(每个片段约29秒)
- 并行处理每个片段:
- 导出片段到临时MP3文件
- 读取文件内容到内存
- 上传内容到转写API
- 删除临时MP3文件
- 将所有片段的转写结果合并为完整文本
我目前用asyncio.to_thread()来避免阻塞事件循环,也严格按照步骤处理临时文件,但警告依然存在。看起来警告可能来自pydub的chunk.export()方法?我实在搞不懂代码里哪里还有未关闭的文件,难道是pydub内部对AudioSegment的处理有隐藏的文件句柄没释放吗?
备注:内容来源于stack exchange,提问作者codeing_monkey
相关产品推荐
相关产品推荐

