← Назад к списку тем

23. gRPC и сериализация

tonic и Protocol Buffers, serde (JSON/bincode/MessagePack), zero-copy сериализация, streaming.

Protocol Buffers и tonic

tonic — async gRPC-фреймворк на tokio и prost (генератор Rust-структур из .proto). Код генерируется на этапе сборки через build.rs и крейт tonic-build.

// proto/user.proto
syntax = "proto3";
package userpb;

message User {
  string id = 1;
  string email = 2;
  repeated string tags = 3;
}

service UserService {
  rpc GetUser(GetUserRequest) returns (User);
  rpc ListUsers(ListUsersRequest) returns (stream User);
}
// build.rs — кодогенерация во время сборки
fn main() -> Result<(), Box<dyn std::error::Error>> {
    tonic_build::compile_protos("proto/user.proto")?;
    Ok(())
}

Сгенерированный код подключается макросом include_proto!, а реализация сервиса — это трейт с async-методами, автоматически сгенерированный prost/tonic из описания service:

mod userpb {
    tonic::include_proto!("userpb");
}

use userpb::{User, GetUserRequest, user_service_server::{UserService, UserServiceServer}};
use tonic::{Request, Response, Status};

#[derive(Default)]
struct MyUserService;

#[tonic::async_trait]
impl UserService for MyUserService {
    async fn get_user(&self, req: Request<GetUserRequest>) -> Result<Response<User>, Status> {
        let id = req.into_inner().id;
        if id.is_empty() {
            return Err(Status::invalid_argument("id is required"));
        }
        Ok(Response::new(User { id, email: "user@example.com".into(), tags: vec![] }))
    }
    // ...
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    tonic::transport::Server::builder()
        .add_service(UserServiceServer::new(MyUserService::default()))
        .serve("0.0.0.0:50051".parse()?)
        .await?;
    Ok(())
}

Streaming RPC — server, client, bidirectional

tonic реализует все четыре паттерна gRPC. Стриминг возвращает impl Stream, обычно через tokio_stream и канал mpsc.

use tokio_stream::wrappers::ReceiverStream;
use tonic::Streaming;

type ListUsersStream = std::pin::Pin<Box<dyn tokio_stream::Stream<Item = Result<User, Status>> + Send>>;

async fn list_users(&self, _req: Request<ListUsersRequest>) -> Result<Response<ListUsersStream>, Status> {
    let (tx, rx) = tokio::sync::mpsc::channel(16);

    tokio::spawn(async move {
        for user in fetch_all_users().await {
            if tx.send(Ok(user)).await.is_err() {
                break; // клиент отменил стрим
            }
        }
    });

    Ok(Response::new(Box::pin(ReceiverStream::new(rx)) as ListUsersStream))
}

// Bidirectional streaming — двусторонний чат
async fn chat(&self, req: Request<Streaming<ChatMessage>>) -> Result<Response<ChatStream>, Status> {
    let mut inbound = req.into_inner();
    let (tx, rx) = tokio::sync::mpsc::channel(16);

    tokio::spawn(async move {
        while let Some(Ok(msg)) = inbound.message().await.transpose() {
            let reply = ChatMessage { text: format!("echo: {}", msg.text) };
            if tx.send(Ok(reply)).await.is_err() { break; }
        }
    });

    Ok(Response::new(Box::pin(ReceiverStream::new(rx)) as ChatStream))
}
⚠️ Подводный камень: Если receiver стрима на клиенте отменяет запрос (drop), tx.send вернёт ошибку — обязательно проверяйте её и прерывайте фоновую задачу, иначе она продолжит работать вхолостую (утечка ресурсов и БД-соединений).

Interceptors и метаданные

tonic использует tower-слои для сквозной функциональности: аутентификация, логирование, дедлайны — реализуются как Interceptor или обычный tower Layer.

use tonic::{Request, Status, service::Interceptor};

#[derive(Clone)]
struct AuthInterceptor;

impl Interceptor for AuthInterceptor {
    fn call(&mut self, req: Request<()>) -> Result<Request<()>, Status> {
        match req.metadata().get("authorization") {
            Some(t) if is_valid(t.to_str().unwrap_or_default()) => Ok(req),
            _ => Err(Status::unauthenticated("missing or invalid token")),
        }
    }
}

let svc = UserServiceServer::with_interceptor(MyUserService::default(), AuthInterceptor);

serde — универсальная сериализация

serde отделяет описание структуры данных (Serialize/Deserialize, выводятся через derive) от конкретного формата — один и тот же тип сериализуется в JSON, YAML, TOML, bincode или MessagePack без изменений кода.

#[derive(serde::Serialize, serde::Deserialize, Debug)]
struct Event {
    #[serde(rename = "eventId")]
    id: u64,
    name: String,
    #[serde(skip_serializing_if = "Option::is_none")]
    parent_id: Option<u64>,
    #[serde(default)]
    tags: Vec<String>,
}

let event = Event { id: 1, name: "deploy".into(), parent_id: None, tags: vec!["prod".into()] };

let json = serde_json::to_string(&event)?;          // текстовый, читаемый, самоописывающий
let bin = bincode::serialize(&event)?;               // компактный бинарный, только Rust-Rust
let msgpack = rmp_serde::to_vec(&event)?;             // бинарный, кросс-языковой аналог JSON
ФорматРазмерСкоростьСовместимостьКогда использовать
JSONбольше всехсредняяуниверсальная, читаема человекомпубличные API, отладка, конфиги
Protocol Buffersкомпактныйвысокаякросс-языковая, требует схему .protogRPC, версионированные контракты
MessagePackкомпактныйвысокаякросс-языковая, без схемызамена JSON там, где важен размер
bincodeминимальныймаксимальнаятолько Rust↔Rust, привязан к версии структурывнутренний IPC, кеш, снапшоты между Rust-сервисами
🚫 Опасно: bincode не имеет самоописывающей схемы — если формат структуры изменится (добавили/удалили поле), старые сериализованные данные не десериализуются корректно, а иногда молча дают мусорные значения. Не используйте bincode для долгоживущих данных на диске без версионирования.

Zero-copy сериализация

Обычные serde-форматы требуют аллокации при десериализации (копирование строк, векторов). Zero-copy подход — данные читаются прямо из буфера без копирования, что критично для low-latency систем.

// serde с заимствованием: &#39;a привязывает данные к времени жизни буфера
#[derive(serde::Deserialize)]
struct LogLine<'a> {
    #[serde(borrow)]
    message: &'a str,   // &str вместо String — без аллокации при парсинге
}

let buf = std::fs::read_to_string("log.json")?;
let line: LogLine = serde_json::from_str(&buf)?; // message ссылается прямо в buf

Для максимальной производительности используют rkyv — библиотеку архивной (zero-copy) сериализации, где сериализованные байты уже являются валидным представлением структуры в памяти, и десериализация не требуется вовсе, только валидация (опционально).

#[derive(rkyv::Archive, rkyv::Serialize, rkyv::Deserialize)]
struct Tick {
    symbol: String,
    price: f64,
}

let bytes = rkyv::to_bytes::<_, 256>(&tick)?;

// доступ к данным без десериализации — прямое чтение из архива
let archived = rkyv::check_archived_root::<Tick>(&bytes[..]).unwrap();
println!("{}", archived.price); // нет malloc, нет копирования полей
✅ Рекомендация: Для сетевых API между независимыми командами/языками используйте protobuf или MessagePack. Zero-copy форматы (rkyv, borrowed serde) оправданы там, где сериализация/десериализация — измеренное узкое место (маркет-дата фиды, шаринг памяти между процессами, hot path).

Обработка ошибок в gRPC

tonic использует Status с кодами, совместимыми со стандартом gRPC (codes.google.com). Как и в HTTP, внутренние детали ошибок не стоит пробрасывать наружу напрямую.

impl From<RepoError> for Status {
    fn from(err: RepoError) -> Self {
        match err {
            RepoError::NotFound => Status::not_found("resource not found"),
            RepoError::Conflict(msg) => Status::already_exists(msg),
            RepoError::Database(e) => {
                tracing::error!(error = %e, "db error in grpc handler");
                Status::internal("internal error")
            }
        }
    }
}

Reflection и отладка

gRPC reflection позволяет инструментам вроде grpcurl и evans опрашивать сервис без наличия .proto-файлов на клиенте — полезно для отладки в проде.

let reflection_service = tonic_reflection::server::Builder::configure()
    .register_encoded_file_descriptor_set(userpb::FILE_DESCRIPTOR_SET)
    .build()?;

tonic::transport::Server::builder()
    .add_service(UserServiceServer::new(MyUserService::default()))
    .add_service(reflection_service)
    .serve(addr)
    .await?;

// $ grpcurl -plaintext localhost:50051 list
// $ grpcurl -plaintext -d '{"id":"1"}' localhost:50051 userpb.UserService/GetUser