Flask+React Native应用中如何获取用户输入并优化状态更新?
解决方案:异步任务交互 + 实时状态推送
Hi there! 针对你提到的两个核心需求——在stuff.main()执行过程中获取用户输入,以及替换轮询提升性能——我整理了一套实用的实现方案,结合异步任务队列和实时推送技术来解决问题:
一、在stuff.main()执行中获取用户输入
当前代码里stuff.main()是同步调用的,这会阻塞请求且无法中途接收用户输入。我们需要把它改成可中断的异步任务,通过中间存储(比如Redis)来传递用户输入:
具体实现步骤
- 引入异步任务队列:用Celery + Redis管理后台任务,让
/start接口立即返回任务ID,而非等待任务完成。 - 改造
stuff.main():拆分成多步骤执行,在需要用户输入的节点暂停,从Redis读取输入后继续。 - 新增输入接收接口:客户端提交用户输入时,将输入与任务ID关联存储到Redis。
服务器端代码示例
from celery import Celery import redis import time from flask import request, jsonify, Response from flask_cors import cross_origin # 初始化Redis和Celery redis_client = redis.Redis(host='localhost', port=6379, db=0) celery = Celery('tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0') # 异步任务:可接收用户输入的main逻辑 @celery.task(bind=True) def run_main_task(self, task_id, url, username, password, stuff1, stuff2): # 第一步:执行无需输入的逻辑 step1_result = stuff.step1(url, username, password) self.update_state(state='RUNNING', meta={'message': '执行第一步完成'}) # 等待用户输入:阻塞读取Redis(比轮询更高效) self.update_state(state='WAITING', meta={'message': '请输入必要信息'}) user_input = redis_client.blpop(f"task:{task_id}:input", timeout=300)[1].decode('utf-8') # 第二步:使用用户输入继续执行 step2_result = stuff.step2(stuff1, stuff2, user_input) self.update_state(state='SUCCESS', meta={'message': '任务完成', 'result': step2_result}) return step2_result # 修改/start接口:提交异步任务并返回任务ID @app.route('/start', methods=['POST']) @cross_origin(origin='*',headers=['Content-Type','Authorization']) def start_main(): data = request.get_json() task_id = data['session'] # 用sessionID作为任务唯一标识 run_main_task.delay(task_id, data['url'], data['user'], data['password'], data['stuff1'], data['stuff2']) return jsonify({'task_id': task_id}) # 新增接收用户输入的接口 @app.route('/submit-input', methods=['POST']) @cross_origin(origin='*',headers=['Content-Type','Authorization']) def submit_input(): data = request.get_json() task_id = data['session'] user_input = data['input'] redis_client.rpush(f"task:{task_id}:input", user_input) return jsonify({'status': '输入已接收'})
客户端调整逻辑
当监听到任务进入WAITING状态时,显示输入界面并提交:
// (配合后续的实时推送逻辑)在状态更新回调中处理等待输入的情况 if(data.state === 'WAITING'){ const userInput = prompt(data.meta.message); if(userInput){ fetch('http://localhost:5000/submit-input', { method: 'POST', headers: {'Content-Type': 'application/json'}, body: JSON.stringify({session: sessionID, input: userInput}) }); } }
二、替换轮询,实现高效实时状态推送
轮询在用户量增大时会导致服务器请求爆炸,推荐用Server-Sent Events (SSE) 或WebSocket实现服务器主动推送,前者适合单向状态更新,后者支持双向交互:
方案1:Server-Sent Events (SSE) (轻量单向推送)
SSE基于HTTP协议,只需一次连接即可持续接收服务器推送,实现简单:
服务器端代码
@app.route('/status-sse/<task_id>') @cross_origin(origin='*') def status_sse(task_id): def generate_status(): task = run_main_task.AsyncResult(task_id) while True: if task.state == 'PENDING': status_data = {'status': '任务初始化中...', 'state': 'PENDING'} elif task.state == 'WAITING': status_data = {'status': task.meta['message'], 'state': 'WAITING'} elif task.state == 'SUCCESS': status_data = {'status': f"任务完成:{task.result}", 'state': 'SUCCESS'} yield f"data: {json.dumps(status_data)}\n\n" break elif task.state == 'FAILURE': status_data = {'status': '任务失败', 'state': 'FAILURE'} yield f"data: {json.dumps(status_data)}\n\n" break else: status_data = {'status': f"执行中:{task.meta['message']}", 'state': 'RUNNING'} yield f"data: {json.dumps(status_data)}\n\n" time.sleep(1) # 控制推送频率 return Response(generate_status(), mimetype='text/event-stream')
客户端代码
React.useEffect(() => { const eventSource = new EventSource(`http://localhost:5000/status-sse/${sessionID}`); eventSource.onmessage = (event) => { const data = JSON.parse(event.data); setStatus(data.status); if(data.state === 'FAILURE'){ setStatus("Sorry, this task has failed"); setTimeout(()=>{navigation.goBack()}, 3000); eventSource.close(); } else if(data.state === 'SUCCESS'){ setTimeout(()=>{navigation.goBack()}, 3000); eventSource.close(); } }; eventSource.onerror = () => { setStatus("Either our server or this device went offline"); setTimeout(()=>{ setStatus("Try again..."); navigation.goBack('Task_Elaborate'); eventSource.close(); }, 5000); }; return () => eventSource.close(); }, [sessionID]);
方案2:WebSocket (双向实时交互)
如果需要更复杂的双向通信(比如客户端主动发送指令),可以用Flask-SocketIO:
服务器端代码
from flask_socketio import SocketIO, emit, join_room socketio = SocketIO(app, cors_allowed_origins="*") @socketio.on('join_task') def handle_join_task(data): task_id = data['session'] join_room(task_id) # 加入任务专属房间,定向推送状态 # 后台监听任务状态 def monitor_task(): task = run_main_task.AsyncResult(task_id) while task.state not in ['SUCCESS', 'FAILURE']: emit('status_update', { 'status': task.meta.get('message', '执行中'), 'state': task.state }, room=task_id) time.sleep(1) # 推送最终状态 if task.state == 'SUCCESS': emit('status_update', {'status': '任务完成', 'state': 'SUCCESS'}, room=task_id) else: emit('status_update', {'status': '任务失败', 'state': 'FAILURE'}, room=task_id) socketio.start_background_task(monitor_task) # 接收用户输入的WebSocket事件 @socketio.on('submit_input') def handle_submit_input(data): task_id = data['session'] user_input = data['input'] redis_client.rpush(f"task:{task_id}:input", user_input) emit('input_received', {'status': '输入已接收'}, room=task_id)
客户端代码
import io from 'socket.io-client'; React.useEffect(() => { const socket = io('http://localhost:5000'); socket.emit('join_task', {session: sessionID}); socket.on('status_update', (data) => { setStatus(data.status); if(data.state === 'FAILURE'){ setStatus("Sorry, this task has failed"); setTimeout(()=>{navigation.goBack()}, 3000); socket.disconnect(); } else if(data.state === 'SUCCESS'){ setTimeout(()=>{navigation.goBack()}, 3000); socket.disconnect(); } else if(data.state === 'WAITING'){ const userInput = prompt(data.status); if(userInput){ socket.emit('submit_input', {session: sessionID, input: userInput}); } } }); socket.on('connect_error', () => { setStatus("Either our server or this device went offline"); setTimeout(()=>{ setStatus("Try again..."); navigation.goBack('Task_Elaborate'); socket.disconnect(); }, 5000); }); return () => socket.disconnect(); }, [sessionID]);
总结
- 用Celery+Redis实现异步任务,让
stuff.main()可以中途接收用户输入,避免请求阻塞。 - 用SSE或WebSocket替换轮询,服务器主动推送状态,大幅降低服务器负载,提升高并发场景下的性能。
内容的提问来源于stack exchange,提问作者Vishal DS
相关产品推荐
相关产品推荐

