【消息队列开发】 实现BrokerServer类——本体服务器

这篇具有很好参考价值的文章主要介绍了【消息队列开发】 实现BrokerServer类——本体服务器。希望对大家有所帮助。如果存在错误或未考虑完全的地方,请大家不吝赐教,您也可以点击"举报违法"按钮提交疑问。

🍃前言

本次开发任务

  • 实现 BrokerServer 类,也就是咱们消息队列的本体服务器。

其实本质上就是一个 TCP 的服务器。
【消息队列开发】 实现BrokerServer类——本体服务器,消息队列开发,服务器,运维,java,spring boot

🎋创建 BrokerServer 类

创建 BrokerServer 类如下:

public class BrokerServer {
	// 当前程序只考虑⼀个虚拟主机的情况.
	private VirtualHost virtualHost = new VirtualHost("default-VirtualHost");
	// key 为 channelId, value 为 channel 对应的 socket 对象.
	private ConcurrentHashMap<String, Socket> sessions = new ConcurrentHashMap<>
	private ServerSocket serverSocket;
	private ExecutorService executorService;
	private volatile boolean runnable = true;
}
  • virtualHost表示服务器持有的虚拟主机.队列,交换机,绑定,消息都是通过虚拟主机管理.
  • sessions ⽤来管理所有的客⼾端的连接. 记录每个客户端的 socket.
  • serverSocket 是服务器自身的 socket
  • executorService 这个线程池用来处理响应
  • runnable 这个标志位用来控制服务器的运行停⽌

🎍启动与停止服务器

代码实现如下:

public BrokerServer(int port) throws IOException {
    serverSocket = new ServerSocket(port);
}

public void start() throws IOException {
    System.out.println("[BrokerServer] 启动!");
    executorService = Executors.newCachedThreadPool();
    try {
        while (runnable) {
            Socket clientSocket = serverSocket.accept();
            // 把处理连接的逻辑丢给这个线程池.
            executorService.submit(() -> {
                processConnection(clientSocket);
            });
        }
    } catch (SocketException e) {
        System.out.println("[BrokerServer] 服务器停止运行!");
        // e.printStackTrace();
    }
}

// 一般来说停止服务器, 就是直接 kill 掉对应进程就行了.
// 此处还是搞一个单独的停止方法. 主要是用于后续的单元测试.
public void stop() throws IOException {
    runnable = false;
    // 把线程池中的任务都放弃了. 让线程都销毁.
    executorService.shutdownNow();
    serverSocket.close();
}

🍀实现处理连接

通过这个方法, 来处理一个客户端的连接.

我们使用 InputStreamOutputStream,由于后面要按照特定格式来读取并解析.

此时就需要用到 DataInputStreamDataOutputStream

在这一个连接中, 可能会涉及到多个请求和响应,我们使用一个while(true)来进行实现

在此循环我们要做的事情有三件:

  1. 读取请求并解析
  2. 根据请求计算响应
  3. 把响应写回客户端

具体处理逻辑,我们后面再仔细实现,

那么我们怎么结束这个循环呢?

注意我们上面使用的是 DataInputStreamDataOutputStream,当没有数据进行读取的时候,就会进行抛出异常而结束循环

最后当连接处理完了, 就需要记得关闭 socket, 一个 TCP 连接中, 可能包含多个 channel. 需要把当前这个 socket 对应的所有 channel 也顺便清理掉.

代码实现如下:

// 通过这个方法, 来处理一个客户端的连接.
// 在这一个连接中, 可能会涉及到多个请求和响应.
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();
        }
    }
}

🎄实现 readRequest 与 writeResponse

【消息队列开发】 实现BrokerServer类——本体服务器,消息队列开发,服务器,运维,java,spring boot

关于读取请求,我们前面定义了一个类为 Request ,此时我们构造相应的对象,并对该对象相应属性进行填充即可。

代码实现如下:

private Request readRequest(DataInputStream dataInputStream) throws IOException {
    Request request = new Request();
    request.setType(dataInputStream.readInt());
    request.setLength(dataInputStream.readInt());
    byte[] payload = new byte[request.getLength()];
    int n = dataInputStream.read(payload);
    if (n != request.getLength()) {
        throw new IOException("读取请求格式出错!");
    }
    request.setPayload(payload);
    return request;
}

关于响应,实现相反,传入的 响应对象 相应的属性返回即可。

最后不要忘了刷新缓冲区

代码实现如下:

private void writeResponse(DataOutputStream dataOutputStream, Response response) throws IOException {
    dataOutputStream.writeInt(response.getType());
    dataOutputStream.writeInt(response.getLength());
    dataOutputStream.write(response.getPayload());
    // 这个刷新缓冲区也是重要的操作!!
    dataOutputStream.flush();
}

🌴实现处理请求

先把请求转换成 BaseArguments , 获取到其中的 channelId 和 rid

再根据不同的 type, 分别处理不同的逻辑. (主要是调用virtualHost中不同的方法).
【消息队列开发】 实现BrokerServer类——本体服务器,消息队列开发,服务器,运维,java,spring boot

针对消息订阅操作,则需要在存在消息的时候通过回调,把响应结果写回给对应的客⼾端.

最后构造成统⼀的响应.

代码实现如下:

private Response process(Request request, Socket clientSocket) throws IOException, ClassNotFoundException, MqException {
    // 1. 把 request 中的 payload 做一个初步的解析.
    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 就是 ExchangeDeclareArguments 对象了.
        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 handleDelivery(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 的 payload 就是一个 SubScribeReturns
                        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 {
        // 当前的 type 是非法的.
        throw new MqException("[BrokerServer] 未知的 type! type=" + request.getType());
    }
    // 3. 构造响应
    BasicReturns basicReturns = new BasicReturns();
    basicReturns.setChannelId(basicArguments.getChannelId());
    basicReturns.setRid(basicArguments.getRid());
    basicReturns.setOk(ok);
    byte[] payload = BinaryTool.toBytes(basicReturns);
    Response response = new Response();
    response.setType(request.getType());
    response.setLength(payload.length);
    response.setPayload(payload);
    System.out.println("[Response] rid=" + basicReturns.getRid() + ", channelId=" + basicReturns.getChannelId()
            + ", type=" + response.getType() + ", length=" + response.getLength());
    return response;
}

🌲实现 clearClosedSession

这里要做的事情, 主要就是遍历上述 sessions hash 表, 把该被关闭的 socket 对应的键值对, 统统删掉

需要注意的是:

  • 我们在进行迭代的时候,不要直接删除,这样会影响集合类的结构

代码实现如下:

private void clearClosedSession(Socket clientSocket) {
    // 这里要做的事情, 主要就是遍历上述 sessions hash 表, 把该被关闭的 socket 对应的键值对, 统统删掉.
    List<String> toDeleteChannelId = new ArrayList<>();
    for (Map.Entry<String, Socket> entry : sessions.entrySet()) {
        if (entry.getValue() == clientSocket) {
            // 不能在这里直接删除!!!
            // 这属于使用集合类的一个大忌!!! 一边遍历, 一边删除!!!
            // sessions.remove(entry.getKey());
            toDeleteChannelId.add(entry.getKey());
        }
    }
    for (String channelId : toDeleteChannelId) {
        sessions.remove(channelId);
    }
    System.out.println("[BrokerServer] 清理 session 完成! 被清理的 channelId=" + toDeleteChannelId);
}

⭕总结

关于《【消息队列开发】 实现BrokerServer类——本体服务器》就讲解到这儿,感谢大家的支持,欢迎各位留言交流以及批评指正,如果文章对您有帮助或者觉得作者写的还不错可以点一下关注,点赞,收藏支持一下文章来源地址https://www.toymoban.com/news/detail-854719.html

到了这里,关于【消息队列开发】 实现BrokerServer类——本体服务器的文章就介绍完了。如果您还想了解更多内容,请在右上角搜索TOY模板网以前的文章或继续浏览下面的相关文章,希望大家以后多多支持TOY模板网!

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处: 如若内容造成侵权/违法违规/事实不符,请点击违法举报进行投诉反馈,一经查实,立即删除!

领支付宝红包 赞助服务器费用

相关文章

  • Windows Server 2012 R2服务器Microsoft 消息队列远程代码执行漏洞CVE-2023-21554补丁KB5025288的安装及问题解决

    近日,系统安全扫描中发现Windows Server 2012 R2服务器存在Microsoft 消息队列远程代码执行漏洞。本文记录补丁安装中遇到的“此更新不适用于你的计算机”问题及解决办法。 一、问题描述: 1、系统安全扫描中发现Windows Server 2012 R2服务器存在Microsoft 消息队列远程代码执行漏洞,

    2024年02月10日
    浏览(53)
  • 极光Java 版本服务器端实现别名消息推送

    REST API 文档:

    2024年02月15日
    浏览(48)
  • Qt实现客户端与服务器消息发送

    里用Qt来简单设计实现一个场景,即: (1)两端:服务器QtServer和客户端QtClient (2)功能:服务端连接客户端,两者能够互相发送消息,传送文件,并且显示文件传送进度。 环境:VS20013 + Qt5.11.2 + Qt设计师 先看效果: 客户端与服务器的基本概念不说了,关于TCP通信的三次握

    2024年02月11日
    浏览(51)
  • Web服务器实现|基于阻塞队列线程池的Http服务器|线程控制|Http协议

    代码地址:WebServer_GitHub_Addr 摘要 本实验通过C++语言,实现了一个基于阻塞队列线程池的多线程Web服务器。该服务器支持通过http协议发送报文,跨主机抓取服务器上特定资源。与此同时,该Web服务器后台通过C++语言,通过原生系统线程调用 pthread.h ,实现了一个 基于阻塞队列

    2024年02月07日
    浏览(65)
  • 微信小程序实现订阅消息功能(Node服务器篇)

            * 源码已经上传到资源处,需要的话点击跳转下载 |  源码下载         在上一篇内容当中在微信小程序中实现订阅消息功能,都在客户端(小程序)中来实现的,在客户端中模拟了服务器端来进行发送订阅消息的功能,那么本篇就将上一篇内容中仅在客户端中实现

    2024年02月03日
    浏览(64)
  • 使用Java服务器实现UDP消息的发送和接收(多线程)

    在本篇博客中,我们将介绍如何使用Java服务器来实现UDP消息的发送和接收,并通过多线程的方式来处理并发请求。UDP(User Datagram Protocol)是一种无连接、不可靠的传输协议,适合于实时性要求高的应用场景,如实时游戏、语音通信等。 步骤: 首先,我们需要导入Java提供的

    2024年02月12日
    浏览(47)
  • SSE与WebSocket分别实现服务器发送消息通知(Golang、Gin)

    服务端推送,也称为消息推送或通知推送,是一种允许应用服务器主动将信息发送到客户端的能力,为客户端提供了实时的信息更新和通知,增强了用户体验。 服务端推送的背景与需求主要基于以下几个诉求: 实时通知:在很多情况下,用户期望实时接收到应用的通知,如

    2024年02月03日
    浏览(52)
  • 消息队列缓存,以蓝牙消息服务为例

    消息队列缓存,支持阻塞、非阻塞模式;支持协议、非协议模式 可自定义消息结构体数据内容 使用者只需设置一些宏定义、调用相应接口即可 这里我用蓝牙消息服务举例 有纰漏请指出,转载请说明。 学习交流请发邮件 1280253714@qq.com 队列采用\\\"生产者-消费者\\\"模式, 当接收数

    2024年02月07日
    浏览(46)
  • 消息队列 -提供上层服务接口

    我们之前已经将 数据库 的操作 和文件的操作 都完成了, 但是对于上层调用来说, 并不关心是于数据库中存储数据还是往文件中存储数据, 因此 我们提供一个类, 封装一下 上述俩个类中的操作, 并将这个类 提供给上层调用 由于之前我们已经分别测试过了,写入数据库与写入文件

    2024年02月14日
    浏览(38)
  • 【微服务】RabbitMQ&SpringAMQP消息队列

    🚩本文已收录至专栏:微服务探索之旅 👍希望您能有所收获 微服务间通讯有同步和异步两种方式: 同步通讯 :就像打电话,可以 立即得到响应 ,但是你却 不能跟多个人 同时通话。 异步通讯 :就像发消息,可以 同时与多个人 发送并接收消息,但是往往 响应会有延迟

    2024年02月16日
    浏览(61)

觉得文章有用就打赏一下文章作者

支付宝扫一扫打赏

博客赞助

微信扫一扫打赏

请作者喝杯咖啡吧~博客赞助

支付宝扫一扫领取红包,优惠每天领

二维码1

领取红包

二维码2

领红包