You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Rust下Actix WebSocket与RabbitMQ(lapin)双向消息代理集成咨询

Actix WebSocket服务集成RabbitMQ实现双向消息代理问题

我已经使用Rust的Actix框架编写了一个WebSocket服务,现在需要将RabbitMQ集成到该项目中,我找到了lapin crate用于RabbitMQ操作,但在将其与Actix框架集成时遇到了问题,我希望基于现有的WebSocket实现,将RabbitMQ的消息代理转发到客户端。

以下是我初步的实现方案,目前还处于非常早期的阶段,如果我的实现路径存在问题请告知我,方便我及时调整方向。


现有项目上下文

WebSocket Actor会和另一个存储Socket与房间信息的Actor通信,核心代码如下:

pub struct WSConn {
    pub id: Uuid,
    room_id: Uuid,
    hb: Instant,
    lobby_address: Addr<Lobby>,
}
impl WSConn {
    pub fn new(lobby: Addr<Lobby>, rabbit: Addr<MyRabbit>) -> Self {
        Self {
            id: Uuid::new_v4(),
            room_id: Uuid::nil(),
            hb: Instant::now(),
            lobby_address: lobby,
        }
    }

    fn hb(&self, context: &mut WebsocketContext<Self>) {
            /// HEARTBEAT CODE GOES HERE ///
    }
}

当Actor启动时,我会向Lobby发送一条连接消息,Lobby会处理将连接添加到自身结构体的全部逻辑。

impl Actor for WSConn {
    type Context = WebsocketContext<Self>;

    fn started(&mut self, ctx: &mut Self::Context) {
        info!("Starting hearbeat");
        self.hb(ctx);

        let wsserver_address = ctx.address();
        info!("A new client has connected with id {}", self.id);

        self.lobby_address.send(Connect {
            address: wsserver_address.recipient(),
            room_id: self.room_id,
            id: self.id
        }).into_actor(self).then(|res, _, ctx| {
            match res {
                Ok(_) => (),
                _ => ctx.stop()
            }
            fut::ready(())
        }).wait(ctx);
    }

    fn stopping(&mut self, _ctx: &mut Self::Context) -> Running {
        info!("Stopping actor");

        self.lobby_address.do_send(Disconnect {
            room_id: self.room_id,
            id: self.id
        });
        Running::Stop
    }
}

Lobby的定义如下:

type Socket = Recipient<WSServerMessage>;

pub struct Lobby {
    pub sessions: DashMap<Uuid, Socket>,//self id to self
    pub rooms: DashMap<Uuid, DashSet<Uuid>>,//room id  to list of users id
}

Lobby的完整代码量较大,此处不贴出,如果需要查看可以提供相关代码。客户端连接后会被分配到默认房间,客户端发送的消息会由StreamHandler处理。

impl StreamHandler<Result<ws::Message, ws::ProtocolError>> for WSConn {
    fn handle(&mut self, item: Result<Message, ProtocolError>, ctx: &mut Self::Context) {
        match item.unwrap() {
            ws::Message::Binary(bin) => ctx.binary(bin),
            ws::Message::Ping(bin) => {
                self.hb = Instant::now();
                ctx.pong(&bin);
            }
            ws::Message::Pong(_) => self.hb = Instant::now(),
            ws::Message::Close(reason) => {
                ctx.close(reason);
                ctx.stop();
            }
            ws::Message::Text(text) => {
                let command = serde_json::from_str::<Command>(&text)
                    .expect(&format!("Can't parse message {}", &text));
                info!("{:?}", command);
                if command.command.starts_with("/") {
                    info!("This is a {} request", command.command);
                    match command.command.as_ref() {
                        "/join" => {
                            info!("Join Room {}", command.payload);
                            let uid = Uuid::from_str(command.payload.as_str().unwrap()).expect("Can't parse message {} to uuid");
                            self.lobby_address.send(Join {
                                current_room: self.room_id,
                                room_id: uid,
                                id: self.id
                            }).into_actor(self).then(|res, _, ctx| {
                                match res {
                                    Ok(_) => (),
                                    _ => ctx.stop()
                                }
                                fut::ready(())
                            }).wait(ctx);
                            self.room_id = uid;
                        },
                        _ => ()
                    }
                } else {
                    info!("Text is {}", text);
                    self.lobby_address.do_send(ClientActorMessage {
                        id: self.id,
                        msg: command,
                        room_id: self.room_id,
                    });
                }
            }
            _ => {
                info!("Something weird happened. Closing");
                self.lobby_address.do_send(Disconnect {
                    room_id: self.room_id,
                    id: self.id
                });
                ctx.stop();}
        }
    }
}

如上述代码所示,客户端发送携带/join命令和有效uuidv4 payload的消息即可加入对应房间,我限制了客户端同一时间只能加入一个房间,加入新房间时会自动退出上一个房间。


目前RabbitMQ集成实现

我计划使用连接池维护RabbitMQ连接,再基于连接创建Channel,首先我定义了存储连接池的结构体:

use actix::{Actor, Context, Handler, StreamHandler};
use deadpool_lapin::{Config, Pool, Runtime};
use deadpool_lapin::lapin::Error;
use lapin::message::Delivery;

use crate::lapin_server::messages::CreateChannel;

pub struct MyRabbit {
    pub pool: Pool,
}

impl MyRabbit {
    pub fn new() -> Self {
        let mut cfg = Config::default();
        cfg.url = Some("amqps://ghqcmhat:KbhPAA309QRg7TjdgFEV14pQRheoh44P@codfish.rmq.cloudamqp.com/ghqcmhat".into());
        let new_pool = cfg.create_pool(Some(Runtime::Tokio1)).expect("Can't create pool");
        MyRabbit {
            pool: new_pool
        }
    }
}

impl Actor for MyRabbit {
    type Context = Context<Self>;
}

将其封装为Actor后,我可以在启动WebSocket服务的同时启动该Actor,这部分逻辑在main函数中实现:

#[actix_web::main]
async fn main() -> std::io::Result<()>{
    std::env::set_var("RUST_LOG", "actix_web=info,info");
    env_logger::init();
    // Start Lobby actor and get his address
    let websocket_lobby = Lobby::default().start();
    let rabbit = MyRabbit::new().start();

    let application_data = web::Data::new(Appdata::new());

    info!("Starting server on 127.0.0.1:8080");
    let server = HttpServer::new(move || {
        App::new()
            .wrap(Logger::default())
            .route("/ws/",web::get().to(websocket_handler))
            .app_data(Data::new(websocket_lobby.clone()))
            .app_data(Data::new(rabbit.clone()))
            .app_data(application_data.clone())
    });
    server.bind("127.0.0.1:8080")?.run().await
}

为了适配新的Actor,我在请求处理器中新增了对应参数:

pub async fn websocket_handler(request: HttpRequest, stream: web::Payload, srv: Data<Addr<Lobby>>, rab: Data<Addr<MyRabbit>>, data: Data<Appdata>) -> Result<HttpResponse, Error> {
    let mut counter = data.counter.lock().unwrap();
    counter.add_assign(1);
    info!("This is request # {}", counter);
    let ws = WSConn::new(srv.get_ref().clone(), rab.get_ref().clone());
    let response = ws::start(ws, &request, stream);
    debug!("Response: {:?}", &response);
    response
}

调整后的WSConn结构体如下:

pub struct WSConn {
    pub id: Uuid,
    room_id: Uuid,
    hb: Instant,
    lobby_address: Addr<Lobby>,
    rabbit_address: Addr<MyRabbit>
}

我知道要从RabbitMQ的主题消费消息,需要exchange名称、类型、routing key和队列名称,因此我也定义了对应的结构体:

pub struct Channel {
    pub queue_name: String,
    pub exchange_name: String,
    pub exchange_type: ExchangeKind,
    pub routing_key: String
}

impl Default for Channel {
    fn default() -> Self {
        Self {
            queue_name: "".to_string(),
            exchange_name: "".to_string(),
            exchange_type: Default::default(),
            routing_key: "".to_string()
        }
    }
}

待解决问题与最终目标

目前我卡在了后续实现环节:

  • 不确定deadpool_lapin是否是适用于该场景的crate
  • 不清楚如何适配lapin官方文档中的示例,官方示例使用async_global_executor::block_on,且通过async_global_executor::spawn创建新线程来消费消息,不知道如何和Actix的Actor模型结合

我最终要实现的效果是WebSocket与RabbitMQ双向消息代理:客户端连接WebSocket后发送如下格式的消息:

{
    "command": "SUBSCRIBE",
    "payload": "topic_name"
}

就可以收到对应RabbitMQ主题下发布的所有消息,发送UNSUBSCRIBE命令即可取消订阅。恳请各位提供相关实现建议,如果需要更多信息可以告知我。


内容的提问来源于stack exchange,提问作者Fabrex

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.29 08:18:03