如何通过HTTP或MQTT调用启停Python音频流脚本?
实现HTTP/MQTT控制音频流启停的方案
现有一个基于websockets、asyncio和pyaudio的Python音频流脚本,目前需手动执行
python script.py启动,按Ctrl+C停止。希望实现通过HTTP或MQTT调用启停该脚本或音频流,请问可行的实现方案是什么?
原脚本代码:
# starts recording stream = p.open( format=FORMAT, channels=CHANNELS, rate=RATE, input=True, frames_per_buffer=FRAMES_PER_BUFFER ) async def send_receive(): print(f'Connecting websocket to url ${URL}') async with websockets.connect( URL, extra_headers=(("Authorization", auth_key),), ping_interval=5, ping_timeout=20 ) as _ws: await asyncio.sleep(0.1) print("Receiving SessionBegins ...") session_begins = await _ws.recv() print(session_begins) print("Sending messages ...") async def send(): while True: try: data = stream.read(FRAMES_PER_BUFFER) data = base64.b64encode(data).decode("utf-8") json_data = json.dumps({"audio_data":str(data)}) await _ws.send(json_data) except websockets.exceptions.ConnectionClosedError as e: print(e) assert e.code == 4008 break except Exception as e: assert False, "Not a websocket 4008 error" await asyncio.sleep(0.01) return True async def receive(): while True: try: result_str = await _ws.recv() print(json.loads(result_str)['text']) except websockets.exceptions.ConnectionClosedError as e: print(e) assert e.code == 4008 break except Exception as e: assert False, "Not a websocket 4008 error" send_result, receive_result = await asyncio.gather(send(), receive()) asyncio.run(send_receive())
方案一:HTTP接口控制(基于FastAPI)
核心思路
- 用FastAPI搭建轻量HTTP服务,提供
/start和/stop两个控制接口 - 维护全局状态变量管理音频流任务的运行状态
- 将原有的音频流+websocket逻辑封装成可异步启动/终止的独立任务
修改后完整代码
import asyncio import websockets import base64 import json import pyaudio from fastapi import FastAPI from fastapi.responses import JSONResponse # 配置参数(根据实际场景修改) FORMAT = pyaudio.paInt16 CHANNELS = 1 RATE = 16000 FRAMES_PER_BUFFER = 3200 WS_URL = "ws://your-websocket-server-url" AUTH_KEY = "your-auth-token" # 全局状态与资源管理 stream = None audio_task = None is_running = False p = pyaudio.PyAudio() app = FastAPI() async def audio_stream_task(): global stream, is_running try: # 初始化录音流 stream = p.open( format=FORMAT, channels=CHANNELS, rate=RATE, input=True, frames_per_buffer=FRAMES_PER_BUFFER ) print(f'Connecting websocket to url {WS_URL}') async with websockets.connect( WS_URL, extra_headers=(("Authorization", AUTH_KEY),), ping_interval=5, ping_timeout=20 ) as _ws: await asyncio.sleep(0.1) print("Receiving SessionBegins ...") session_begins = await _ws.recv() print(session_begins) print("Sending messages ...") async def send_audio(): while is_running: try: data = stream.read(FRAMES_PER_BUFFER) data = base64.b64encode(data).decode("utf-8") json_data = json.dumps({"audio_data": str(data)}) await _ws.send(json_data) except websockets.exceptions.ConnectionClosedError as e: print(e) break except Exception as e: print(f"Send error: {str(e)}") break await asyncio.sleep(0.01) return True async def receive_result(): while is_running: try: result_str = await _ws.recv() print(json.loads(result_str)['text']) except websockets.exceptions.ConnectionClosedError as e: print(e) break except Exception as e: print(f"Receive error: {str(e)}") break await asyncio.gather(send_audio(), receive_result()) finally: # 清理资源 if stream is not None: stream.stop_stream() stream.close() is_running = False print("Audio stream stopped") @app.post("/start") async def start_stream(): global audio_task, is_running if is_running: return JSONResponse(content={"status": "error", "message": "Stream already running"}, status_code=400) is_running = True audio_task = asyncio.create_task(audio_stream_task()) return JSONResponse(content={"status": "success", "message": "Stream started"}) @app.post("/stop") async def stop_stream(): global audio_task, is_running if not is_running: return JSONResponse(content={"status": "error", "message": "Stream not running"}, status_code=400) is_running = False if audio_task is not None: await audio_task return JSONResponse(content={"status": "success", "message": "Stream stopped"}) if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000)
使用说明
- 安装依赖:
pip install fastapi uvicorn websockets pyaudio - 启动服务:
python script.py - 调用控制接口:
- 启动流:
curl -X POST http://localhost:8000/start - 停止流:
curl -X POST http://localhost:8000/stop
- 启动流:
方案二:MQTT消息控制(基于paho-mqtt)
核心思路
- 用paho-mqtt客户端订阅指定控制主题(如
audio/control) - 接收
start/stop指令来触发音频流任务的启停 - 通过全局状态变量维护任务生命周期,确保资源正确释放
修改后完整代码
import asyncio import websockets import base64 import json import pyaudio import paho.mqtt.client as mqtt # 配置参数(根据实际场景修改) FORMAT = pyaudio.paInt16 CHANNELS = 1 RATE = 16000 FRAMES_PER_BUFFER = 3200 WS_URL = "ws://your-websocket-server-url" AUTH_KEY = "your-auth-token" MQTT_BROKER = "your-mqtt-broker-address" MQTT_PORT = 1883 CONTROL_TOPIC = "audio/control" # 全局状态与资源管理 stream = None audio_task = None is_running = False p = pyaudio.PyAudio() mqtt_client = mqtt.Client(client_id="audio-stream-controller") async def audio_stream_task(): global stream, is_running try: stream = p.open( format=FORMAT, channels=CHANNELS, rate=RATE, input=True, frames_per_buffer=FRAMES_PER_BUFFER ) print(f'Connecting websocket to url {WS_URL}') async with websockets.connect( WS_URL, extra_headers=(("Authorization", AUTH_KEY),), ping_interval=5, ping_timeout=20 ) as _ws: await asyncio.sleep(0.1) print("Receiving SessionBegins ...") session_begins = await _ws.recv() print(session_begins) print("Sending messages ...") async def send_audio(): while is_running: try: data = stream.read(FRAMES_PER_BUFFER) data = base64.b64encode(data).decode("utf-8") json_data = json.dumps({"audio_data": str(data)}) await _ws.send(json_data) except websockets.exceptions.ConnectionClosedError as e: print(e) break except Exception as e: print(f"Send error: {str(e)}") break await asyncio.sleep(0.01) return True async def receive_result(): while is_running: try: result_str = await _ws.recv() print(json.loads(result_str)['text']) except websockets.exceptions.ConnectionClosedError as e: print(e) break except Exception as e: print(f"Receive error: {str(e)}") break await asyncio.gather(send_audio(), receive_result()) finally: if stream is not None: stream.stop_stream() stream.close() is_running = False print("Audio stream stopped") def on_mqtt_message(client, userdata, msg): global audio_task, is_running payload = msg.payload.decode("utf-8").strip().lower() if payload == "start": if not is_running: is_running = True audio_task = asyncio.create_task(audio_stream_task()) print("Received start command, starting stream") else: print("Stream already running, ignore start command") elif payload == "stop": if is_running: is_running = False print("Received stop command, stopping stream") else: print("Stream not running, ignore stop command") async def mqtt_loop(): mqtt_client.on_message = on_mqtt_message mqtt_client.connect(MQTT_BROKER, MQTT_PORT, 60) mqtt_client.subscribe(CONTROL_TOPIC) while True: mqtt_client.loop(timeout=1.0) await asyncio.sleep(0.1) if __name__ == "__main__": loop = asyncio.get_event_loop() try: loop.run_until_complete(asyncio.gather(mqtt_loop())) except KeyboardInterrupt: if is_running: is_running = False if audio_task is not None: loop.run_until_complete(audio_task) mqtt_client.disconnect() p.terminate()
使用说明
- 安装依赖:
pip install paho-mqtt websockets pyaudio - 启动脚本:
python script.py - 发送控制消息:
- 启动流:向主题
audio/control发送消息start - 停止流:向主题
audio/control发送消息stop
- 启动流:向主题
关键注意事项
- 全局状态变量在单线程asyncio环境下无需加锁,多线程场景需添加线程锁保证安全
- 停止任务时通过
is_running标志位让循环自然退出,避免强制终止协程导致资源泄漏 - 必须确保pyaudio流在任务结束后正确关闭,防止音频设备被占用
- 生产环境部署时,HTTP方案可配合Nginx做反向代理,MQTT方案需保证broker的稳定性
内容的提问来源于stack exchange,提问作者Abhishek Kumar
相关产品推荐
相关产品推荐

