项目背景

和各位读者大致介绍下具体场景,线上的小程序中开放一些语音麦克风的房间,让用户进入房间之后可以互相通过语音聊天的方式进行互动。

这里分享一下相关的技术设计方案。这款系统的核心点设计在于如何能让一个用户发出的语音通知到其他用户上边O ~ y h ; ! 5 Y l。语音数据在客户端同事的处理下最终变成了io数据流请求到了后端,后端只需要将这b d 8 / + 0 ~ $ e些数据流传达给各个不同的终端即可达到广播通知的效果。

单机版架构

最初期上线的时候,为了赶速度,快速试错,所以简单地采用了单机版架构去设! _ o \ F 0 j 6计。结合技术栈为 SpringBoot,WebSocket,MySQL技术) I S

线上一间语音房间的同时在线M V _人数并不会特别多,大概在15-50人的区间段内,t K ^ H系统核心代码是通过SpringBoot内部的WebSocket技术去进行数据的主动推送。

设计思路

整体的设计图比较简单,基本就是一台服务器存储WebSocket连接,如下图所示:

用户进行WebSocket初始化连接的时候需要一个连接分配和存储的过程:

早期的存储是存放在了服^ E H ! D j + h务器本地的一个Map集合中。

N ` 3 + m w 6 ) @WeS ( g ( v Y b 3 sbSocket进行连, O q接的时候就会往内存中写入一条数据信息,m e # e 2 5当链接断开的时候,就将内存中k M 4的数据移除。然后进行语音广播的时候需要结合WebSocket内部的广播发送功能进行通知

看似设计比较简单,但是在后期业务变得庞大的时候出现了瓶颈。因为随着参加语音活动用户的增C 4 q 8 ?加,越来越多的WebSocketSession对象需要被存储到内存当中,这种有状态性的存储对于单机扩容不灵活。

设计缺陷

1.假设原先的服务器扩容到了A,B两台机. z L 1 0 9器,A用户在A机器上边建9 R Q K = f \ = 9立了WebSocketSession,B用户在B机器上边建立的WebSo– 1 ccketSession连接。此时如果A想要和B进行对话发送,需要先查找到具体WebSocketSession存放在哪台机器上边。

2.当用户出现了网络异常,临时断开连接进行重连的时候,也可能会出现1所说的问题。

集群架构

设计思路

一旦出现需要发送语音通知的时候,发送一条广播的mq消息,每个机器都接收到消息之后,触发自己的广播操作U v t Z _即可。

RocketMq的接入系统设计里C ] g u M ^ * $面mq采用的是广播模式,这和我们通常使用的集群模式有一定的区别。

消息队列Rn ^ 2 h $ T 0ocketMQ版是基于发布或订阅模型的消息系统。消费者,即消息的订m R O y { o阅方订阅关注的Top0 9 f = b ?ic,以获取并消费消息。由于消费者应用一般是分布式系统,以P 4 # . V 6 = s集群方式部署,因此消息队列RocketMQb 8 h版约定以下概念:

  • 集群:使用相同Group Ih b d 9 1 m 9D的消费者属于同一个集群。同一个集群下的消费者消费逻辑必须完全一致(包括Tag的使用)。
  • 集群消费:当使用集群消费模式时,消息队列RocketMQ版认为任意一条消息只需要被集群内的任意一个消费者处理即可。
  • 广播消费:当使用广播消费模式时,消息队列RocketMQ版会将每条消息推送给集群内所有注册过的消费者,保证消息至少被每个消费者消费一次。

集群消费模式适用场景 适用于消费端集群化部署,每条消息只需要被处理一次的场景。此外,由于消费Z m { / I _ P进度在服务端维护,可靠性更高。具体消费示例如下图所示。

注意事项

  • 集群消费模式下,每一条消息都只# g o会被分发到一台机器上处理。如果需要被P U ]集群下的每一台机器都处理,请使用广播模式。
  • 集群消费模式下,不保证每一次失败重投的消息路由到同一台机器h G t o j \ G j ?上。

广播消费模式适用场景 适用于消费端集群化部署,每条消息需# O C r F要被集群下的f t S W v每个消费者处理的场景。具体消费示例如下图所示。

[ 4 Q C g F意事项

  • 广播消费模式下不支持顺序消息。
  • 广[ H I N + 8 I C播消L n T E Y m 4 R 3费模式下不支持重置消费位点。
  • 每条消息都需要被B [ @ + t l e + Y相同订阅逻辑的多台机器处理。
  • 消费进度在客户端维护,出现重复消费l g 4 : @的概率稍大于集群模式。
  • 广播模式下,消息队列RocketMQ版保证每条消息至少被每台客户端消费一次,但是并不会重投消费失败的消息,因此业务方需要关注消费失败的情况。
  • 广播模式下,客户端每一次重启都会从最新消息消费。客户端在被停止期间发送至服务端的消息将会被自动跳过,请谨慎选择。
  • 广播模式下,每条消息都会被大量的客户端重复处理,因此推荐尽可能使用集群模式。
  • 广播模式下服务端不维护消费进度,t u X所以消息队列RocketMQ版控制台不支持消息堆积查询、消息堆积} F u报警和订阅关系查询T d s Da + R o n能。

这里面的应用场景需要对集群内部对每个消费者都对服务器内存中的socket连接进行session是否存在对判断,因此需要采用mq的广播模式。

关于mq部分的接入代码

CZ 7 m honsumer模块的配置:

  1. packageorg.idea.web.socket.config;
  2. importorg.springframework_ { D A l.boot.context.properties.CU ; I = s +onfiguW ? * \ @ )rationProperties;
  3. /! 8 R | _**
  4. *@Q U HAuthC @ B W _ # 9 ~ .orlinhao
  5. *@Datecreatedin10:30上午2021/5/10
  6. */
  7. @H S ^ G + h E WConfigurationPropertieS n y x E \s(prefix="rocketz , e Xmq.consumer")
  8. publicclassMqConsumerConfig{
  9. privatebooleanisOn;
  10. privateStringgroupName;
  11. prK * @ uivateStringnameSrvAddr;
  12. privateStL P 2ringtopics;
  13. privateIntegerconsumeThreadMin;
  14. privateIntegerconsumeThreadMax;
  15. privateIntegerconsumeMessageBatchMaxSize;
  16. /**
  17. getter和setter部分省略
  18. **/
  19. }

Producer模块的配置展示:

  1. packageorg.idea.web.socket.config;
  2. importorg.sprins d j ?g8 2 m Q l } | +framework.boot.context.properties.ConfigurationPropertien i P y / F _s;
  3. /**
  4. *@Authorlinhao
  5. *@Datecrk _ ] z Jeatedin10:26上午2021/5/10
  6. */
  7. @ConfigurationProperties(prefix="rocketd $ ]mq.producer")
  8. publicclassMqProducerConfig{
  9. pH l l [ 1 2 : |rivatebooleanisOn;
  10. privateStY d |ringgroupName;
  11. privateStringnameSrv] s e Y GAddr;
  12. privateIntegermaxMessageSize;
  13. privateIntegersendMsgTimeout;
  14. privateIntegerretryTimesWhenSendFailed;
  15. /**: Q _ O Y
  16. getter和setter部分! 7 ! A _ j G \ L省略
  17. **/
  18. }

RocketMq内部的消费端Bean配置

  1. packageorg.idea.web.socket.mq;
  2. importlombok.extern.slf4j.Slf4j;
  3. importorg.apache.rocketmq.c\ % O |lient.consumer.DefaultMQPushConsumer;
  4. impon @ K -rtorg.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
  5. importorg.apachea r c #.rocked z E ^ o U ~ dtmq.client.exception.MQClientException;
  6. importorg.apache.rocketmq.common.consumer.CA q 8 D m o h _onsumeFromWhere;
  7. importorg.apache.rocketmq.common.protocol.heartbeat.MessageModel;
  8. importorg.idea.web.socket.config.MqConsumerConfig;
  9. importorg.idea.web.socket.config.d ` Y 7 t V g tMqProducerConfig;
  10. importorg.springframework.boot.autoconfigur_ c q 8e.AutoConfigureAfter;
  11. importo^ 2 &rg.springframework.boot.autoconfigure.u p F o Y o + x GAutoConfigureBefore;
  12. importorg.springframework.boot.autoconfigure.cos u \ $ x - vndition.ConditionalOnMissingBean;
  13. importorg.spa ^ K N ( ,ringframework.boot.context.properties.EnableConfigu] j F U 7 Y b Y IrationProperties;
  14. importorg.springfrl m j f c J v 1amework.context.annotation.Bean;
  15. importorg.springframework.context.annotation.Configuration;
  16. importjavax.annotatio0 K * 9 Bn.Resource;
  17. /**
  18. *@Authorlinhao
  19. *@B 6 [ L ) pDatecreatedin10:! t s % 734上午2021/5/10
  20. */
  21. @ConfigP } 0 H luration
  22. @Slf4j
  23. @EnableConfigurationProp# d Serties({MqConA } b H m &sumerCoI \ ]nfig.class})
  24. publicclassMqConsumerAutoConfig{
  25. @Resource
  26. privateMqConsumerConfigmqC8 5 ~ NonsumerConfig;
  27. @Resource
  28. //这个接口需要手动实现顺序消费的逻辑每次获取到消息队列的第一条数据
  29. privateMessageListeE $ V s 1 onerHandlermessageListenerConcurrently;
  30. @B$ t ~ ( D R \ 3ean
  31. @s E )ConditionalOnMissingBean
  32. publicDefaultMQPushConsumerdefaultMQPushConsumer(){
  33. Defs I t w 9 #aultMQP7 Z = $ U v M }u_ Z n J @ + 7 6 {shConsumerconsumer=neW * ~ M 9 2 E = %wDefaultMQPushConsumer();
  34. consumer.setNamesrvAk J e pddr(mqConsumerConfig.getNameSrvAddr());
  35. consumer.setConsumerGroup(mqConsumerConfig.getGroupName()} + & x G Y Q r q);
  36. consumer.setConsumeThreadMin(mqConsumerConfig.gQ O & S | yetConsumeThreadMin());
  37. consumer.setConsumeThreadMax(mqConsumerConfig.getConsumeThreadMax());
  38. consumer.registerMessageLi* B L =stener(messageListenerConcurrentl[ 8 /y);
  39. consumer4 ; # R ? M X ).setCo/ 5 _ f a b rnsumeFro$ n V * u @ U lmWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFS ` D { P A 8 yFSET);
  40. //消费模型是什么?
  41. consumer.setMessageModel(MessageModel.BROADCASTING);
  42. //默认一次拉取一条消费
  43. consumer.setConsumeM= - h T ~essageBat6 H ] r U q z ` UchMaxSize(mqConsumerConfig.getConsumeMessageBatchMaxSize());
  44. //*表示订阅所有的tag
  45. try{
  46. conT i \ bsJ v [umer.subscribe(y V J ] 4mqConsumerConfig.getTopics(),"*");
  47. consumer.start();
  48. log.info("【MqConsumerAutoConfig】mqconsumerisstarted!");
  49. }catch(Exceptione){
  50. log.error("k & d + 4 S rmqstartfail,eis",e);
  51. }
  52. returnconsumer;
  53. }
  54. }

Ro_ _ w u T o j ) :cketMq的服务生产者Bean配置

  1. packageorg.idea.web.socket.mq;
  2. importlombok.extern.slf4j.Slf4j;
  3. importorg.apache.rocketmq.client.producer.DefaultMQProducer;
  4. i( j . z \ Dmportorg.idea.weN E f r r :b.socket.config.MqProducerConfig;
  5. importorg.springframework.boot.autou c KcoP B % =nfigure.AutoConfigureAfter;
  6. importorg.springframewC L p | Mork.boot.autoconfigure.AutoConfy = N 3 h _igk 6 ( 8 m gureBeforeb p n o A j | - r;
  7. importorg.springfQ p ? B u 0 U Tramework.boot.autoconfigure.condition.ConditionalOnMissingBean;
  8. importorg.springframework.boot.context.properties.EnableConfigurationProperties;
  9. importorg.springframewor! A _ J B w Fk.contex/ j 9 \ { ` Y u $t.annotation.Bean;
  10. importorg.springframework.context.annotation.Configuration;
  11. importjavax.ann* R U l { : [ ~ ;otation.Resource;
  12. /**
  13. *@Authorlinhao
  14. *@Datecreatedin11:05上午2021/5/10
  15. */
  16. @Configuration
  17. @e e rSlf4j
  18. @EnableConfigurationProperties({MqProducerCh H Y F sonfig.class})q D C b
  19. publicclash 5 /sMqProd! F K lucerAutoConfig{
  20. @Re= } f \ ) m K OsourH 1 mce
  21. privV i . & P a ` uateMqPrA k H C t 9 EoducerConfigmqProducerConfig;
  22. @Bean
  23. @Condit+ d C b +ionalOnMissingBean
  24. //意味着DefaultMQP@ Y j v e E qroducer的配置可以被覆盖
  25. publicDefa 2 A \ [ n w 3 _aultMQProducerdefaultMQProducer(){
  26. DefaultMQProducerproducer=newDef] j ] a L Q daultMQProducer(mqProducerConfig.getGroupName());
  27. producer.setNamesrvAddr(mqProducerConfig.getNameSrvAddr());
  28. //没有则自动创建topic的key
  29. //producer.setCreateTopicKey("T : G l Q 7 oAUTO_CREATE_TOPIC_KEY");
  30. producer.setMaxMessageSize(mqProducerConfig.getMaxMessageSize());
  31. producer.setSendMsgTimeout(mqe x S k y xProducerConfig.getSendMsgTimeout());K e ( u L = n
  32. producer.setRetryTimesWhenSendFailed(mqProducerConfig.m 1 @ A C kgetRetryTimesWhenSendFailed());
  33. try{
  34. producer.start();
  35. log.inf[ = D K a :o("【MqProducerAutoConfig】mqproducerisstarted!");
  36. }catch(Exceptione){
  37. log.error("[MqProducerAutoConfig]startfail,eis",e);
  38. }
  39. retp F + Q `urnproducer;
  40. }
  41. }@ 7 + I w O C

然后是对RocketMq内部发送消息事件的一层函数封装

  1. packageorg.idea.web.socket.mq;
  2. importcom.alibaba.fastjson.JSON;
  3. importlombok.extern.slf4j.Slf4j;
  4. importorg.apache.commons.lang3.Stri| v 0 EngUtils;
  5. importorg.N 3 5 g N ;apache.rocketmq.client.producer.DefaultMQProducer;
  6. importorg.apache.rocketmq.client.producer.SendResult;
  7. ix / \ O ] J | . Zmportorg.apache.rocketmq.common.mes$ 5 Asa! $ j K j * v Wge.Message;
  8. i- | V \mportorg.apac% d 8 n Whe.rocketmq.remoting.common.RemotingHelper;
  9. importorg.idea.web.# K ) Y V n gsocket.config.MqPC \ z i - H T ZroducerConfig;
  10. importorg.idea.web.socket.dto.BroadcastMqDTO;
  11. importorg.spring_ w )framt X Q ) p ` J Rework.stereotype.ComO w d k G ] Y m (ponent;
  12. importjn 1 ` $ f e Tavax.ann& v + 9 k ? h * notation.Resource;
  13. importjava.io.UnsupportedEncodingException;
  14. /**
  15. *消息广播发送端
  16. *
  17. *@Authorlinhao
  18. *@Datecreatedin10:43下午2021/5/9
  19. */F Z & M r J
  20. @Component
  21. @Slf4j
  22. publicclassBroadcastMqProducer{
  23. @Resource
  24. privateDefaultMQProducerdefau` R # \ s P r !ltMQProducer;
  25. @Resource
  26. prib F ; n $ ^ ~ rvateMqProducerCom h 8 ? F 2nfigmqProducerConfig;
  27. privatestaticStringTOPIC="ws-topic";
  28. privatestaticStringTAGS="ws-tag";
  29. publicstatic$ x G _ 0 - C jIntegerALL_USER_RECEIVE_TYPE=1a A ~;
  30. publicstaticIntegerONE_USER_RECEIVE_TYPE=2;
  31. /**
  32. *点对点之间的消息发送
  33. *
  34. *@paramdestSessionKey
  35. *@C g ! h r h $ K #parammsg
  36. *@return
  37. */
  38. publicSendResultsendWebSocketToUser(StringdestSessionKey,Stringmsg){
  39. if(StringUtils.isEmptm t N 5y(msg)){
  40. log.error("[sendWebSocketToUser]msgcannotbe$ @ p ! Snull!");
  41. returnnull;
  42. }
  43. Messagemessage=null;
  44. SendResultsendResult=null;
  45. try{
  46. BroadcastMG i L , x o \ !qDTObroadcastMqDTO=newBroadcastMqDTO();
  47. broadcr V ( U f _ O AastMqDTO.setEventType(ONE_USER_RECEIVE_TYPE);
  48. broadcastMqDTO.setMessage(msg);
  49. broadcastMqDTl * ) v * Q R * +O.setSessionKey(destSessionKey);
  50. messag: 8 Ie=ns S CewMessage(TOZ 3 K [ } ! / R nPIC,TAGS,(JSON.toJSONString(broadce f = ? # D ) !astMqDTO))K h A Z _ , j A.getBytes(RemotingHelper.DEFAULT_CHARSET));
  51. sendResult=defaultMQProducer.send(message);
  52. }catch(Exceptione){
  53. log.error("[sendWebSocketBroadcastMsg]eis",e);
  54. }
  55. returnsendResult;
  56. }
  57. /**
  58. *广播消息发送
  59. *
  60. *@parammsg
  61. *@return
  62. */
  63. publicSendResultC f e !sendWebSocket; M $ d , ] qBroadcastMsg(StringmS F e !sg_ a n M){
  64. if(StringUtils.isEmpty(msg)){
  65. log.error("[sendWebSocketBroadcastMsg]msgcannotbenull!");
  66. returnnull;
  67. }
  68. Messagemessage5 - = * 8 Q Q=null;
  69. SendResultsendResult=null;
  70. try{
  71. BroadcastMqDTObrF c soadcastMqDTO=newBroadcastMqDTO();
  72. broadcastMqDTO.setEventType(ALL_USER_RECEIVE_TYPE);O d w A ! u 3
  73. broadcastMqDTO.setMess@ ` . qage(msg);
  74. mes6 X e v ] vsaged k U 8=newV C LMessage(TOPIC,TAGS,(JSO0 T % t UN.toJSONString(broadcastMqDTO)).getBytes(RemotingHelper.DEFAULT_CHARSET));
  75. sendResult=defaultMQProducer.send(a m A M N Gmessage);
  76. }catch(Exceptione){
  77. log.error("[sendWebSocketBroadcastMsg]eis",e);
  78. }
  79. returnsendRW m ) ) 6 t vesult;
  80. }
  81. }

对消息的! r i订阅模块实现代码如下:

  1. packageorx } Xg.idea.web.socket.mq;
  2. importcom.alibaba.fastjson.JSON;
  3. importcom.oracle.tool9 z W - z rs.packager.Log;
  4. importlombok.extern.slf4j.Slf4jy n ; =;
  5. importorg.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
  6. importorg.apache.rocketmq.client.consumer.listeR W D R M $ _ Lner.Coo q 2 9 \ M r 3 mnsumeConO = ] a scurrentlyStatus;, 3 7 q - 7
  7. importl 9 n ( 1 I jorg.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
  8. importorg.apache.rocketmq.co) @ * w dmmon.message.MessageExt;
  9. importorg.idea.web.socketo 1 g o.dto.BroadcastMqDTO;
  10. importorg, H # 6 X H = Z P.idea.web.socket.manager.SocketManager;
  11. importorg.springframework.messaging.simp.SimpMessagingTemplate;
  12. importorg.springframework.stereotype.Component;
  13. importorg.springframework.util.CollectionUtils;
  14. importorg.springframework.web.socket.WebSocketSession;
  15. importjavax.annotation.Resource;
  16. importjava.util.L$ 4 } 3ist;
  17. importstaticorg.idea.web.socket.mq.BroadcastMqProd) 9 Q V * ] 3ucer.ALL_USER_RECEIVE_TYPE;
  18. importstaticorg.idea.web.socket.mq] $ O , H ].BroadcastMqProducer.ONE_USER_RECEIVE_TYPE;
  19. /**
  20. *@Authorlinhao
  21. *@Datecreatedin10:59上午2021/5/10
  22. */
  23. @Component
  24. @Slf4jb y w ; { } \ 5 M
  25. publicclassMessageListenerHandlerimplementsMessageListenerConcurrently{
  26. @Resource
  27. privateSocketManagersocketManager;
  28. @Resource
  29. privateSimpMessagingTemplatetemplate;
  30. @Override
  31. publicConsumeConcurrentlyStatusconsumeMessage(List<MessageExt>list,ConsumeConcurrentl; B x S _ _yContextconsumeConcurrentlyContext){
  32. if(CollectionUtils.isEmpty(list)){
  33. Log.info("receiA 4 q tveemptymV f o V j W ysg");
  34. returnCo| S - b S . @ C ,nsumeConcurrentlyStatus.CONSUME_SUCCESS;
  35. }
  36. MessageExtmessageExt=list.get(0);
  37. byte[]bytes` ` h=messageExt.getBodyx F q 6 6 U ; g J();
  38. Stringjson=newString(bytes);! ] w h
  39. BroadcastMqDTObroadcastMqDTO=JSON.parseObject(json,BroadcastMqDTO.class)c Q | g B s q c k;
  40. log.info("[MessageListenerHandler]broadcastMqDTOis"+broadcastMqDTO);
  41. if(ALL_USER_RECEIVE_TYPE.equals(broadcastMqDTO.getEventType())){
  42. log.info("[consumeMessage]广播发送消\ S Y :息:触发----》消息内容为:"+broadcastMqDTO);
  43. template.convertAndSend("/topic/sendTc . V f 3 v | 1 {opic",broadcastMqDTO);
  44. }elseif(ONE_USER_RECEIVE_TYPE.equals(broadcastMqDTO.getEventType())){
  45. StringsessionKey=broadcastMqDTO.getSel T ? k J dssionKey();
  46. WebSocketSessionZ \ c . W \webSocketSession=socketManagh J s j e ` )er.get(sessi{ 4 { h ionKey6 ] Y y S n 6 2);
  47. if(webSock& ; N h h ] I +etSession!=n2 _ Cull){
  48. template.convertAndSendToUser(sessionKey,"/queue/sendUser",broadcasY j F _ j / 7 k otMqDTO.getMessage());
  49. log.info("[consumeMessage]点对点8 = N a M发送消息;触发----》消息内容为:"+broadcastMqDTO);
  50. }
  51. }
  52. returnConsumeConcurrentlyStatus.CONSUME_SUCCESS;
  53. }
  54. }

整体设计结构如下图:

于是按照这个结构进行了一版本的紧急开发迭代,原先的单台服务器扩展为了服务集群。

业务拓展后续产品经理提出一个需求,要求支持在同一间房内的两个用户之间发送悄悄话功能。这就需要我们进行一个点对点之间传输通讯的功能了。因此需要在mq通知到每台机器的时候加一个本地Sesw ] (sion遍历的逻辑,如果当前机器存有用户token对应的session变0 1 W ~ * _ `量,那么就单独针对那个Session进行WebSocket的发送通知。

设计弊端u ^ +一旦某台机器出现了异常崩溃,那么就意味着这台机器上的所有语音连接可能会出现中断情况。目前这一块的问题也在考虑解决,计划是将WebSocketV } N oSession存入到分布式缓存的redis中保证数据e { l 9 T a Q |可靠存储,但是在后续尝试的时q & & o 7 : P %候发现WebSocketSession对象没有实现序列化= p & 0 q n接口,在存储到Redis的时] A 3 # ) O r d %候会出现异常。目前这个问题还在寻找解决思路中F k b . y t –,不知道各位读者6 J ^j C 5 N ^友们有什么好的思路。

遇到的问题点用户请求直接访问到了我们的内部服务器,如果在请求的中间加入一台nginx做负载均衡则需要在nginx中配置一些额外信息。

项目的源代码比较多,这里我把核心部分的代码整理了一份,感兴趣的朋友可以到我的gitee上边去下载:

https://gitee.co. J m J om/IdeaHome_admin/socket-framework

【责任编辑:庞桂玉 TEL:(010)68476606】

点赞 0

发表评论

您的电子邮箱地址不会被公开。 必填项已用*标注