爆款云主机2核4G限时秒杀,88元/年起!
查看详情

活动

天翼云最新优惠活动,涵盖免费试用,产品折扣等,助您降本增效!
热门活动
  • 618智算钜惠季 爆款云主机2核4G限时秒杀,88元/年起!
  • 免费体验DeepSeek,上天翼云息壤 NEW 新老用户均可免费体验2500万Tokens,限时两周
  • 云上钜惠 HOT 爆款云主机全场特惠,更有万元锦鲤券等你来领!
  • 算力套餐 HOT 让算力触手可及
  • 天翼云脑AOne NEW 连接、保护、办公,All-in-One!
  • 中小企业应用上云专场 产品组合下单即享折上9折起,助力企业快速上云
  • 息壤高校钜惠活动 NEW 天翼云息壤杯高校AI大赛,数款产品享受线上订购超值特惠
  • 天翼云电脑专场 HOT 移动办公新选择,爆款4核8G畅享1年3.5折起,快来抢购!
  • 天翼云奖励推广计划 加入成为云推官,推荐新用户注册下单得现金奖励
免费活动
  • 免费试用中心 HOT 多款云产品免费试用,快来开启云上之旅
  • 天翼云用户体验官 NEW 您的洞察,重塑科技边界

智算服务

打造统一的产品能力,实现算网调度、训练推理、技术架构、资源管理一体化智算服务
智算云(DeepSeek专区)
科研助手
  • 算力商城
  • 应用商城
  • 开发机
  • 并行计算
算力互联调度平台
  • 应用市场
  • 算力市场
  • 算力调度推荐
一站式智算服务平台
  • 模型广场
  • 体验中心
  • 服务接入
智算一体机
  • 智算一体机
大模型
  • DeepSeek-R1-昇腾版(671B)
  • DeepSeek-R1-英伟达版(671B)
  • DeepSeek-V3-昇腾版(671B)
  • DeepSeek-R1-Distill-Llama-70B
  • DeepSeek-R1-Distill-Qwen-32B
  • Qwen2-72B-Instruct
  • StableDiffusion-V2.1
  • TeleChat-12B

应用商城

天翼云精选行业优秀合作伙伴及千余款商品,提供一站式云上应用服务
进入甄选商城进入云市场创新解决方案
办公协同
  • WPS云文档
  • 安全邮箱
  • EMM手机管家
  • 智能商业平台
财务管理
  • 工资条
  • 税务风控云
企业应用
  • 翼信息化运维服务
  • 翼视频云归档解决方案
工业能源
  • 智慧工厂_生产流程管理解决方案
  • 智慧工地
建站工具
  • SSL证书
  • 新域名服务
网络工具
  • 翼云加速
灾备迁移
  • 云管家2.0
  • 翼备份
资源管理
  • 全栈混合云敏捷版(软件)
  • 全栈混合云敏捷版(一体机)
行业应用
  • 翼电子教室
  • 翼智慧显示一体化解决方案

合作伙伴

天翼云携手合作伙伴,共创云上生态,合作共赢
天翼云生态合作中心
  • 天翼云生态合作中心
天翼云渠道合作伙伴
  • 天翼云代理渠道合作伙伴
天翼云服务合作伙伴
  • 天翼云集成商交付能力认证
天翼云应用合作伙伴
  • 天翼云云市场合作伙伴
  • 天翼云甄选商城合作伙伴
天翼云技术合作伙伴
  • 天翼云OpenAPI中心
  • 天翼云EasyCoding平台
天翼云培训认证
  • 天翼云学堂
  • 天翼云市场商学院
天翼云合作计划
  • 云汇计划
天翼云东升计划
  • 适配中心
  • 东升计划
  • 适配互认证

开发者

开发者相关功能入口汇聚
技术社区
  • 专栏文章
  • 互动问答
  • 技术视频
资源与工具
  • OpenAPI中心
开放能力
  • EasyCoding敏捷开发平台
培训与认证
  • 天翼云学堂
  • 天翼云认证
魔乐社区
  • 魔乐社区

支持与服务

为您提供全方位支持与服务,全流程技术保障,助您轻松上云,安全无忧
文档与工具
  • 文档中心
  • 新手上云
  • 自助服务
  • OpenAPI中心
定价
  • 价格计算器
  • 定价策略
基础服务
  • 售前咨询
  • 在线支持
  • 在线支持
  • 工单服务
  • 建议与反馈
  • 用户体验官
  • 服务保障
  • 客户公告
  • 会员中心
增值服务
  • 红心服务
  • 首保服务
  • 客户支持计划
  • 专家技术服务
  • 备案管家

了解天翼云

天翼云秉承央企使命,致力于成为数字经济主力军,投身科技强国伟大事业,为用户提供安全、普惠云服务
品牌介绍
  • 关于天翼云
  • 智算云
  • 天翼云4.0
  • 新闻资讯
  • 天翼云APP
基础设施
  • 全球基础设施
  • 信任中心
最佳实践
  • 精选案例
  • 超级探访
  • 云杂志
  • 分析师和白皮书
  • 天翼云·创新直播间
市场活动
  • 2025智能云生态大会
  • 2024智算云生态大会
  • 2023云生态大会
  • 2022云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心

      根据源码,模拟实现 RabbitMQ - 网络通讯设计,自定义应用层协议,实现 BrokerServer (8)

      首页 知识中心 云计算 文章详情页

      根据源码,模拟实现 RabbitMQ - 网络通讯设计,自定义应用层协议,实现 BrokerServer (8)

      2024-03-29 09:48:26 阅读次数:53

      rabbitmq,分布式

      一、网络通讯协议设计


      1.1、交互模型

      目前我们需要考虑的交互模型:生产者消费者都是客户端,都需要通过 网络 和 BrokerServer 进行通信

      根据源码,模拟实现 RabbitMQ - 网络通讯设计,自定义应用层协议,实现 BrokerServer (8)

      此处我们使⽤ TCP 协议, 来作为通信的底层协议. 同时在这个基础上⾃定义应⽤层协议, 完成客⼾端对服 务器这边功能的远程调⽤.

      TCP 是有连接的(Connection),创建 / 断开 TCP 连接成本还是挺高的(需要三次握手啥的),那么这里就是用 Channel 来表示 Connection 内部的 “逻辑上” 的连接,使得 “一个管道,多个网线传输” 的效果,使得 TCP连接得到复用

      Ps:要远程调用的功能就是在 VirtualHost 中 public 的方法.

      1.2、自定义应用层协议

      1.2.1、请求和响应格式约定

      之前我们定义的 Message 对象,本体就是二进制的数据,因此这里不方便使用 JSON 这种文本协议 / 格式.

      因此这里使用 二进制 的方式来设定协议.

      请求如下:

      /**
       * 表示一个网络通信中的请求对象,按照自定义协议的格式展开
       */
      public class Request {
      
          private int type;
          private int length;
          private byte[] payload;
      
          public int getType() {
              return type;
          }
      
          public void setType(int type) {
              this.type = type;
          }
      
          public int getLength() {
              return length;
          }
      
          public void setLength(int length) {
              this.length = length;
          }
      
          public byte[] getPayload() {
              return payload;
          }
      
          public void setPayload(byte[] payload) {
              this.payload = payload;
          }
      }
      

      根据源码,模拟实现 RabbitMQ - 网络通讯设计,自定义应用层协议,实现 BrokerServer (8)

      响应如下:

      /**
       * 这个对象表示一个响应,是根据自定义应用层协议来的
       */
      public class Response {
      
          private int type;
          private int length;
          private byte[] payload;
      
          public int getType() {
              return type;
          }
      
          public void setType(int type) {
              this.type = type;
          }
      
          public int getLength() {
              return length;
          }
      
          public void setLength(int length) {
              this.length = length;
          }
      
          public byte[] getPayload() {
              return payload;
          }
      
          public void setPayload(byte[] payload) {
              this.payload = payload;
          }
      }
      

      根据源码,模拟实现 RabbitMQ - 网络通讯设计,自定义应用层协议,实现 BrokerServer (8)

      1.2.2、参数说明

      1)type是一个整形,用来表示当前这个请求和响应是用来干啥的(对应 VirtualHost 中的核心 API),取值如下:

      • 0x1 创建 channel
      • 0x2 关闭 channel
      • 0x3 创建 exchange
      • 0x4 销毁 exchange
      • 0x5 创建 queue
      • 0x6 销毁 queue
      • 0x7 创建 binding
      • 0x8 销毁 binding
      • 0x9 发送 message
      • 0xa 订阅 message
      • 0xb 返回 ack
      • 0xc 服务器给客⼾端推送的消息.(被订阅的消息) 响应独有的.

      2)length 就是用来描述 payload 长度(防止粘包问题)

      3)payload 就是具体要传输的二进制数据。数据具体是什么,会根据当前是请求还是响应,以及当前的 type 的不同取值来确定。

      比如 type 是 0x3(创建交换机),同时当前是一个请求,此时 payload 里的内容,就相当于 exchangeDeclare 的 参数 的序列化的结果.

      比如 type 是 0x3(创建交换机),同时当前是一个响应,此时 payload 里的内容,就是 exchangDeclare 的 返回结果 的序列化内容.

      根据源码,模拟实现 RabbitMQ - 网络通讯设计,自定义应用层协议,实现 BrokerServer (8)

      1.2.3、具体例子

      栗子如下:

      1)请求

      当前需要远程调用 exchangeDeclare 方法,那么我们就需要传递核心 API 以下参数

      根据源码,模拟实现 RabbitMQ - 网络通讯设计,自定义应用层协议,实现 BrokerServer (8)

      使用一个公共的父类包装每次 请求 中公共(每个请求都要传输)的参数

      /**
       * 这个类用来表示方法的公共参数/辅助字段
       * 后续每个方法会有一些不同的参数,不同的参数再用不同的子类来表示
       */
      public class BasicArguments implements Serializable {
      
          // 表示一次 请求/响应 的身份标识,让请求和响应能对的上
          protected String rid;
          // 表示这次通信使用的 channel 的身份标识
          protected String channelId;
      
          public String getRid() {
              return rid;
          }
      
          public void setRid(String rid) {
              this.rid = rid;
          }
      
          public String getChannelId() {
              return channelId;
          }
      
          public void setChannelId(String channelId) {
              this.channelId = channelId;
          }
          
      }
      

      创建 ExchangeDeclareArguments 类(当前这个类将来会被序列化成 request 类中的 payload),继承 BasicArguments(公共参数),实现 Serializable 接口(避免序列化问题),要传递的参数如下:

      public class ExchangeDeclareArguments extends BasicArguments implements Serializable {
      
          private String exchangeName;
          private ExchangeType exchangeType;
          private boolean durable;
          private boolean autoDelete;
          private Map<String, Object> arguments;
      
          public String getExchangeName() {
              return exchangeName;
          }
      
          public void setExchangeName(String exchangeName) {
              this.exchangeName = exchangeName;
          }
      
          public ExchangeType getExchangeType() {
              return exchangeType;
          }
      
          public void setExchangeType(ExchangeType exchangeType) {
              this.exchangeType = exchangeType;
          }
      
          public boolean isDurable() {
              return durable;
          }
      
          public void setDurable(boolean durable) {
              this.durable = durable;
          }
      
          public boolean isAutoDelete() {
              return autoDelete;
          }
      
          public void setAutoDelete(boolean autoDelete) {
              this.autoDelete = autoDelete;
          }
      
          public Map<String, Object> getArguments() {
              return arguments;
          }
      
          public void setArguments(Map<String, Object> arguments) {
              this.arguments = arguments;
          }
      }
      

      2)响应

      当前 VirtualHost 中的核心 API 返回值都是 Boolean 类型,因此我们使用一个公共类来封装响应(当前这个类将来会被序列化成 response 类中的 payload 参数)

      public class BasicReturns implements Serializable {
      
          //用来标识唯一的请求和响应
          protected String rid;
          //标识一个 channel
          protected String channelId;
          //标识当前这个远程调用方法的返回值
          protected boolean ok;
      
          public String getRid() {
              return rid;
          }
      
          public void setRid(String rid) {
              this.rid = rid;
          }
      
          public String getChannelId() {
              return channelId;
          }
      
          public void setChannelId(String channelId) {
              this.channelId = channelId;
          }
      
          public boolean isOk() {
              return ok;
          }
      
          public void setOk(boolean ok) {
              this.ok = ok;
          }
      }
      

      Ps:其他核心 API 自定义应用层协议也一样

       

      1.2.4、特殊栗子

      0xa 订阅 message ,这个核心 API 比较特殊,参数中有回调函数

      根据源码,模拟实现 RabbitMQ - 网络通讯设计,自定义应用层协议,实现 BrokerServer (8)

       1)请求

      创建 BasicConsumeArguments 类(当前这个类将来会被序列化成 request 类中的 payload) 表示要传递的参数,需要注意的是 Consumer 这个回调,在发送的请求中不需要携带这个参数(实际上也携带不了)

      Ps:因为服务器收到这个订阅消息请求之后,就直接取拿队列中的消息,接着直接反馈给客户端,客户端拿到消息后才执行回调方法(要拿这个消息干什么事)。

      这就类似于你去商店订阅报纸,接着拿到报纸以后,你要对这个报纸做什么,商店是不知道的~~

      public class BasicConsumeArguments extends BasicArguments implements Serializable {
      
          private String consumerTag;
          private String queueName;
          private boolean autoAck;
      
          //注意! 这里的 Consumer 回调函数不用发送给服务器(实际上也发送不了)
          //因为服务器收到这个订阅消息请求之后,就直接取拿队列中的消息,接着直接反馈给客户端
          //客户端拿到消息后才执行回调方法
          //这就类似于你去商店订阅报纸,接着拿到报纸以后,你要对这个报纸做什么,商店是不知道的~~
      
      
          public String getConsumerTag() {
              return consumerTag;
          }
      
          public void setConsumerTag(String consumerTag) {
              this.consumerTag = consumerTag;
          }
      
          public String getQueueName() {
              return queueName;
          }
      
          public void setQueueName(String queueName) {
              this.queueName = queueName;
          }
      
          public boolean isAutoAck() {
              return autoAck;
          }
      
          public void setAutoAck(boolean autoAck) {
              this.autoAck = autoAck;
          }
      }
      

      2)响应

      创建 SubScribeReturns 类(当前这个类将来会被序列化成 response 类中的 payload 参数) 来描述响应, 这个响应中不光要携带 BasicReturns (返回的公共响应参数),还需要带上回调中消息的参数,如下:

      public class SubScribeReturns extends BasicReturns implements Serializable {
      
          private String consumerTag;
          private BasicProperties basicProperties;
          private byte[] body;
      
          public String getConsumerTag() {
              return consumerTag;
          }
      
          public void setConsumerTag(String consumerTag) {
              this.consumerTag = consumerTag;
          }
      
          public BasicProperties getBasicProperties() {
              return basicProperties;
          }
      
          public void setBasicProperties(BasicProperties basicProperties) {
              this.basicProperties = basicProperties;
          }
      
          public byte[] getBody() {
              return body;
          }
      
          public void setBody(byte[] body) {
              this.body = body;
          }
      }
      

      1.3、实现 BrokerServer

      这里的写法就和以前写过的 TCP 回显服务器很类似了,只是根据请求计算响应的方式不同

      1.3.1、属性和构造

          private ServerSocket serverSocket = null;
      
          //当前考虑一个 BrokerServer 上只有一个 虚拟主机
          private VirtualHost virtualHost = new VirtualHost("default");
          //使用 哈希表 来标识当前所有会话(哪个客户端正在和服务器进行通信)
          //key 是 channelId, value 为对应的 Socket 对象
          private ConcurrentHashMap<String, Socket> sessions = new ConcurrentHashMap<>();
          //用线程池来处理多个客户端请求
          private ExecutorService executorService = null;
          //引入一个 Boolean 变量控制服务器是否继续运行
          private volatile boolean runnable = true;
      
          public BrokerServer(int port) throws IOException {
              serverSocket = new ServerSocket(port);
          }
      

      1.3.2、启动 BrokerServer

          public void start() throws IOException {
              System.out.println("[BrokerServer] 启动!");
              executorService = Executors.newCachedThreadPool();
              while(runnable) {
                  Socket clientSocket = serverSocket.accept();
                  //处理连接的逻辑给线程池
                  executorService.submit(() -> {
                      processConnection(clientSocket);
                  });
              }
          }
      

      1.3.3、停止 BrokerServer

          /**
           * 停止服务器,一般是直接 kill 就可以了
           * 此处这个单独的方法,主要是为了后续的单元测试
           */
          public void stop() throws IOException {
              runnable = false;
              //放弃线程池中的任务,并销毁线程
              executorService.shutdown();
              serverSocket.close();
          }
      

      1.3.4、处理每一个客户端连接

          private void processConnection(Socket clientSocket) {
              try (InputStream inputStream = clientSocket.getInputStream();
                   OutputStream outputStream = clientSocket.getOutputStream()) {
                  // 这里需要按照特定格式来读取并解析. 此时就需要用到 DataInputStream 和 DataOutputStream
                  try (DataInputStream dataInputStream = new DataInputStream(inputStream);
                       DataOutputStream dataOutputStream = new DataOutputStream(outputStream)) {
                      while (true) {
                          // 1. 读取请求并解析.
                          Request request = readRequest(dataInputStream);
                          // 2. 根据请求计算响应
                          Response response = process(request, clientSocket);
                          // 3. 把响应写回给客户端
                          writeResponse(dataOutputStream, response);
                      }
                  }
              } catch (EOFException | SocketException e) {
                  // 对于这个代码, DataInputStream 如果读到 EOF , 就会抛出一个 EOFException 异常.
                  // 需要借助这个异常来结束循环
                  System.out.println("[BrokerServer] connection 关闭! 客户端的地址: " + clientSocket.getInetAddress().toString()
                          + ":" + clientSocket.getPort());
              } catch (IOException | ClassNotFoundException | MqException e) {
                  System.out.println("[BrokerServer] connection 出现异常!");
                  e.printStackTrace();
              } finally {
                  try {
                      // 当连接处理完了, 就需要记得关闭 socket
                      clientSocket.close();
                      // 一个 TCP 连接中, 可能包含多个 channel. 需要把当前这个 socket 对应的所有 channel 也顺便清理掉.
                      clearClosedSession(clientSocket);
                  } catch (IOException e) {
                      e.printStackTrace();
                  }
              }
          }

      1.3.5、读取请求和写响应

          private Request readRequest(DataInputStream dataInputStream) throws IOException {
              Request request = new Request();
              request.setType(dataInputStream.readInt());
              request.setLength(dataInputStream.readInt());
              byte[] body = new byte[request.getLength()];
              int n = dataInputStream.read(body);
              if(n != request.getLength()) {
                  throw new IOException("读出请求格式出错!");
              }
              request.setPayload(body);
              return request;
          }
      
          private void writeResponse(DataOutputStream dataOutputStream, Response response) throws IOException {
              dataOutputStream.write(response.getType());
              dataOutputStream.write(response.getLength());
              dataOutputStream.write(response.getPayload());
              dataOutputStream.flush();
          }
      

      1.3.6、根据请求计算响应

      这里就是根据不同的 type 类型,来远程调用 VirtualHost 中不同的核心 API(需要特别注意订阅消息功能的回调函数)

          private Response process(Request request, Socket clientSocket) throws IOException, ClassNotFoundException, MqException {
              //1.将 request 初步解析成 BasicArguments
              BasicArguments basicArguments = (BasicArguments) BinaryTool.fromBytes(request.getPayload());
              System.out.println("[Request] rid=" + basicArguments.getRid() + ", channelId=" + basicArguments.getChannelId() +
                      ", type=" + request.getType() + ", length=" + request.getLength());
              //2.根据 type 的值,进一步区分接下来要干什么
              boolean ok = true;
              if (request.getType() == 0x1) {
                  //创建 channel
                  sessions.put(basicArguments.getChannelId(), clientSocket);
                  System.out.println("[BrokerServer] 创建 channel 完成!channelId=" + basicArguments.getChannelId());
              } else if(request.getType() == 0x2) {
                  //销毁 channel
                  sessions.remove(basicArguments.getChannelId());
                  System.out.println("[BrokerServer] 销毁 channel 完成!channelId=" + basicArguments.getChannelId());
              } else if(request.getType() == 0x3) {
                  //创建交换机,此时 payLoad 就是 ExchangDeclareArguments 了
                  ExchangeDeclareArguments arguments = (ExchangeDeclareArguments) basicArguments;
                  ok = virtualHost.exchangeDeclare(arguments.getExchangeName(), arguments.getExchangeType(),
                          arguments.isDurable(), arguments.isAutoDelete(), arguments.getArguments());
              } else if(request.getType() == 0x4) {
                  ExchangeDeleteArguments arguments = (ExchangeDeleteArguments) basicArguments;
                  ok = virtualHost.exchangeDelete(arguments.getExchangeName());
              } else if(request.getType() == 0x5) {
                  QueueDeclareArguments arguments = (QueueDeclareArguments) basicArguments;
                  ok = virtualHost.queueDeclare(arguments.getQueueName(), arguments.isDurable(),
                          arguments.isExclusive(), arguments.isAutoDelete(), arguments.getArguments());
              } else if(request.getType() == 0x6) {
                  QueueDeleteArguments arguments = (QueueDeleteArguments) basicArguments;
                  ok = virtualHost.queueDelete(arguments.getQueueName());
              } else if(request.getType() == 0x7) {
                  QueueBindArguments arguments = (QueueBindArguments) basicArguments;
                  ok = virtualHost.queueBind(arguments.getQueueName(), arguments.getExchangeName(), arguments.getBindingKey());
              } else if(request.getType() == 0x8) {
                  QueueUnBindArguments arguments = (QueueUnBindArguments) basicArguments;
                  ok = virtualHost.queueUnBind(arguments.getQueueName(), arguments.getExchangeName());
              } else if(request.getType() == 0x9) {
                  BasicPublishArguments arguments = (BasicPublishArguments) basicArguments;
                  ok = virtualHost.basicPublish(arguments.getExchangeName(), arguments.getRoutingKey(), arguments.getBasicProperties(), arguments.getBody());
              } else if(request.getType() == 0xa) {
                  BasicConsumeArguments arguments = (BasicConsumeArguments) basicArguments;
                  ok = virtualHost.basicConsume(arguments.getConsumerTag(), arguments.getQueueName(), arguments.isAutoAck(), new Consumer() {
                      //这个回调函数要做的就是,把服务器收到的消息可以直接推送回对应的消费者客户端
                      @Override
                      public void handlerDelivery(String consumerTag, BasicProperties basicProperties, byte[] body) throws MqException, IOException {
                          //首先需要知道收到的消息要发给哪个客户端
                          //此处 consumerTag 其实就是 channelId,根据 channelId 去 sessions 中查询,既可以得到对应的
                          //socket 对象了,从而往里面发送数据
                          //1.根据 channelId 找到 socket 对象
                          Socket clientSocket = sessions.get(consumerTag);
                          if(clientSocket == null || clientSocket.isClosed()) {
                              throw new MqException("[BrokerServer] 订阅消息的客户端已经关闭!");
                          }
                          //2.构造响应数据
                          SubScribeReturns subScribeReturns = new SubScribeReturns();
                          subScribeReturns.setChannelId(consumerTag);
                          subScribeReturns.setRid("");//由于这里只有响应,没有请求,不需要去对应,rid 暂时不需要
                          subScribeReturns.setOk(true);
                          subScribeReturns.setConsumerTag(consumerTag);
                          subScribeReturns.setBasicProperties(basicProperties);
                          subScribeReturns.setBody(body);
                          byte[] payload = BinaryTool.toBytes(subScribeReturns);
                          Response response = new Response();
                          // 0xc 表示服务器给消费者客户端推送的消息数据
                          response.setType(0xc);
                          response.setLength(payload.length);
                          response.setPayload(payload);
                          //3.把数据写回给客户端
                          //  注意!此处的 dataOutputStream 这个对象不能 close
                          //  如果把 dataOutputStream 关闭, 就会直接把 clientSocket 里的 outputStream 也关了
                          //  此时就无法继续往 socket 中写后续的数据了
                          DataOutputStream dataOutputStream = new DataOutputStream(clientSocket.getOutputStream());
                          writeResponse(dataOutputStream, response);
                      }
                  });
              } else if(request.getType() == 0xb) {
                  //调用 basicAck 确认消息
                  BasicAckArguments arguments = (BasicAckArguments) basicArguments;
                  ok = virtualHost.basicAck(arguments.getQueueName(), arguments.getMessageId());
              } else {
                  throw new MqException("[BrokerServer] 未知 type!type=" + request.getType());
              }
              //构造响应
              BasicReturns basicReturns = new BasicReturns();
              basicReturns.setRid(basicArguments.getRid());
              basicReturns.setChannelId(basicArguments.getChannelId());
              basicReturns.setOk(ok);
              byte[] payload = BinaryTool.toBytes(basicReturns);
              Response response = new Response();
              response.setType(request.getType());
              response.setLength(request.getLength());
              response.setPayload(payload);
              System.out.println("[Response] rid=" + basicReturns.getRid() + ", channelId=" + basicReturns.getChannelId()
                      + ", type=" + response.getType() + ", length=" + response.getLength());
              return response;
          }
      

      1.3.7、清除 channel

      清理 map 中对应的(clientSocket) session 信息

          private void clearClosedSession(Socket clientSocket) {
              List<String> toDeleteChannelId = new ArrayList<>();
              for(Map.Entry<String, Socket> entry : sessions.entrySet()) {
                  if(entry.getValue() == clientSocket) { //这里一个 key 可能对应多个相同的 Socket
                      //在集合类中不能一边用迭代器一边删除,会破坏迭代器结构的!
                      //sessions.remove(entry.getKey());
                      //因此这里先记录下 key
                      toDeleteChannelId.add(entry.getKey());
                  }
              }
              for(String channelId : toDeleteChannelId) {
                  sessions.remove(channelId);
              }
              System.out.println("[BrokerServer] 清理 session 完毕!channelId=" + toDeleteChannelId);
          }
      
      版权声明:本文内容来自第三方投稿或授权转载,原文地址:https://blog.csdn.net/CYK_byte/article/details/132446025,作者:陈亦康,版权归原作者所有。本网站转在其作品的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如因作品内容、版权等问题需要同本网站联系,请发邮件至ctyunbbs@chinatelecom.cn沟通。

      上一篇:python的多进程并行计算---分布式计算、联邦学习

      下一篇:RabbitMQ - 单机部署(超详细)

      相关文章

      2025-05-06 09:18:49

      【Linux 从基础到进阶】Ceph分布式存储系统搭建

      随着数据量的爆炸式增长,传统的存储解决方案逐渐暴露出扩展性差、成本高、管理复杂等问题。Ceph是一种高性能、可扩展的开源分布式存储系统,能够为对象存储、块存储和文件系统提供统一的存储平台。

      2025-05-06 09:18:49
      分布式 , 存储 , 高可用性
      2025-04-18 07:10:53

      LDAP基础理论

      分布式目录服务是一种用于存储和管理大量数据的系统,其中数据以层次结构的方式组织,并在多个服务器之间进行分布。它提供了一种集中式访问和管理数据的方法,使得用户可以通过网络连接到任何一个服务器来查询、添加、修改或删除存储在该目录中的信息。

      2025-04-18 07:10:53
      LDAP , 分布式 , 属性 , 服务器 , 用户 , 目录
      2025-04-15 09:24:56

      探秘Redis分布式锁:实战与注意事项

      Redis的Watch命令可以实现乐观锁,这是一种保护数据完整性的机制。在分布式环境中,当多个客户端并发地操作相同的键时,乐观锁有助于防止数据竞争和冲突。

      2025-04-15 09:24:56
      Redis , Redisson , 分布式 , 实例 , 客户端 , 获取
      2025-04-15 09:19:55

      分布式事务大揭秘:使用MQ实现最终一致性

      在单体应用中,事务的管理相对简单,可以通过数据库的事务机制来保证数据的一致性和完整性。然而,在分布式系统中,由于涉及到多个不同的服务和数据源,保证事务的一致性就变得复杂了。

      2025-04-15 09:19:55
      RocketMQ , 一致性 , 事务 , 分布式 , 发送 , 消息 , 系统
      2025-04-11 07:11:40

      java中常用的缓存框架

      java中常用的缓存框架

      2025-04-11 07:11:40
      API , Cache , Java , 分布式 , 缓存
      2025-03-28 07:42:50

      分布式存储技术

      分布式存储技术是一种数据存储技术,它通过网络将企业中每台机器上的磁盘空间利用起来,并将这些分散的存储资源构成一个虚拟的存储设备,实现数据的分散存储。

      2025-03-28 07:42:50
      分布式 , 存储 , 数据 , 节点 , 虚拟化
      2025-03-12 09:33:43

      【分布式理论13】分布式存储:数据存储难题与解决之道

      【分布式理论13】分布式存储:数据存储难题与解决之道

      2025-03-12 09:33:43
      RAID , 分布式 , 存储 , 数据 , 磁盘
      2025-03-12 09:32:22

      docker安装rabbitmq详解

      docker安装rabbitmq详解

      2025-03-12 09:32:22
      rabbitmq , 客户端 , 镜像
      2025-03-11 09:36:54

      【分布式理论12】事务协调者高可用:分布式选举算法

      【分布式理论12】事务协调者高可用:分布式选举算法

      2025-03-11 09:36:54
      分布式 , 算法 , 节点
      2025-03-06 09:15:52

      Spring Boot + Shiro 实现 Session 持久化实现思路及遗留问题

      Spring Boot + Shiro 实现 Session 持久化实现思路及遗留问题

      2025-03-06 09:15:52
      Session , 分布式 , 服务器 , 用户 , 登陆
      查看更多
      推荐标签

      作者介绍

      天翼云小翼
      天翼云用户

      文章

      33561

      阅读量

      5252953

      查看更多

      最新文章

      【Linux 从基础到进阶】Ceph分布式存储系统搭建

      2025-05-06 09:18:49

      LDAP基础理论

      2025-04-18 07:10:53

      分布式事务大揭秘:使用MQ实现最终一致性

      2025-04-15 09:19:55

      分布式存储技术

      2025-03-28 07:42:50

      【分布式理论13】分布式存储:数据存储难题与解决之道

      2025-03-12 09:33:43

      【分布式理论12】事务协调者高可用:分布式选举算法

      2025-03-11 09:36:54

      查看更多

      热门文章

      Android移动设备远程接入ZooKeeper分布式集群

      2023-04-18 14:14:56

      分布式版本控制系统——git

      2023-06-13 08:29:57

      python学习——分布式进程

      2023-05-08 10:00:50

      Elasticsearch分布式架构原理(二)

      2023-06-01 06:30:49

      分布式系统常见的事务处理机制

      2023-05-23 01:22:38

      分布式-技术专区-Redis分布式锁原理实现

      2023-05-29 10:45:37

      查看更多

      热门标签

      系统 测试 用户 分布式 Java java 计算机 docker 代码 数据 服务器 数据库 源码 管理 算法
      查看更多

      相关产品

      弹性云主机

      随时自助获取、弹性伸缩的云服务器资源

      天翼云电脑(公众版)

      便捷、安全、高效的云电脑服务

      对象存储

      高品质、低成本的云上存储服务

      云硬盘

      为云上计算资源提供持久性块存储

      查看更多

      随机文章

      分布式-技术专区-Redis分布式锁原理实现

      使用Java构建高效的分布式缓存系统

      Redisson分布式锁主从一致性问题解决

      基于grpc从零开始搭建一个准生产分布式应用(系列)

      Kafka安装记录

      RabbitMQ【路由模式(概念、编写生产者 、编写消费者 ) 通配符模式(概念、编写生产者、编写消费者) 】(四)-全面详解(学习总结---从入门到深化)

      • 7*24小时售后
      • 无忧退款
      • 免费备案
      • 专家服务
      售前咨询热线
      400-810-9889转1
      关注天翼云
      • 旗舰店
      • 天翼云APP
      • 天翼云微信公众号
      服务与支持
      • 备案中心
      • 售前咨询
      • 智能客服
      • 自助服务
      • 工单管理
      • 客户公告
      • 涉诈举报
      账户管理
      • 管理中心
      • 订单管理
      • 余额管理
      • 发票管理
      • 充值汇款
      • 续费管理
      快速入口
      • 天翼云旗舰店
      • 文档中心
      • 最新活动
      • 免费试用
      • 信任中心
      • 天翼云学堂
      云网生态
      • 甄选商城
      • 渠道合作
      • 云市场合作
      了解天翼云
      • 关于天翼云
      • 天翼云APP
      • 服务案例
      • 新闻资讯
      • 联系我们
      热门产品
      • 云电脑
      • 弹性云主机
      • 云电脑政企版
      • 天翼云手机
      • 云数据库
      • 对象存储
      • 云硬盘
      • Web应用防火墙
      • 服务器安全卫士
      • CDN加速
      热门推荐
      • 云服务备份
      • 边缘安全加速平台
      • 全站加速
      • 安全加速
      • 云服务器
      • 云主机
      • 智能边缘云
      • 应用编排服务
      • 微服务引擎
      • 共享流量包
      更多推荐
      • web应用防火墙
      • 密钥管理
      • 等保咨询
      • 安全专区
      • 应用运维管理
      • 云日志服务
      • 文档数据库服务
      • 云搜索服务
      • 数据湖探索
      • 数据仓库服务
      友情链接
      • 中国电信集团
      • 189邮箱
      • 天翼企业云盘
      • 天翼云盘
      ©2025 天翼云科技有限公司版权所有 增值电信业务经营许可证A2.B1.B2-20090001
      公司地址:北京市东城区青龙胡同甲1号、3号2幢2层205-32室
      • 用户协议
      • 隐私政策
      • 个人信息保护
      • 法律声明
      备案 京公网安备11010802043424号 京ICP备 2021034386号