并发服务器:第7部分 - Rust
摘要
本文是关于并发服务器系列文章的一部分,介绍了如何使用Rust实现并发网络服务器,涵盖了顺序、线程和事件驱动方法,并提供了代码示例。
<p>这是关于编写并发网络服务器系列文章的第七部分。在这一部分中,我们将讨论如何在Rust编程语言中解决前面部分所描述的挑战。</p>
<p>本系列的所有文章:</p>
<ul class="simple">
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-1-introduction/">第1部分 - 介绍</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-2-threads/">第2部分 - 线程</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-3-event-driven/">第3部分 - 事件驱动</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-4-libuv/">第4部分 - libuv</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-5-redis-case-study/">第5部分 - Redis案例研究</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2018/concurrent-servers-part-6-callbacks-promises-and-asyncawait/">第6部分 - 回调、Promise和async/await</a></li>
<li>第7部分 - Rust(本部分)</li>
</ul>
<p>自前几部分发布以来,已经过去了几年。我最近重新审视了它们,以确保所呈现的信息仍然相关,并且所有代码示例都能在现代工具链中构建和运行。我强烈建议在阅读本部分之前先回顾前几部分。</p>
<p>本文假设读者对Rust编程语言有基本的了解。我们只会在遇到入门书籍或教程中不常见的代码时解释Rust的构造。</p>
<div class="section" id="setting-the-baseline-a-sequential-state-machine-server">
<h2>设定基准 - 顺序状态机服务器</h2>
<p>系列的前几部分重点介绍了一个实现简单状态机协议的套接字服务器。有关协议的完整描述,请参见<a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-1-introduction/">第1部分</a>。让我们首先展示如何在基本的顺序Rust服务器中实现该协议:</p>
<div class="highlight"><pre><span></span><span class="k">use</span><span class="w"> </span><span class="n">async_socket_server</span>::<span class="n">serve_connection</span><span class="p">;</span><span class="w"></span>
<span class="k">use</span><span class="w"> </span><span class="n">std</span>::<span class="n">net</span>::<span class="n">TcpListener</span><span class="p">;</span><span class="w"></span>
<span class="k">fn</span> <span class="nf">main</span><span class="p">()</span><span class="w"> </span>-> <span class="nc">std</span>::<span class="n">io</span>::<span class="nb">Result</span><span class="o"><</span><span class="p">()</span><span class="o">></span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="kd">let</span><span class="w"> </span><span class="n">port</span><span class="w"> </span><span class="o">=</span><span class="w"> </span><span class="k">match</span><span class="w"> </span><span class="n">std</span>::<span class="n">env</span>::<span class="n">args</span><span class="p">().</span><span class="n">nth</span><span class="p">(</span><span class="mi">1</span><span class="p">)</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="nb">Some</span><span class="p">(</span><span class="n">s</span><span class="p">)</span><span class="w"> </span><span class="o">=></span><span class="w"> </span><span class="n">s</span><span class="p">,</span><span class="w"></span>
<span class="w"> </span><span class="nb">None</span><span class="w"> </span><span class="o">=></span><span class="w"> </span><span class="s">"9090"</span><span class="p">.</span><span class="n">to_string</span><span class="p">(),</span><span class="w"></span>
<span class="w"> </span><span class="p">};</span><span class="w"></span>
<span class="w"> </span><span class="kd">let</span><span class="w"> </span><span class="n">addr</span><span class="w"> </span><span class="o">=</span><span class="w"> </span><span class="fm">format!</span><span class="p">(</span><span class="s">"127.0.0.1:{port}"</span><span class="p">);</span><span class="w"></span>
<span class="w"> </span><span class="kd">let</span><span class="w"> </span><span class="n">listener</span><span class="w"> </span><span class="o">=</span><span class="w"> </span><span class="n">TcpListener</span>::<span class="n">bind</span><span class="p">(</span><span class="n">addr</span><span class="p">)</span><span class="o">?</span><span class="p">;</span><span class="w"></span>
<span class="w"> </span><span class="fm">println!</span><span class="p">(</span><span class="s">"Serving on port {port}"</span><span class="p">);</span><span class="w"></span>
<span class="w"> </span><span class="k">loop</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="kd">let</span><span class="w"> </span><span class="p">(</span><span class="n">stream</span><span class="p">,</span><span class="w"> </span><span class="n">addr</span><span class="p">)</span><span class="w"> </span><span class="o">=</span><span class="w"> </span><span class="n">listener</span><span class="p">.</span><span class="n">accept</span><span class="p">()</span><span class="o">?</span><span class="p">;</span><span class="w"></span>
<span class="w"> </span><span class="fm">println!</span><span class="p">(</span><span class="s">"connection received from {}"</span><span class="p">,</span><span class="w"> </span><span class="n">addr</span><span class="p">);</span><span class="w"></span>
<span class="w"> </span><span class="k">if</span><span class="w"> </span><span class="kd">let</span><span class="w"> </span><span class="nb">Err</span><span class="p">(</span><span class="n">e</span><span class="p">)</span><span class="w"> </span><span class="o">=</span><span class="w"> </span><span class="n">serve_connection</span><span class="p">(</span><span class="n">stream</span><span class="p">)</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="fm">eprintln!</span><span class="p">(</span><span class="s">"error serving connection: {}"</span><span class="p">,</span><span class="w"> </span><span class="n">e</span><span class="p">);</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"> </span><span class="k">else</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="fm">println!</span><span class="p">(</span><span class="s">"peer done {addr}"</span><span class="p">);</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"></span>
<span class="p">}</span><span class="w"></span>
</pre></div>
<p>函数<tt class="docutils literal">serve_connection</tt>定义如下:</p>
<div class="highlight"><pre><span></span><span class="k">pub</span><span class="w"> </span><span class="k">enum</span> <span class="nc">ProcessingState</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="n">WaitForMsg</span><span class="p">,</span><span class="w"></span>
<span class="w"> </span><span class="n">InMsg</span><span class="p">,</span><span class="w"></span>
<span class="p">}</span><span class="w"></span>
<span class="k">pub</span><span class="w"> </span><span class="k">fn</span> <span class="nf">serve_connection</span><span class="p">(</span><spa
查看缓存全文
缓存时间: 2026/08/16 03:26
# 并发服务器:第7部分 - Rust
来源:https://eli.thegreenplace.net/2026/concurrent-servers-part-7-rust
这是关于编写并发网络服务器系列文章的第七部分。在本部分中,我们将探讨前几部分中描述的挑战如何在 Rust 编程语言中得到解决。
本系列所有文章:
- 第1部分 - 简介 (https://eli.thegreenplace.net/2017/concurrent-servers-part-1-introduction/)
- 第2部分 - 线程 (https://eli.thegreenplace.net/2017/concurrent-servers-part-2-threads/)
- 第3部分 - 事件驱动 (https://eli.thegreenplace.net/2017/concurrent-servers-part-3-event-driven/)
- 第4部分 - libuv (https://eli.thegreenplace.net/2017/concurrent-servers-part-4-libuv/)
- 第5部分 - Redis 案例研究 (https://eli.thegreenplace.net/2017/concurrent-servers-part-5-redis-case-study/)
- 第6部分 - 回调、Promise 和 async/await (https://eli.thegreenplace.net/2018/concurrent-servers-part-6-callbacks-promises-and-asyncawait/)
- 第7部分 - Rust(本文)
自前几部分发布以来,已过去多年。我最近重新审视了它们,以确保呈现的信息仍然相关,并且所有代码示例都能使用现代工具链构建和运行。我强烈建议在阅读本文之前回顾前面的部分。
本文假定您对 Rust 编程语言有基本的熟悉。它只会在我们遇到不会出现在入门书籍或教程中的代码时,才解释 Rust 的构造。
## 设定基线 - 一个顺序状态机服务器
本系列的前几部分重点关注一个实现简单状态机协议的套接字服务器。完整协议描述请参见[第1部分](https://eli.thegreenplace.net/2017/concurrent-servers-part-1-introduction/)。
让我们首先展示如何在基本的顺序 Rust 服务器中实现此协议:
```rust
use async_socket_server::serve_connection;
use std::net::TcpListener;
fn main() -> std::io::Result<()> {
let port = match std::env::args().nth(1) {
Some(s) => s,
None => "9090".to_string(),
};
let addr = format!("127.0.0.1:{port}");
let listener = TcpListener::bind(addr)?;
println!("Serving on port {port}");
loop {
let (stream, addr) = listener.accept()?;
println!("connection received from {}", addr);
if let Err(e) = serve_connection(stream) {
eprintln!("error serving connection: {}", e);
} else {
println!("peer done {addr}");
}
}
}
```
其中 `serve_connection` 函数定义如下:
```rust
pub enum ProcessingState {
WaitForMsg,
InMsg,
}
pub fn serve_connection(mut stream: TcpStream) -> std::io::Result<()> {
stream.write_all(b"*")?;
let mut state = ProcessingState::WaitForMsg;
let mut buf = [0u8; 1024];
loop {
let n = stream.read(&mut buf)?;
if n == 0 {
// 客户端关闭了连接。
break;
}
for byte in &buf[..n] {
match state {
ProcessingState::WaitForMsg => {
if *byte == b'^' {
state = ProcessingState::InMsg;
}
}
ProcessingState::InMsg => {
if *byte == b'$' {
state = ProcessingState::WaitForMsg;
} else {
let newbyte = byte.wrapping_add(1);
stream.write_all(&[newbyte])?;
}
}
}
}
}
Ok(())
}
```
作为提醒,此服务器版本是*顺序*的,因为它一次只接受一个客户端;主循环在 `serve_connection` 上阻塞,直到其完成(客户端关闭连接),然后才返回接受下一个客户端。
## 每个客户端一个线程
显然,逐个处理客户端是不可行的。在[第2部分](https://eli.thegreenplace.net/2017/concurrent-servers-part-2-threads/)中,我们讨论了使用操作系统线程并发处理客户端的方法。让我们从 Rust 中无界的“每个客户端一个线程”解决方案开始:
```rust
use async_socket_server::serve_connection;
use std::{net::TcpListener, thread};
fn main() -> std::io::Result<()> {
let port = match std::env::args().nth(1) {
Some(s) => s,
None => "9090".to_string(),
};
let addr = format!("127.0.0.1:{port}");
let listener = TcpListener::bind(addr)?;
println!("Serving on port {port}");
loop {
let (stream, addr) = listener.accept()?;
println!("connection received from {}", addr);
let res = thread::Builder::new().spawn(move || {
if let Err(e) = serve_connection(stream) {
eprintln!("error serving connection: {}", e);
} else {
println!("peer done {addr}");
}
});
if let Err(e) = res {
eprintln!("error spawning thread: {}", e);
}
}
}
```
`spawn` 方法返回一个 `Result<JoinHandle<()>>`;成功时,我们允许句柄在循环迭代结束时被丢弃。在 Rust 中,这会*分离*该线程;我们实际上并不等待它完成。对于我们的代码示例来说这是合理的,因为循环是*无限*的;它永远不会终止。失控线程的可能性只是第2部分中讨论的无界线程方法的问题之一。解决方案是使用固定线程池。
## 线程池
在深入代码之前,先简要说明一下设计:线程池是一组固定的线程,它们等待“作业”并将其处理到完成。在我们的例子中,“作业”是为特定客户端执行 `serve_connection`。
有许多方法可以实现线程池;对于我们的用例,我选择了一组线程,它们都获得一个共享的*通道*,主线程向该通道发送作业。工作线程从通道中取出下一个作业,将其服务到完成,然后返回等待下一个作业。代码如下所示:
```rust
struct Job {
stream: TcpStream,
addr: SocketAddr,
}
fn worker(receiver: Receiver<Job>) {
while let Ok(job) = receiver.recv() {
if let Err(e) = serve_connection(job.stream) {
eprintln!("error serving connection from {}: {}", job.addr, e);
} else {
println!("peer done {}", job.addr);
}
}
}
```
`Receiver` 是什么?它是来自 `crossbeam_channel` crate 的类型:
```rust
use crossbeam_channel::{Receiver, bounded};
```
Rust 内置的 `std` 中的通道是*MPSC*(多生产者、单消费者),但我们的作业队列需要一个支持多个消费者(工作线程)的通道。虽然 `std` 确实有 [mpmc](https://doc.rust-lang.org/std/sync/mpmc/fn.channel.html),但这是一个实验性 API,在撰写本文时仅在 nightly 版本中可用。因此,我选择包含 `crossbeam_channel` crate,它为此示例提供了经过良好测试的 MPMC 通道 [\[1\]](https://eli.thegreenplace.net/2026/concurrent-servers-part-7-rust#footnote-1)。
这是主函数:
```rust
use async_socket_server::serve_connection;
use crossbeam_channel::{Receiver, bounded};
use std::{
io,
net::{SocketAddr, TcpListener, TcpStream},
thread,
};
const NUM_WORKERS: usize = 256;
const JOB_QUEUE_CAPACITY: usize = NUM_WORKERS;
fn main() -> io::Result<()> {
let port = match std::env::args().nth(1) {
Some(s) => s,
None => "9090".to_string(),
};
let addr = format!("127.0.0.1:{port}");
let listener = TcpListener::bind(addr)?;
println!("Serving on port {port}");
let (tx, rx) = bounded(JOB_QUEUE_CAPACITY);
for _ in 0..NUM_WORKERS {
let receiver = rx.clone();
thread::Builder::new().spawn(move || worker(receiver))?;
}
drop(rx);
loop {
let (stream, addr) = listener.accept()?;
println!("connection received from {addr}");
// 队列已满时会阻塞此循环,
// 使其在工作线程可用之前无法接受更多连接。
if tx.send(Job { stream, addr }).is_err() {
return Err(io::Error::other("all connection workers stopped"));
}
}
}
```
请注意,我们的作业通道是*有界*的——它具有固定大小。这有助于自然地实现*背压*机制——如果太多客户端连接,后续客户端将不得不等待——主循环在 `tx.send` 上阻塞,并且在作业从通道中清除之前不会接受更多的套接字连接。
## 异步、事件驱动服务器
在本系列的第4、5和6部分中,我们讨论了*事件驱动*或*异步*服务器。让我们看看如何在 Rust 中实现。具体来说,[第6部分](https://eli.thegreenplace.net/2018/concurrent-servers-part-6-callbacks-promises-and-asyncawait/)介绍了从回调到 Promise 再到 async/await 机制的渐进过程;Rust 支持所有这些,并且——正如您所预期的——现代代码通常使用 async/await 编写,同时将所有 Promise 的细节(在 Rust 中称为 *futures*)隐藏在下面。
闲话少说,这是用异步 Rust 编写的简单状态机协议:
```rust
use async_socket_server::ProcessingState;
use tokio::io::{self, AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
#[tokio::main]
async fn main() -> io::Result<()> {
let port = match std::env::args().nth(1) {
Some(s) => s,
None => "9090".to_string(),
};
let addr = format!("127.0.0.1:{port}");
let listener = TcpListener::bind(addr).await?;
println!("Serving on port {port}");
loop {
let (socket, addr) = listener.accept().await?;
println!("connection received from {addr:?}");
tokio::spawn(async move {
if let Err(e) = serve_connection_async(socket).await {
eprintln!("error serving connection from {addr:?}: {e}");
} else {
println!("peer done {addr:?}");
}
});
}
}
```
Rust 对异步编程采取了一种有趣的方法:它在核心语言中支持一些基本构建块(如 futures 以及 `async` 和 `await` 关键字),但将实际的异步引擎实现(实现事件循环的东西)留给了外部 crate。
目前 Rust 中异步编程最流行的 crate 是 [tokio](https://tokio.rs/),所以我们在这里使用它。在阅读了第6部分的 JavaScript 代码后,上面的 Rust 代码片段应该会显得相当熟悉,只是可能需要解释一下显式的 tokio 任务“spawn”。代码不是将回调入队到 `listener.accept` 返回的连接上,而是生成一个 tokio 任务,这可以看作是一个[绿色线程](https://en.wikipedia.org/wiki/Green_thread),因此使用了类似的术语 [\[2\]](https://eli.thegreenplace.net/2026/concurrent-servers-part-7-rust#footnote-2)。
这些任务不能发起阻塞调用;因此,它们应该使用 tokio 的 I/O 工具而不是通常的、阻塞的 `std` 工具。实际上,我们必须实现 `serve_connection` 的异步版本才能使其工作:
```rust
async fn serve_connection_async(mut stream: tokio::net::TcpStream) -> std::io::Result<()> {
stream.write_all(b"*").await?;
let mut state = ProcessingState::WaitForMsg;
let mut buf = [0u8; 1024];
loop {
let n = stream.read(&mut buf).await?;
if n == 0 {
// 客户端关闭了连接。
break;
}
for byte in &buf[..n] {
match state {
ProcessingState::WaitForMsg => {
if *byte == b'^' {
state = ProcessingState::InMsg;
}
}
ProcessingState::InMsg => {
if *byte == b'$' {
state = ProcessingState::WaitForMsg;
} else {
let newbyte = byte.wrapping_add(1);
stream.write_all(&[newbyte]).await?;
}
}
}
}
}
Ok(())
}
```
请注意这段代码与之前的 `serve_connection` 有多么相似;唯一真正的区别是套接字读写上的 `await` 调用 [\[3\]](https://eli.thegreenplace.net/2026/concurrent-servers-part-7-rust#footnote-3),以及涉及的类型。例如,代替同步示例中使用的 `std::net::TcpStream`,这里我们使用 `tokio::net::TcpStream`。Tokio 有一个底层依赖项叫做 [mio](https://github.com/tokio-rs/mio),用于处理各种 I/O 的非阻塞 API。它包装了操作系统特定的事件循环,如 `epoll`,以高效地完成此操作。
## 异步素性测试服务器
虽然本系列大部分时间都使用简单的状态机服务器作为驱动示例,但第6部分将重点转移到了一个用于素性测试、模拟长时间计算任务的服务器上。让我们看看如何在 Rust 中使用 tokio 实现这一点:
```rust
use tokio::io::{self, AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
#[tokio::main]
async fn main() -> io::Result<()> {
let port = match std::env::args().nth(1) {
Some(s) => s,
None => "8070".to_string(),
};
let addr = format!("127.0.0.1:{port}");
let listener = TcpListener::bind(addr).await?;
println!("Serving on port {port}");
loop {
let (socket, addr) = listener.accept().await?;
println!("connection received from {addr:?}");
tokio::spawn(async move {
if let Err(e) = serve_client(socket).await {
eprintln!("error serving connection from {addr:?}: {e}");
} else {
println!("peer done {addr:?}");
}
});
}
}
async fn serve_client(mut stream: tokio::net::TcpStream) -> std::io::Result<()> {
let mut buf = [0u8; 1024];
loop {
let n = stream.read(&mut buf).await?;
if n == 0 {
// 客户端关闭了连接。
return Ok(());
}
// 将读取的缓冲区解析为 u64。
let num = std::str::from_utf8(&buf[..n])
.map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?
.trim()
.parse::<u64>()
.map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
let answer = if isprime(num, true) {
"prime"
} else {
"composite"
};
stream
.write_all((answer.to_string() + "\n").as_bytes())
.await?;
}
}
```
这段代码在概念上与之前的片段非常相似;`isprime` 如下:
```rust
// 检查 n 是否为素数,返回布尔值。delay 参数是可选的;
// 如果为 true,函数将在计算答案之前阻塞 n 毫秒。这用于模拟长时间运行的计算。
fn isprime(n: u64, delay: bool) -> bool {
if delay {
std::thread::sleep(std::time::Duration::from_millis(n));
}
if n < 2 {
return false;
}
if n % 2 == 0 {
return n == 2;
}
let mut r = 3;
while r * r <= n {
if n % r == 0 {
return false;
}
r += 2;
}
true
}
```
请注意,这个示例演示了一个可能阻塞的任务(在这种情况下通过 `sleep` 模拟)。正如 [tokio 文档](https://docs.rs/tokio/latest/tokio/task/index.html#blocking-and-yielding)所解释的那样,这在异步上下文中可能是一个问题。一个潜在的解决方案是将阻塞任务分派到单独的线程池,并使用 [tokio channels](https://docs.rs/tokio/1.53.1/tokio/sync/) 与之通信;这类似于我们在上面线程池示例中采取的方法。
第6部分还包括此服务器的一个版本,该版本在本地 Redis 实例上缓存数据;目标是演示当添加额外的回调层时事件驱动代码的复杂性,以及 async/await 如何帮助缓解这种情况。这是此服务器的 Rust 版本,使用了 `redis` crate(该 crate 显式启用了 `tokio` 组件以支持异步调用):
```rust
use redis::AsyncTypedCommands;
use redis::aio::MultiplexedConnection;
use tokio::io::{self, AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
#[derive(Clone)]
struct AppState {
redis_connection: MultiplexedConnection,
}
const REDIS_URL: &str = "redis://127.0.0.1";
#[tokio::main]
async fn main() -> io::Result<()> {
let port = match std::env::args().nth(1) {
Some(s) => s,
None => "8070".to_string(),
};
let addr = format!("127.0.0.1:{port}");
let app_state = AppState {
redis_connection: {
redis::Client::open(REDIS_URL)
.map_err(io::Error::other)?
.get_multiplexed_async_connection()
.await
.map_err(io::Error::other)?
},
};
let listener = TcpListener::bind(addr).await?;
println!("Serving on port {port}");
loop {
let (socket, addr) = listener.accept().await?;
println!("connection received from {addr:?}");
let app_state = app_state.clone();
tokio::spawn(async move {
if let Err(e) = serve_client(socket, app_state).await {
eprintln!("error serving connection from {addr:?}: {e}");
} else {
println!("peer done {addr:?}");
}
});
}
}
async fn serve_client(
mut stream: tokio::net::TcpStream,
mut app_state: AppState,
) -> std::io::Result<()> {
let mut buf = [0u8; 1024];
loop {
let n = stream.read(&mut buf).await?;
if n == 0 {
// 客户端关闭了连接。
return Ok(());
}
// 将读取的缓冲区解析为 u64。
let num = std::str::from_utf8(&buf[..n])
.map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?
.trim()
.parse::<u64>()
.map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
// 首先尝试 Redis 缓存。如果在缓存中找到,发送缓存的答案。
let cachekey = format!("primecache:{num}");
match app_state.redis_connection.get(cachekey.clone()).await {
Ok(Some(cached)) => {
stream.write_all((cached.clone() + "\n").as_bytes()).await?;
println!("cached num {num} is {cached}");
continue;
}
Ok(None) => {
// 未在缓存中找到,继续计算...
}
Err(e) => {
eprintln!("redis get error: {}", e);
// 缓存出错,继续计算...
}
}
let answer = if isprime(num, true) {
"prime"
} else {
"composite"
};
// 将结果存储在缓存中。
let _ = app_state.redis_connection
.set(cachekey, answer)
.await;
stream
.write_all((answer.to_string() + "\n").as_bytes())
.await?;
}
}
```
此代码与之前的素性测试服务器片段非常相似,但现在包括与 Redis 缓存交互的逻辑。关键部分是使用 `tokio::spawn` 处理每个连接,以及在异步上下文中使用 Redis 客户端(`MultiplexedConnection`)。
相似文章
谁在运行你的 Rust Future?动手实践入门异步 Rust
这是一套动手实践教程系列,旨在弥合理解异步 Rust 内部机制(Future、poll、Pin、执行器)与使用 Tokio 部署实际异步代码之间的差距,面向熟悉 JavaScript 异步和基础 Rust 的开发者。
@debasishg:我关于Rust底层系统设计系列的第一部分现已发布 - 这部分涵盖:• 如何根据谁接触什么来布局共享的Rust结构体…
关于Rust底层系统设计系列的第一部分介绍了缓存感知的数据布局技术,包括字段分区以避免伪共享,重点涉及多线程结构体和128字节规则,并以SPSC环形缓冲区为例。
不会编译的数据竞争
本文解释了作者如何利用ruxe库中的类型级不相交技术,教会Rust的类型系统拒绝可能导致数据竞争的并行reducer管道。
从Go迁移到Rust
一份为Go开发者迁移到Rust编写的全面指南,专注于后端服务,对比正确性、运行时和人体工程学方面的权衡,并提供关于渐进式迁移的实用建议。
Ursula:基于线程每核心、多Raft架构的HTTP事件流Rust运行时
Ursula是一个开源、自托管的分布式服务器,用于可重放、仅追加的事件时间线,运行于HTTP和SSE之上,采用线程每核心、多Raft架构,并搭配S3存储以实现低延迟和持久性。