如何在aiortc中通过数据通道主动发送消息?
解决aiortc全局主动发送消息的线程与asyncio问题
错误原因分析
你的问题核心是两个点:
- aiortc的
dc.send()是异步方法,不能在普通线程(比如你的Thread-4)里直接同步调用,因为该线程没有绑定asyncio事件循环,导致抛出RuntimeError。 - 给
sendMessage加了async/await后,在普通线程里直接调用协函数却没有用asyncio的方式调度,所以出现“coroutine was never awaited”警告。
解决方案:线程安全的全局发送函数
下面是适配aiortc的全局发送函数实现,支持在async函数或普通线程(比如ROS回调线程)里调用:
import asyncio from aiortc import RTCPeerConnection # 全局变量:保存数据通道和aiortc所在的事件循环 dc = None event_loop = None def send_message(msg): """全局同步发送函数,可在任意线程调用""" if dc is not None and event_loop is not None: # 将异步发送任务提交到aiortc的事件循环中执行 asyncio.run_coroutine_threadsafe(_async_send(msg), event_loop) async def _async_send(msg): """实际执行异步发送的内部函数""" if dc is not None and dc.readyState == "open": await dc.send(msg) async def offer(request): global dc, event_loop # 保存当前运行的asyncio事件循环(必须在async函数内调用) event_loop = asyncio.get_running_loop() pc = RTCPeerConnection() ... @pc.on("datachannel") def on_datachannel(channel): global dc dc = channel # 监听通道关闭,及时重置dc @channel.on("close") def on_channel_close(): global dc dc = None # 监听连接状态变化,连接失效时重置dc @pc.on("connectionstatechange") def on_connection_change(): global dc if pc.connectionState in ["failed", "closed"]: dc = None ... # 其他offer逻辑
关键细节说明
- 保存事件循环:在
offer这个async函数里调用asyncio.get_running_loop(),拿到aiortc运行的事件循环引用,普通线程需要通过这个引用调度异步任务。 - 异步发送封装:把实际的发送逻辑放到
_async_send异步函数里,用await dc.send(msg)正确调用aiortc的异步API。 - 线程安全调度:用
asyncio.run_coroutine_threadsafe把异步任务提交到事件循环,这是普通线程调用异步函数的标准方式,不会阻塞当前线程。 - 状态校验:发送前检查
dc是否存在且处于open状态,避免往已关闭的通道发送消息。
调用示例
- 在async函数里调用:
async def some_async_task(): send_message("来自async任务的消息")
- 在普通线程(比如ROS回调、Thread-4)里调用:
def ros_battery_callback(msg): battery_level = msg.percentage send_message(f"电池电量:{battery_level}%")
内容的提问来源于stack exchange,提问作者Peter Gaston
相关产品推荐
相关产品推荐

