【问题标题】:Asynchronous web socket application server with two event loops具有两个事件循环的异步 Web 套接字应用程序服务器
【发布时间】:2015-12-09 19:22:41
【问题描述】:

我正在尝试制作一个分布式 RPC 类型的 Web 应用程序,它使用 websockets 作为其主界面。我想使用排队系统(如 RabbitMQ)来分配通过 websocket 连接请求的昂贵作业。

基本上,流程是这样的:

  1. 客户端通过 websocket 连接向服务器发送作业
  2. 服务器会将此消息发送到 RabbitMQ 交换以由工作人员处理
  3. worker 将执行作业并将作业结果添加到响应队列中
  4. 服务器会检查响应队列并通过 websocket 连接将作业结果发送回客户端。

据我所知,在服务器上我需要两个共享内存的事件循环。 websocket 服务器需要监听传入的作业,而 RabbitMQ 消费者需要监听作业结果以发送回客户端。

我可以在这里使用哪些合适的技术?我考虑了以下几点:

  • 多线程应用程序并在每个线程上启动一个事件循环
  • 通过 shm(共享内存)使用两个进程
  • 使用两个通过套接字进行通信的进程(一个 unix 套接字,甚至可能将工作人员设置为特殊的 websocket 客户端)
  • 挂钩到 websocket 服务器的事件循环以检查结果队列

我对 websockets 和分布式计算都是新手,所以我真的不知道其中哪一个(或者我没有想到的)最适合我。

【问题讨论】:

    标签: asynchronous websocket ipc message-queue distributed-computing


    【解决方案1】:

    据我所知,在服务器上我需要两个共享内存的事件循环。 websocket 服务器需要监听传入的作业,而 RabbitMQ 消费者需要监听作业结果以发送回客户端。

    由于您可以让多个客户端同时发送作业,因此您需要一个多线程服务器。除非您的应用程序会为每个客户端处理客户端。现在有多种方法可以实现多线程服务器,每种方法都有自己的优点/缺点。看看多线程通过:

    1. 每个请求一个线程(+:吞吐量可能最大化,-:线程的创建成本很高,必须管理并发)
    2. 每个客户端一个线程(+:更少的线程管理开销,-:不能扩展到许多连接,仍然管理并发)
    3. 一个线程池(+:避免创建线程的开销,最多可扩展至 N 个并发连接(N = 线程池的大小),-:管理 N 个之间的并发线程)

    您可以选择上述方法之一(我会为每个客户端选择一个线程,因为它相对容易实现,并且您拥有数万个客户端的机会相对较小) .

    请注意,这是一种多线程方法,不是事件驱动方法!但是由于您不限于一个线程(在这种情况下,它应该是事件驱动的,以便能够处理多个客户端“并发”)我不会选择那个选项,因为它更多难以实施。 (程序员有时会谈到事件驱动方法中的“回调地狱”)。

    这就是我将如何实现它(每个客户端一个线程,Java)

    基本上,流程是这样的:

    1. 客户端通过 websocket 连接向服务器发送作业

    服务器部分:

    public class Server {
        private static ServerSocket server_skt;
        private static ... channel; // channel to communicate with the rabbitMQ distributed priority queue.
    
        // Constructor
        Server(int port) {
            server_skt = new ServerSocket(port);
            
            /*
             * Set up connection with the distributed queue
             * channel = ...;
             */
        }
    
        public static void main(String argv[]) {
            Server server = new Server(5555); // Make server instance
    
            while(true) {
                // Always waiting for new clients to connect
                try {
                    System.out.println("Waiting for a client to connect...");
                    // Spawn new thread for communication with client (hence one thread per client approach)
                    new CommunicationThread(server_skt.accept(), server.channel).start(); // Will listen for new jobs and send them
                } catch(IOException e) {
                    System.out.println("Exception occured :" + e.getStackTrace());
                }
            }
        }
    }
    
    1. 服务器会将此消息发送到 RabbitMQ 交换以由工作人员处理
    2. ...
    3. 服务器会检查响应队列并通过 websocket 连接将作业结果发送回客户端。
    public class CommunicationThread extends Thread {
        private Socket client_socket;
        private InputStream client_in;
        private OutputStream client_out;
        private ... channel; // Channel to communicate with rabbitMQ
        private ... resultQueue;
    
        public CommunicationThread(Socket socket, ... channel) { // replace ... by type of the rabbitMQ channel
            try {
                this.client_socket = socket;
                this.client_in = client_socket.getInputStream();
                this.client_out = client_socket.getOutputStream(); 
                this.channel = channel;
                this.resultQueue = ...;
    
                System.out.println("Client connected : " + client_socket.getInetAddress().toString());
            } catch(IOException e) {
                System.out.println("Could not initialize communication properly. -- CommunicationThread.\n");
            }
        }
        
        public yourJobType readJob() {
            // Read input from client (e.g. read a String from "client_in")
            // Make a job from it (e.g. map String to a job)
            // return the job
        }
    
        @Override
        public void run() {
            while(active) {
                
                /*
                 * Always listen for incoming jobs (sent by client) and for results (to be sent back to client)
                 */
                
                // Read client input (only if available, else it would be blocking!)
                if(client_in.available() > 0) {
                    yourJobType job = readJob();
                    channel.basicPublish(...); // Send job to rabbitMQ
                }
                
                /* Check result queue (THIS is why reading client input MUST be NON-BLOCKING! Else while loop could be blocked on reading input
                 * and the result queue won't be checked until next job arrives)
                 */
                
                ResultType next_result = resultQueue.poll(); // Could be "null" if the queue is empty
                if(next_result != null) {
                    // There is a result
                    client_out.write(next_result.toByteArray());
                    client_out.flush();
                }
            }
            
            client_in.close();
            client_out.close();
        }
    }
    

    请注意,从结果队列中读取时,请务必仅读取该客户端发送的作业的结果

    如果您有一个包含(所有客户端的)作业结果的结果队列,并且您检索到类似于上面代码中的结果,那么该结果可能是另一个客户端的作业的结果,因此将结果发送回错误的客户。

    要解决此问题,您可以poll() 带有过滤器和通配符 (*) 的结果队列,或者为每个客户端设置一个结果队列,因此知道从我们的队列中检索到的结果将被发送到相应的客户端。

    (*) :您可以为每个客户分配一个 ID。当从客户端接收作业时,将作业与客户端 ID 配对(例如,在元组 中)并将其放入队列中。并对结果执行相同的操作(将结果与客户端 ID 配对并将其放入结果队列中)。然后在CommunicationThreadrun() 方法中,您将不得不只为 .

    重要提示:您还必须为每项工作分配一个 ID!因为发送作业 A 然后发送作业 B 并不能保证作业 A 的结果会在作业 B 的结果之前出现。(作业 B 可能比作业 A 耗时更少,因此可以在作业之前将结果发送回客户端A 的结果)。

    (PS : 看你如何实现workers(由服务器执行,每个worker一个线程执行?还是由其他进程执行?))


    上述答案是一种可能的多线程解决方案。我只讨论了服务器部分,客户端部分应该发送作业并等待结果(如何实现这取决于您的目标,客户端首先发送所有作业然后接收每个作业的结果还是可以这是混合的吗?)。

    还有其他方法可以实现它,但对于分布式计算的初学者来说,我认为这是最简单的解决方案(使用线程池,......会变得更棘手)。

    【讨论】:

    • 谢谢!这非常有帮助(一年多后;))
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-04-15
    • 2014-02-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多