基于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
相关产品推荐
相关产品推荐

