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

如何在Rust中借助Enum实现事件发射器与处理器(含Axum+Tokio)

Rust事件发送与处理机制实现方案(Axum+Tokio)

需求概述

需要实现基于Rust的事件处理系统,核心事件枚举定义如下:

enum Events {
    Message(String),
    Channel(i64),
    Disconnect(i16),
    Connect(JSON),
    HeartBeat
}

要求通过Axum的POST接口接收事件,匹配枚举类型触发对应处理逻辑,同时基于Tokio实现多线程并发处理。


实现步骤

1. 依赖准备

在Cargo.toml中添加必要依赖:

[dependencies]
axum = "0.7"
tokio = { version = "1.0", features = ["full"] }
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }

2. 事件枚举序列化适配

为了支持HTTP请求的JSON解析,给事件枚举添加Serde序列化/反序列化支持:

use serde::{Deserialize, Serialize};

#[derive(Debug, Deserialize, Serialize)]
#[serde(tag = "type", content = "data")] // 通过`type`字段区分事件变体
enum Events {
    Message(String),
    Channel(i64),
    Disconnect(i16),
    Connect(serde_json::Value), // 替换原JSON为serde_json::Value
    HeartBeat,
}

3. Axum接收POST事件

搭建Axum服务,接收事件并通过Tokio异步任务处理:

use axum::{extract::Json, routing::post, Router};
use std::sync::Arc;

#[tokio::main]
async fn main() {
    // 初始化日志
    tracing_subscriber::fmt::init();

    // 定义共享状态(可添加数据库连接、HTTP客户端等资源)
    let shared_state = Arc::new(SharedState {});

    // 构建路由
    let app = Router::new()
        .route("/events", post(receive_event))
        .with_state(shared_state);

    // 启动服务器
    let addr = std::net::SocketAddr::from(([127, 0, 0, 1], 3000));
    tracing::info!("listening on {}", addr);
    axum::Server::bind(&addr)
        .serve(app.into_make_service())
        .await
        .unwrap();
}

// 共享状态结构体,按需扩展字段
struct SharedState {
    // 示例:http_client: Arc<HttpClient>,
}

// 接收事件的接口handler
async fn receive_event(
    Json(event): Json<Events>,
    state: Arc<SharedState>,
) -> axum::response::Json<serde_json::Value> {
    // 用Tokio spawn异步任务,实现多线程并发处理
    tokio::spawn(handle_event(event, state.clone()));

    // 立即返回响应,不阻塞请求线程
    axum::response::Json(serde_json::json!({"status": "ok", "message": "event received"}))
}

4. 事件处理器实现

参考twilight-rs的匹配模式,编写事件处理逻辑:

use std::error::Error;

async fn handle_event(
    event: Events,
    state: Arc<SharedState>,
) -> Result<(), Box<dyn Error + Send + Sync>> {
    match event {
        Events::Message(msg) => {
            tracing::info!("Received message: {}", msg);
            // 自定义消息处理逻辑,比如存储、回复等
        }
        Events::Channel(channel_id) => {
            tracing::info!("Channel event triggered for ID: {}", channel_id);
            // 自定义频道逻辑
        }
        Events::Disconnect(code) => {
            tracing::warn!("Disconnected with code: {}", code);
            // 自定义断开连接处理逻辑
        }
        Events::Connect(data) => {
            tracing::info!("Connected with metadata: {}", data);
            // 自定义连接初始化逻辑
        }
        Events::HeartBeat => {
            tracing::debug!("Heartbeat received");
            // 自定义心跳响应逻辑
        }
    }

    Ok(())
}

Tokio多线程处理说明

完全可以通过Tokio实现多线程并发处理:

  • #[tokio::main]默认启用多线程运行时,会自动管理线程池资源
  • tokio::spawn会将任务调度到线程池执行,不会阻塞Axum的请求处理线程
  • 若需精细控制任务调度,可使用JoinSet管理任务集合,或自定义Tokio线程池配置

测试示例

用curl发送测试请求:

curl -X POST http://localhost:3000/events \
  -H "Content-Type: application/json" \
  -d '{"type": "Message", "data": "Hello, Rust!"}'

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 04:15:00