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

基于Python aiortc的WebRTC实时音频流传输异常问题排查

使用aiortc实现WebRTC实时音频流传输的问题

我正尝试用Python的aiortc库基于WebRTC实现客户端间的实时音频流传输,架构设计如下:

  • 发送端:使用继承自MediaStreamTrack的自定义类MicrophoneAudioTrack,其recv()方法返回从麦克风捕获的音频字节数据。
  • 接收端:预期通过AudioReceiver类处理传入轨道,该类会调用轨道的recv()方法。

目前遇到的问题:recv()方法在发送端本地运行完全正常,但将轨道发送至其他客户端后,接收端无法正确处理轨道数据。

代码实现

import asyncio
import json
import websockets
import sounddevice as sd
import numpy as np
from aiortc import RTCPeerConnection, RTCSessionDescription, MediaStreamTrack, RTCIceServer, RTCConfiguration

class MicrophoneAudioTrack(MediaStreamTrack):
    kind = "audio"

    def __init__(self):
        super().__init__()  # 初始化父类MediaStreamTrack
        self.stream = sd.InputStream(
            channels=1,  # 单声道
            samplerate=48000,  # 采样率
            dtype=np.int16,  # 数据类型
            callback=self.audio_callback
        )
        self.audio_queue = asyncio.Queue()
        self.stream.start()

    def audio_callback(self, indata, frames, time, status):
        if status:
            print(status)
        self.audio_queue.put_nowait(indata.copy())  # 将音频数据存入队列

    async def recv(self):
        try:
            print('called recv')
            audio_data = await self.audio_queue.get()  # 等待队列中的数据
            return audio_data.tobytes()  # 以字节格式返回数据
        except asyncio.CancelledError:
            print("接收数据已取消。")

class AudioReceiver:
    def __init__(self):
        self.track = None
        self.output_stream = sd.OutputStream(
            channels=1,  # 单声道
            samplerate=48000,  # 采样率
            dtype=np.int16  # 数据类型
        )
        self.output_stream.start()

    async def handle_track(self, track):
        print("进入handle_track方法")
        self.track = track
        while True:
            try:
                audio_data = await track.recv()  # 等待接收帧数据
                # 处理接收到的音频数据
                print("已接收音频数据。")
                audio_array = np.frombuffer(audio_data, dtype=np.int16)
                self.output_stream.write(audio_array)  # 播放音频
            except Exception as e:
                print(f"发生错误: {e}")

async def signaling(ws, pc):
    """WebSocket信令处理函数"""
    async for message in ws:
        data = json.loads(message)

        if "sdp" in data:
            desc = RTCSessionDescription(sdp=data["sdp"], type=data["type"])

            if desc.type == "offer":
                # 响应SDP Offer
                await pc.setRemoteDescription(desc)
                answer = await pc.createAnswer()
                await pc.setLocalDescription(answer)

                await ws.send(json.dumps({
                    "sdp": pc.localDescription.sdp,
                    "type": pc.localDescription.type
                }))
            elif desc.type == "answer":
                # 处理SDP Answer
                await pc.setRemoteDescription(desc)

        elif "candidate" in data:
            candidate = data.get("candidate")
            if candidate:
                ice_candidate = {
                    "candidate": candidate,
                    "sdpMid": data["sdpMid"],
                    "sdpMLineIndex": data["sdpMLineIndex"]
                }
                await pc.addIceCandidate(ice_candidate)

async def run_client():
    # 连接信令服务器
    ws = await websockets.connect("ws://localhost:8000/ws")

    # 创建带ICE服务器的RTCPeerConnection(STUN/TURN)
    ice_servers = [
        RTCIceServer(urls="stun:stun.l.google.com:19302")
    ]
    configuration = RTCConfiguration(iceServers=ice_servers)
    pc = RTCPeerConnection(configuration=configuration)

    # 创建麦克风音频轨道对象并添加到RTCPeerConnection
    microphone_track = MicrophoneAudioTrack()
    pc.addTrack(microphone_track)

    # 跟踪ICE连接状态的标志
    ice_connection_ready = asyncio.Event()

    # 创建用于轨道播放的对象
    audio_receiver = AudioReceiver()

    @pc.on("track")
    async def on_track(track):
        if isinstance(track, MediaStreamTrack):
            await ice_connection_ready.wait()
            print(f"正在接收{track.kind}轨道")
            asyncio.ensure_future(audio_receiver.handle_track(track))

    @pc.on("icecandidate")
    async def on_icecandidate(candidate):
        if candidate:
            await ws.send(json.dumps({
                "candidate": candidate.candidate,
                "sdpMid": candidate.sdpMid,
                "sdpMLineIndex": candidate.sdpMLineIndex
            }))

    @pc.on("iceconnectionstatechange")
    async def on_iceconnectionstatechange():
        print(f"ICE连接状态已变更为{pc.iceConnectionState}")
        if pc.iceConnectionState in ("connected", "completed"):
            ice_connection_ready.set()  # 连接就绪时设置标志

    # 检查服务器上是否已有保存的SDP Offer
    await ws.send(json.dumps({"request_offer": True}))
    response = await ws.recv()
    data = json.loads(response)

    if "sdp" in data and data["type"] == "offer":
        # 接收SDP Offer并创建Answer
        desc = RTCSessionDescription(sdp=data["sdp"], type="offer")
        await pc.setRemoteDescription(desc)
        answer = await pc.createAnswer()
        await pc.setLocalDescription(answer)

        # 发送SDP Answer到服务器
        await ws.send(json.dumps({
            "sdp": pc.localDescription.sdp,
            "type": pc.localDescription.type
        }))
    else:
        # 创建SDP Offer初始化会话
        offer = await pc.createOffer()
        await pc.setLocalDescription(offer)
        await ws.send(json.dumps({
            "sdp": pc.localDescription.sdp,
            "type": "offer"
        }))

    # 继续交换ICE候选和信令
    await signaling(ws, pc)

    # 完成后关闭连接
    await pc.close()

# 启动主事件循环
asyncio.run(run_client())

已确认发送端的recv()方法可正常运行,恳请各位提供解决建议或指导!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 06:07:32