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后端。
- 解决步骤:
- 安装官方推荐的异步服务器:
pip install eventlet - 修改启动代码,指定使用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
相关产品推荐
相关产品推荐

