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

如何通过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)

使用说明

  1. 安装依赖:pip install fastapi uvicorn websockets pyaudio
  2. 启动服务:python script.py
  3. 调用控制接口:
    • 启动流: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()

使用说明

  1. 安装依赖:pip install paho-mqtt websockets pyaudio
  2. 启动脚本:python script.py
  3. 发送控制消息:
    • 启动流:向主题audio/control发送消息start
    • 停止流:向主题audio/control发送消息stop

关键注意事项

  1. 全局状态变量在单线程asyncio环境下无需加锁,多线程场景需添加线程锁保证安全
  2. 停止任务时通过is_running标志位让循环自然退出,避免强制终止协程导致资源泄漏
  3. 必须确保pyaudio流在任务结束后正确关闭,防止音频设备被占用
  4. 生产环境部署时,HTTP方案可配合Nginx做反向代理,MQTT方案需保证broker的稳定性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 08:07:02