学习链接

webrtc实现视频会议,web多人视频通话,websocket通信

  • 源码:https://gitee.com/zzhua195/video-meeting-basic

关联:

webrtcP2P音视频通话(一)

webrtcP2P音视频通话(二)

webrtc视频会议学习(三)

本地简单示例

客户端代码

client.html
<!DOCTYPE html>
<html lang="en">

<head>
    <meta charset="UTF-8">
    <meta name="viewport" content="width=device-width, initial-scale=1.0">
    <title>Document</title>
    <style>
        body {
            margin: 0;
        }

        .video-wrapper {
            video {
                width: 400px;
                height: 300px;
                border: 1px dashed black;
            }
        }
    </style>
</head>

<body>

    <div class="video-wrapper">
        <video id="localVideo"></video>
        <video id="remoteVideo"></video>
    </div>

    <div>
        目标用户id: <input type="text" id="targetIdIpt" value="2">
        <br />
        <button id="startCallBtn" onclick="startCall">发起通话</button></button>
    </div>
</body>
<script>
    let localVideo = document.querySelector('#localVideo')
    let remoteVideo = document.querySelector('#remoteVideo')
    let targetIdIpt = document.querySelector('#targetIdIpt')

    const params = location.href.substring(location.href.indexOf('?') + 1)
    const _params = {}
    params.split('&').forEach(item => {
        const t = item.split('=')
        _params[t[0]] = t[1]
    })

    const { id: currUserId } = _params

    startCallBtn.addEventListener('click', startCall)

    let localStream = null
    let socket = null
    let peer = null

    async function startCamera() {
        localStream = await navigator.mediaDevices.getUserMedia({ video:true, /*  audio:true  */  })
        localVideo.srcObject = localStream
        localVideo.play()
    }
    // 开启摄像头
    startCamera()

    function initWs() {
        socket = new WebSocket('ws://localhost:9095/ws/' + currUserId)

        socket.onopen = () => {
            console.log('连接建立');
        }

        socket.onmessage = (e) => {
            console.log('收到消息: ', e.data);
            let {code, data} = JSON.parse(e.data)
            console.log('收到消息的code', code);
            if (code == 'connect_success') {
                console.log('连接成功');
            } else if (code == 'offer') {
                acceptCall(data)
            } else if (code == 'icecandidate') {
                console.log('收到icecandidate消息',peer,data)
                let {candidate} = data
                peer.addIceCandidate(candidate)
            } else if (code == 'answer') {
                console.log('对方已同意', data);
                acceptAnswer(data)
            }
        }

        socket.onerror = (err) => {
            console.log('连接发生错误');
        }

        socket.onclose = () => {
            console.log('连接关闭');
        }
    }

    // 初始化websocket
    initWs()


    // 发起通话
    async function startCall() {

        peer = new RTCPeerConnection({})

        localStream.getTracks().forEach(track => {
            peer.addTrack(track, localStream)
        })

        peer.addEventListener('track', e => {
            console.log('发起方ontrack触发', e);
            remoteVideo.srcObject = e.streams[0]
            remoteVideo.play()
        })

        // 这个事件要被触发的前提是, 要先执行上面的把音视频流的轨道给到peer。然后 调用peer.setLocalDescription时会多次连续触发icecandidate事件
        peer.addEventListener('icecandidate', e => {

            // RTCPeerConnectionIceEvent {isTrusted: true, candidate: RTCIceCandidate, type: 'icecandidate', target: RTCPeerConnection, currentTarget: RTCPeerConnection, …}
            console.log('发起方 icecandidate触发', e);

            if (e.candidate) {

                // e.candidate RTCIceCandidate {candidate: 'candidate:1764731264 1 udp 2122260223 192.168.134.…760 typ host generation 0 ufrag KPfd network-id 1', sdpMid: '1', sdpMLineIndex: 1, foundation: '1764731264', component: 'rtp', …}
                console.log('发起方 e.candidate', e.candidate);

                // 内网穿透
                const message = {
                    code: 'icecandidate',
                    data: {
                        targetId: targetIdIpt.value,
                        candidate: e.candidate
                    }
                }
                socket.send(JSON.stringify(message))
            }
        })

        sendOffer()
    }

    async function sendOffer() {
        // RTCSessionDescription {type: 'offer', sdp: 'v=0\r\no=- 8833742509318497131 2 IN IP4 127.0.0.1\r\ns…770 msid:- 8e3b686c-7779-48a1-8714-8a7a0caf0136\r\n'}
        let offer = await peer.createOffer()
        console.log('创建了offer', offer);

        // 设置offer, 会连续触发多次 peer 的 icecandidat 事件, 最后1次的candidate是null
        peer.setLocalDescription(offer)

        const message = {
            code: 'offer',
            data: {
                targetId: targetIdIpt.value,
                offer: offer
            }
        }

        socket.send(JSON.stringify(message))
    }

    function acceptAnswer({fromId, answer}) {
        console.log('设置answer', answer);
        peer.setRemoteDescription(answer)
    }

    function acceptCall({fromId, offer}) {
        console.log('接受通话请求', localStream);

        peer = new RTCPeerConnection({})

        localStream.getTracks().forEach(track => {
            peer.addTrack(track, localStream)
        })

        peer.addEventListener('track', e => {
            console.log('接收方 ontrack触发', e);
            remoteVideo.srcObject = e.streams[0]
            remoteVideo.play()
        })

        // 接收方获得offer {type: 'offer', sdp: 'v=0\r\no=- 1381205837170471750 2 IN IP4 127.0.0.1\r\ns…8adcbdb976 87647395-dbcb-4db4-b105-83145f853af1\r\n'}
        console.log('接收方获得offer', offer);

        peer.addEventListener('icecandidate', e => {

            console.log('接收方 icecandidate触发', e);

            if (e.candidate) {

                console.log('接收方 e.candidate', e.candidate);

                // 内网穿透
                const message = {
                    code: 'icecandidate',
                    data: {
                        targetId: fromId,
                        candidate: e.candidate
                    }
                }

                socket.send(JSON.stringify(message))
            }
        })

        // 设置远程发过来的offer, 会触发 peer 的 track事件 监听函数
        peer.setRemoteDescription(offer)

        sendAnswer(fromId)

    }

    async function sendAnswer(fromId) {
        let answer = await peer.createAnswer()
        peer.setLocalDescription(answer)

        const message = {
            code: 'answer',
            data: {
                targetId: fromId,
                answer: answer
            }
        }

        socket.send(JSON.stringify(message))
    }

    
</script>

</html>

服务端代码

SignalWsServer

就是将消息转发给目标对象

@Slf4j
@Component
@ServerEndpoint("/ws/{id}")
public class SignalWsServer {

    private static ConcurrentHashMap<String, Session> SESSION_MAP = new ConcurrentHashMap<>();

    private Session session;

    private String id;

    @OnOpen
    public void onOpen(@PathParam("id") String id, Session session) {

        log.info("连接::onOpen->{}, {}", id, session);

        SESSION_MAP.put(id, session);
        this.id = id;

        this.session = session;

        this.session.getUserProperties().put("id", id);

        try {
            this.session.getBasicRemote().sendText(
                    JsonUtil.obj2Json(MapBuilder.newHashMap().put("code", "connect_success").build())
            );
        } catch (IOException e) {
            e.printStackTrace();
        }

    }

    @OnMessage
    public void onMessage(String msg) throws Exception {

        log.info("收到客户端 【{}】 的消息::{}", id, msg);

        JSONObject jsonObject = JSON.parseObject(msg);

        String code = String.valueOf(jsonObject.get("code"));

        if ("offer".equals(code)) {
            JSONObject data = jsonObject.getJSONObject("data");
            String targetId = data.getString("targetId");
            Object offer = data.get("offer");
            Session session = SESSION_MAP.get(targetId);
            if (session == null) {
                log.error("目标用户不存在");
                return;
            }
            Map<String, Object> obj = MapBuilder.newHashMap()
                    .put("code", "offer")
                    .put("data", MapBuilder.newHashMap().put("fromId", id).put("offer", offer).build())
                    .build();
            sendMsg(session, obj);
        } else if ("icecandidate".equals(code)) {
            JSONObject data = jsonObject.getJSONObject("data");
            String targetId = data.getString("targetId");

            Object candidate = data.get("candidate");

            Session session = SESSION_MAP.get(targetId);
            if (session == null) {
                log.error("目标用户不存在");
                return;
            }

            Map<String, Object> obj = MapBuilder.newHashMap()
                    .put("code", "icecandidate")
                    .put("data", MapBuilder.newHashMap().put("fromId", id).put("candidate", candidate).build())
                    .build();

            sendMsg(session, obj);

        }  else if ("answer".equals(code)) {
            JSONObject data = jsonObject.getJSONObject("data");
            String targetId = data.getString("targetId");

            Object answer = data.get("answer");

            Session session = SESSION_MAP.get(targetId);
            if (session == null) {
                log.error("目标用户不存在");
                return;
            }

            Map<String, Object> obj = MapBuilder.newHashMap()
                    .put("code", "answer")
                    .put("data", MapBuilder.newHashMap().put("fromId", id).put("answer", answer).build())
                    .build();

            sendMsg(session, obj);

        }

    }

    private void sendMsg(Session session, Object obj) {
        try {
            session.getBasicRemote().sendText(JsonUtil.obj2Json(obj));
        } catch (IOException e) {
            e.printStackTrace();
        }
    }

    @OnClose
    public void onClose() {

        log.info("关闭::onClose->{}",session);

        String id = session.getUserProperties().get("id").toString();
        SESSION_MAP.remove(id);
    }

    @OnError
    public void onError(Throwable ex) throws IOException {
        log.info("发生错误::onError->{}", ex);
        if (this.session.isOpen()) {
            session.close();
        }
    }

}

MapBuilder
public class MapBuilder<K, V> {

    private HashMap<K, V> map = new HashMap();

    public MapBuilder put(K k, V v) {
        map.put(k, v);
        return this;
    }

    public static MapBuilder newHashMap() {
        MapBuilder mapBuilder = new MapBuilder();
        return mapBuilder;
    }

    public Map build() {
        return map;
    }
}

JsonUtil
public class JsonUtil {

    static ObjectMapper mapper = new ObjectMapper();

    public static String obj2Json(Object obj) {
        try {
            return mapper.writeValueAsString(obj);
        } catch (JsonProcessingException e) {
            e.printStackTrace();
            return null;
        }
    }

}
WsConfig
@Configuration
public class WsConfig {

    @Bean
    public ServerEndpointExporter serverEndpointExporter() {
        return new ServerEndpointExporter();
    }

}

效果

在这里插入图片描述

源代码

把https://gitee.com/zzhua195/video-meeting-basic源代码也贴一下,

图解(先看这个)

先看这个,了解过程,再看代码
在这里插入图片描述

前端代码

Single.vue

一对一通信

<template>
<div class="container">
  <div>
    <video ref="localVideo" autoplay></video>
    <video ref="remoteVideo" autoplay></video>
  </div>
  <div>
    <button @click="handleStart">start</button>
  </div>
  <div>
    <input v-model="targetId" type="text">
    <button @click="handleCall">call</button>
  </div>
</div>
</template>

<script setup>
import { ref } from 'vue'

const params = location.href.substring(location.href.indexOf('?') + 1)
const _params = {}
params.split('&').forEach(item => {
  const t = item.split('=')
  _params[t[0]] = t[1]
})

const { id: currentUserId } = _params

const localVideo = ref()
const remoteVideo = ref()
let localStream = null
let socket = null
const targetId = ref(null)
let pc = null

async function handleStart(){
  localStream = await navigator.mediaDevices.getUserMedia({video: true, audio: true})
  localVideo.value.srcObject = localStream
}

function initWebsocket(){
  socket = new WebSocket(`wss://192.168.1.8:8080?id=${currentUserId}`)
  socket.onopen = e => {
    console.log('open_success')
  }

  socket.onmessage = e => {
    const message  = e.data
    const { code, data } = JSON.parse(message)
    if(code === 'connect_success'){
      console.log('connect success')
    }else if(code === 'offer'){
      // 接收offer
      getOffer(data)
    }else if(code === 'answer'){
      // 接收answer
      getAnswer(data)
    }else if(code === 'icecandidate'){
      const { fromId, candidate } = data
      pc.addIceCandidate(candidate)
    }
  }

  socket.onerror = e => {
    console.error(e.data)
  }
}

initWebsocket()

// ====== 发送方 ======= //
function handleCall(){
  pc = new RTCPeerConnection({})
  localStream.getTracks().forEach(track => pc.addTrack(track, localStream))
  pc.addEventListener('track', e => {
    remoteVideo.value.srcObject = e.streams[0]
  })
  pc.addEventListener('icecandidate', e => {
    if(e.candidate){
      // 内网穿透
      const message = {
        code: 'icecandidate',
        data: {
          targetId: targetId.value,
          candidate: e.candidate
        }
      }
      socket.send(JSON.stringify(message))
    }
  })
  sendOffer()
}

async function sendOffer(){
  const description = await pc.createOffer()
  pc.setLocalDescription(description)
  const message = {
    code: 'offer',
    data: {
      targetId: targetId.value,
      offer: description
    }
  }
  socket.send(JSON.stringify(message))
}

function getAnswer({fromId, answer}){
  pc.setRemoteDescription(answer)
}


// ====== 接收方 ======= //
function getOffer({fromId, offer}){
  pc = new RTCPeerConnection({})
  localStream.getTracks().forEach(track => pc.addTrack(track, localStream))
  pc.addEventListener('track', e => {
    remoteVideo.value.srcObject = e.streams[0]
  })
  pc.addEventListener('icecandidate', e => {
    if(e.candidate){
      // 内网穿透
      const message = {
        code: 'icecandidate',
        data: {
          targetId: fromId,
          candidate: e.candidate
        }
      }
      socket.send(JSON.stringify(message))
    }
  })
  pc.setRemoteDescription(offer)
  sendAnswer(fromId)
}

async function sendAnswer(targetId){
  const description =  await pc.createAnswer()
  pc.setLocalDescription(description)
  const message = {
    code: 'answer',
    data: {
      targetId,
      answer: description
    }
  }
  socket.send(JSON.stringify(message))
}
</script>

<style scoped>
video {
  width: 500px;
  height: 400px;
  border: 1px dashed black;
}
</style>
Multiple.vue

一对多通信

<template>
  <div class="container">
    <div ref="videoContainer">
      <video ref="localVideo" src="" autoplay></video>
    </div>
    <div>
      <button @click="handleJoin">join</button>
    </div>
  </div>
</template>

<script setup>
import {ref} from 'vue'

const params = location.href.substring(location.href.indexOf('?') + 1)
const _params = {}
params.split('&').forEach(item => {
  const t = item.split('=')
  _params[t[0]] = t[1]
})

const {id: currentUserId} = _params
const videoContainer = ref()
let socket = null
let pcMap = new Map()
let localStream = null
const localVideo = ref()

function initWebsocket() {
  socket = new WebSocket(`wss://192.168.1.8:8080?id=${currentUserId}`)
  socket.onopen = e => {
    console.log('open_success')
  }

  socket.onmessage = e => {
    const message = e.data
    const {code, data} = JSON.parse(message)
    if (code === 'connect_success') {
      console.log('connect success')
    } else if (code === 'offer') {
      // 接收offer
      getOffer(data)
    } else if (code === 'answer') {
      // 接收answer
      getAnswer(data)
    } else if (code === 'icecandidate') {
      const {fromId, candidate} = data
      const pc = pcMap.get(fromId)
      pc.addIceCandidate(candidate)
    } else if (code === 'join_group') {
      const {fromId} = data
      sendOffer(fromId)
    }
  }

  socket.onerror = e => {
    console.error(e.data)
  }
}

initWebsocket()

async function handleJoin() {
  localStream = await navigator.mediaDevices.getUserMedia({video: true, audio: true})
  localVideo.value.srcObject = localStream
  const message = {
    code: 'join_group',
    data: {
      groupId: 1
    }
  }
  socket.send(JSON.stringify(message))
}

async function sendOffer(targetId) {
  const pc = new RTCPeerConnection({})
  pcMap.set(targetId, pc)
  localStream.getTracks().forEach(track => pc.addTrack(track, localStream))
  const video = document.createElement('video')
  video.autoplay = true
  videoContainer.value.appendChild(video)
  pc.addEventListener('track', e => {
    // 更新远程的视频
    video.srcObject = e.streams[0]
  })

  pc.addEventListener('icecandidate', e => {
    if (e.candidate) {
      const message = {
        code: 'icecandidate',
        data: {
          targetId,
          candidate: e.candidate
        }
      }
      socket.send(JSON.stringify(message))
    }
  })
  const description = await pc.createOffer()
  await pc.setLocalDescription(description)
  const message = {
    code: 'offer',
    data: {
      targetId,
      offer: description
    }
  }
  socket.send(JSON.stringify(message))
}

// 接收offer
async function getOffer({fromId: targetId, offer}) {
  const pc = new RTCPeerConnection({})
  pcMap.set(targetId, pc)
  localStream.getTracks().forEach(track => pc.addTrack(track, localStream))
  const video = document.createElement('video')
  video.autoplay = true
  videoContainer.value.appendChild(video)
  pc.addEventListener('track', e => {
    // 更新远程的视频
    video.srcObject = e.streams[0]
  })

  pc.addEventListener('icecandidate', e => {
    if (e.candidate) {
      const message = {
        code: 'icecandidate',
        data: {
          targetId,
          candidate: e.candidate
        }
      }
      socket.send(JSON.stringify(message))
    }
  })
  await pc.setRemoteDescription(offer)
  const description = await pc.createAnswer()
  await pc.setLocalDescription(description)
  const message = {
    code: 'answer',
    data: {
      targetId,
      answer: description
    }
  }
  socket.send(JSON.stringify(message))
}

async function getAnswer({fromId, answer}){
  const pc = pcMap.get(fromId)
  await pc.setRemoteDescription(answer)
}
</script>

<style>
video {
  width: 500px;
  height: 400px;
  border: 1px dashed black;
}
</style>

后端代码

websocket.js - node
// https://www.npmjs.com/package/ws?activeTab=readme

const {WebSocketServer} = require('ws')
const {createServer} = require('https')
const {readFileSync} = require('fs')

const server = createServer({
  cert: readFileSync('./cert/server.pem'),
  key: readFileSync('./cert/server.key')
});
const wss = new WebSocketServer({server});

// 在线人员
const groupOnlineMembers = new Set()

wss.on('connection', function connection(ws, request) {
  const _params = {}
  const params = request.url.substring(request.url.indexOf('?') + 1)
  params.split('&').forEach(item => {
    const t = item.split('=')
    _params[t[0]] = t[1]
  })

  const {id: currentUserId} = _params
  ws._userId = currentUserId

  ws.on('error', console.error);

  ws.on('close', e => {
    groupOnlineMembers.delete(currentUserId)
  })

  ws.on('message', function message(res) {
    const message = res.toString()
    const {code, data} = JSON.parse(message)
    if (code === 'offer') {
      const {targetId, offer} = data
      wss.clients.forEach(client => {
        if (client._userId === targetId) {
          const message = {
            code: 'offer',
            data: {
              fromId: currentUserId,
              offer
            }
          }
          client.send(JSON.stringify(message))
        }
      })
    } else if (code === 'answer') {
      // 转发answer
      const {targetId, answer} = data
      wss.clients.forEach(client => {
        if (client._userId === targetId) {
          const message = {
            code: 'answer',
            data: {
              fromId: currentUserId,
              answer
            }
          }
          client.send(JSON.stringify(message))
        }
      })
    } else if (code === 'icecandidate') {
      const {targetId, candidate} = data
      wss.clients.forEach(client => {
        if (client._userId === targetId) {
          const message = {
            code: 'icecandidate',
            data: {
              fromId: currentUserId,
              candidate
            }
          }
          client.send(JSON.stringify(message))
        }
      })
    } else if (code === 'join_group') {
      const { groupId } = data
      groupOnlineMembers.add(currentUserId)
      wss.clients.forEach(client => {
        if (groupOnlineMembers.has(client._userId) && client._userId !== currentUserId) {
          // 该用户在线
          const message = {
            code: 'join_group',
            data: {
              fromId: currentUserId
            }
          }
          client.send(JSON.stringify(message))
        }
      })
    }
  });

  const message = {
    code: 'connect_success',
    data: ''
  }
  ws.send(JSON.stringify(message));
});

server.listen(8080);
WebSocket - java
@Component
@ServerEndpoint("/websocket/{appId}")
@Slf4j
public class WebSocket {
    private String appId;
    private static final ConcurrentHashMap<String, Session> sessionPool = new ConcurrentHashMap<String, Session>();

    private static final Set<String> groupList = new HashSet<>();

    /**
     * 链接成功调用的方法
     */
    @OnOpen
    public void onOpen(Session session, @PathParam(value = "appId") String appId) {
        try {
            log.debug(appId);
            this.appId = appId;
            sessionPool.put(appId, session);
            session.getId();
            log.debug("【websocket消息】有新的连接,总数为:" + sessionPool.size());
            JSONObject result = new JSONObject();
            result.put("code", "connect_success");
            result.put("data", "success");
            this.sendMessage(appId, result.toJSONString());
        } catch (Exception e) {
            e.printStackTrace();
            log.error(e.getMessage());
        }
    }

    /**
     * 链接关闭调用的方法
     */
    @OnClose
    public void onClose(Session session, CloseReason closeReason) {
        try {
            sessionPool.remove(this.appId);
            groupList.remove(this.appId);
            log.debug("【websocket消息】连接断开,总数为:" + sessionPool.size());
        } catch (Exception e) {
            e.printStackTrace();
            log.error(e.getMessage());
        }
    }

    /**
     * 收到客户端消息后调用的方法
     */
    @OnMessage
    public void onMessage(String message, Session session) {
        log.debug(appId);
        log.debug("【websocket消息】收到客户端消息:" + message);
        if (ObjectUtil.isNotEmpty(message)) {
            JSONObject json = JSON.parseObject(message);
            String code = json.getString("code");
            if ("offer".equals(code)) {
                // 发送offer
                JSONObject data = json.getJSONObject("data");
                String targetId = data.getString("targetId");
                JSONObject description = data.getJSONObject("description");
                JSONObject response = new JSONObject();
                response.put("code", "offer");
                JSONObject result = new JSONObject();
                result.put("fromId", this.appId);
                result.put("description", description);
                response.put("data", result);
                this.sendMessage(targetId, response.toJSONString());
            } else if ("answer".equals(code)) {
                // 发送answer
                JSONObject data = json.getJSONObject("data");
                String targetId = data.getString("targetId");
                JSONObject description = data.getJSONObject("description");

                JSONObject response = new JSONObject();
                response.put("code", "answer");
                response.put("errcode", "0");
                JSONObject result = new JSONObject();
                result.put("fromId", this.appId);
                result.put("description", description);
                response.put("data", result);
                this.sendMessage(targetId, response.toJSONString());
            } else if ("candidate".equals(code)) {
                // 发送candidate
                JSONObject data = json.getJSONObject("data");
                String targetId = data.getString("targetId");
                JSONObject candidate = data.getJSONObject("candidate");

                JSONObject response = new JSONObject();
                response.put("code", "candidate");
                JSONObject result = new JSONObject();
                result.put("fromId", this.appId);
                result.put("candidate", candidate);
                response.put("data", result);
                this.sendMessage(targetId, response.toJSONString());
            } else if ("join_group_video".equals(code)) {
                groupList.add(this.appId);
                JSONObject data = json.getJSONObject("data");
                String groupId = data.getString("targetId");
                JSONObject response = new JSONObject();
                response.put("code", "join_group_video");
                JSONObject result = new JSONObject();
                result.put("fromId", this.appId);
                response.put("data", result);
                this.sendMoreMessage(groupList, response.toJSONString());
            }
        }
    }

    /**
     * 发送错误时的处理
     */
    @OnError
    public void onError(Session session, Throwable error) {
        log.info("用户错误,原因:" + error.getMessage());
        error.printStackTrace();
    }

    // 此为单点消息
    public void sendMessage(String appId, String message) {
        Session session = sessionPool.get(appId);
        if(session == null){
            return;
        }
        synchronized(session){
            if (session.isOpen()) {
                try {
                    log.info("【websocket消息】 单点消息:" + message);
                    session.getAsyncRemote().sendText(message);
                } catch (Exception e) {
                    e.printStackTrace();
                    log.error(e.getMessage());
                }
            }
        }
    }

    // 此为单点消息(多人)
    public void sendMoreMessage(Set<String> appIds, String message) {
        for (String appId : appIds) {
            if(appId.equals(this.appId)){
                continue;
            }
            Session session = sessionPool.get(appId);
            if (session != null && session.isOpen()) {
                try {
                    log.info(appId + "【websocket消息】 单点消息:" + message);
                    session.getAsyncRemote().sendText(message);
                } catch (Exception e) {
                    e.printStackTrace();
                }
            }
        }

    }
}

更多推荐