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

Flask Socket.io应用Emit阻塞求助:线程未结束前端无法接收消息

Flask Socket.io流式响应阻塞问题排查与解决

问题现象

Flask Socket.io应用中,AI模型以流式方式生成响应,定期调用emit推送音频数据块。虽然每个emit执行无报错,但前端必须等待所有前序线程完成后,才能一次性收到所有消息。尝试过线程、线程池、concurrent.futures、后台线程、Celery等方案均无效,推测问题出在配置层面。

核心代码

import os
import re
import io
import logging
import soundfile as sf
from threading import Thread
from flask import Flask
from flask_cors import CORS
from flask_socketio import SocketIO

BASE_DIR = os.path.dirname(os.path.abspath(__file__))
# 假设global_models、user_data、save_location、logger等变量已提前定义

app = Flask(__name__, template_folder=os.path.join(BASE_DIR, 'templates'),
            static_folder=os.path.join(BASE_DIR, 'static'))

CORS(app, resources={r"/*": {"origins": "*"}})
app.config['CORS_HEADERS'] = 'Content-Type'
socketio = SocketIO(app, cors_allowed_origins="*", ping_timeout=120)

class CaseChat:
    def __init__(self):
        self.model = global_models.model
        self.spec_generator = global_models.spec_generator
        self.model_gtts = global_models.model_gtts
        self.bytes_audio = []

    def predict(self, **kwargs):
        recipient_id = kwargs['recipient_id']
        history = kwargs.get('history', [])
        self.bytes_audio = []
        try:
            history.append({'role': 'user', 'content': kwargs.get('chunks') + str(kwargs.get("interview_time", ' '))})
            inputs = global_models.tokenizer.apply_chat_template(history, add_generation_prompt=True, tokenize=False)
            inputs = inputs.replace('''\n\nCutting Knowledge Date: December 2023\nToday Date: 26 Jul 2024\n\n''', "")
            inputs = global_models.tokenizer(inputs, return_tensors='pt').to('cuda')
            streamer = TextIteratorStreamer(global_models.tokenizer, skip_prompt=True, skip_special_tokens=True)
            generation_kwargs = dict(inputs, streamer=streamer, max_new_tokens=512)
            thread_generate = Thread(target=global_models.model.generate, kwargs=generation_kwargs)
            thread_generate.start()
            new_text = ""
            for output in streamer:
                if user_data[recipient_id]['interrupt']:
                    break
                if not output.strip("\n"):
                    continue
                if not len(new_text):
                    if new_text.startswith('assistant'):
                        new_text = new_text[9:]
                    elif new_text.startswith('user'):
                        new_text = new_text[4:]
                new_text += output
                if len(new_text.split(' ')) > 5:
                    history = self.response(history, new_text, kwargs.get('s_id'))
                    new_text = ""
            if len(new_text) and not user_data[recipient_id]['interrupt']:
                history = self.response(history, new_text, kwargs.get('s_id'))
            audio_saver(path=f"{save_location}/{recipient_id}/{user_data[recipient_id]['time_duration']}-AI.wav",
                        audio=self.bytes_audio)
            return history
        except Exception as e:
            logging.error(f"Prediction failed: {str(e)}")
            history.pop(-1)
            return history

    def response(self, history, new_text, s_id):
        new_text = re.sub(r'\([^)]*\)', '', new_text)
        self.speak(new_text, s_id)
        if 'assistant' != history[-1]['role']:
            history.append({'role': 'assistant', 'content': new_text})
        else:
            history[-1]['content'] += new_text
        return history

    def speak(self, text, s_id):
        try:
            ai_text = text.split("AI:")[-1].strip()

            self.spec_generator.eval()

            parsed = self.spec_generator.parse(ai_text)
            spectrogram = self.spec_generator.generate_spectrogram(tokens=parsed)
            audio = self.model_gtts.convert_spectrogram_to_audio(spec=spectrogram)
            audio = audio.detach()
            audio_np = audio.cpu().numpy().squeeze()

            with io.BytesIO() as bytesio:
                sf.write(bytesio, audio_np, 22050, format="wav", subtype="PCM_16")
                bytesio.seek(0)
                audio_data = bytesio.read()

                socketio.emit('response', audio_data, to=s_id)
                print("Text streamed outwards: " + text)
                self.bytes_audio.append(audio_data)
                return "Request received and task started", 202
        except Exception as e:
            logger.exception("TTS conversion failed", e)
            return f"Error: {e}", 500

if __name__ == '__main__':
    socketio.run(app, host='0.0.0.0', port=8080, debug=True, use_reloader=False, allow_unsafe_werkzeug=True,
                 log_output=True)

问题根源与解决方案

1. 同步服务器阻塞消息推送

当前使用默认的Werkzeug同步服务器,即使开启线程,其同步特性会阻塞请求上下文,导致emit的消息无法即时推送到前端。必须切换到支持异步的Socket.IO后端。

  • 解决步骤:
    1. 安装官方推荐的异步服务器:
      pip install eventlet
      
    2. 修改启动代码,指定使用eventlet模式:
      if __name__ == '__main__':
          socketio.run(app, host='0.0.0.0', port=8080, debug=False, use_reloader=False, 
                       allow_unsafe_werkzeug=True, log_output=True, async_mode='eventlet')
      

2. 请求线程被流式循环阻塞

predict方法在Socket.IO的请求线程中执行,主线程持续遍历streamer,阻塞了Socket.IO的消息处理循环,导致emit的消息无法被即时发送。

  • 解决步骤:
    将流式处理逻辑放到Socket.IO的后台任务中,释放请求线程:
    @socketio.on('start_chat')  # 假设触发AI生成的Socket事件是start_chat
    def handle_start_chat(**kwargs):
        # 使用start_background_task自动处理上下文传递
        socketio.start_background_task(target=CaseChat().predict, kwargs=kwargs)
    

3. 调试模式干扰异步行为

调试模式(debug=True)可能会引发异步逻辑异常,建议先关闭调试模式测试,确认问题是否消失。

4. 额外验证点

  • 确保speak方法中to=s_id的s_id是正确的房间/客户端ID,避免消息推送目标错误;
  • 检查前端Socket.IO客户端是否正确监听response事件,没有遗漏接收逻辑。

内容的提问来源于stack exchange,提问作者Aleef

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 01:41:00