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

使用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())

整体处理流程

  1. 使用pydub加载音频文件(MP3直接加载,MP4先通过ffmpeg转成MP3再加载)
  2. 将音频分割为带500ms重叠的片段(每个片段约29秒)
  3. 并行处理每个片段:
    • 导出片段到临时MP3文件
    • 读取文件内容到内存
    • 上传内容到转写API
    • 删除临时MP3文件
  4. 将所有片段的转写结果合并为完整文本

我目前用asyncio.to_thread()来避免阻塞事件循环,也严格按照步骤处理临时文件,但警告依然存在。看起来警告可能来自pydub的chunk.export()方法?我实在搞不懂代码里哪里还有未关闭的文件,难道是pydub内部对AudioSegment的处理有隐藏的文件句柄没释放吗?

备注:内容来源于stack exchange,提问作者codeing_monkey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 18:59:38