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

22. Базы данных

sqlx, diesel, connection pooling, транзакции, миграции, async DB drivers, repository pattern.

sqlx — async SQL с проверкой запросов на этапе компиляции

sqlx — асинхронная библиотека без ORM-абстракции: вы пишете обычный SQL, а sqlx проверяет его через реальное подключение к БД (или офлайн-кеш) прямо во время компиляции — макрос query! сверяет имена колонок и типы с фактической схемой.

use sqlx::postgres::PgPoolOptions;

#[tokio::main]
async fn main() -> Result<(), sqlx::Error> {
    let pool = PgPoolOptions::new()
        .max_connections(20)
        .connect("postgres://app:pass@localhost/app")
        .await?;

    // query! проверяет SQL и типы полей на этапе компиляции
    let user = sqlx::query!(
        "SELECT id, name, email FROM users WHERE id = $1",
        42_i64
    )
    .fetch_one(&pool)
    .await?;

    println!("{} <{}>", user.name, user.email);
    Ok(())
}

Для маппинга в конкретный тип структуры используется query_as!:

#[derive(Debug, sqlx::FromRow)]
struct User {
    id: i64,
    name: String,
    email: String,
}

let users: Vec<User> = sqlx::query_as!(User, "SELECT id, name, email FROM users")
    .fetch_all(&pool)
    .await?;
🔑 Ключевое: Макросы query!/query_as! требуют доступной БД во время сборки (или файла .sqlx, сгенерированного командой cargo sqlx prepare — используется в CI без живой БД).

diesel — ORM с безопасным query builder

diesel — синхронный (с недавним async-адаптером через diesel-async) ORM, где схема БД описывается Rust-типами, а запросы строятся через типобезопасный DSL — некорректный запрос не скомпилируется.

// schema.rs генерируется командой `diesel print-schema`
diesel::table! {
    users (id) {
        id -> Int8,
        name -> Text,
        email -> Text,
    }
}

#[derive(Queryable, Selectable)]
#[diesel(table_name = crate::schema::users)]
struct User {
    id: i64,
    name: String,
    email: String,
}

use self::schema::users::dsl::*;
use diesel::prelude::*;

let results = users
    .filter(email.like("%@example.com"))
    .limit(10)
    .select(User::as_select())
    .load(&mut conn)?;
Критерийsqlxdiesel
Модельсырой SQL, проверка типов через реальную БДquery builder на Rust-типах, схема как код
Asyncнативно async с самого началаисторически sync, async через отдельный адаптер
Гибкость SQLполная — любой SQL, включая CTE, оконные функцииограничена DSL, сложные запросы через raw_sql
Миграциивстроенный CLI sqlx-clidiesel_migrations, зрелая система
Когда выбратьсложные запросы, микросервисы с async I/Oкомандам, которые ценят compile-time схему и DSL

Connection pooling — управление пулом соединений

Открытие TCP-соединения и TLS-хендшейка с БД дорого — пул переиспользует соединения между запросами. Неправильный размер пула — частая причина деградации в проде: слишком большой пул перегружает БД, слишком маленький — создаёт очередь ожидания на стороне приложения.

use sqlx::postgres::PgPoolOptions;
use std::time::Duration;

let pool = PgPoolOptions::new()
    .max_connections(20)                       // формула: cores * 2 + effective_spindle_count
    .min_connections(2)                        // прогретые соединения на старте
    .acquire_timeout(Duration::from_secs(3))    // не ждать бесконечно свободное соединение
    .idle_timeout(Duration::from_secs(600))      // закрывать неиспользуемые соединения
    .max_lifetime(Duration::from_secs(1800))     // принудительная ротация (балансировщики, failover)
    .connect(&database_url)
    .await?;

// метрики пула для мониторинга
tracing::info!(
    size = pool.size(),
    idle = pool.num_idle(),
    "pool stats"
);
⚠️ Подводный камень: При нескольких репликах сервиса общий лимит соединений к БД — это max_connections × число_реплик. С PgBouncer в transaction pooling режиме учитывайте, что prepared statements и session-level настройки (например, SET) работают иначе — используйте statement_cache_capacity(0) при необходимости.

Транзакции — атомарность и изоляция

Транзакция в sqlx автоматически откатывается при drop, если явно не был вызван commit — это защищает от забытого rollback при раннем возврате через ?.

async fn transfer_funds(
    pool: &sqlx::PgPool,
    from: i64,
    to: i64,
    amount: i64,
) -> Result<(), sqlx::Error> {
    let mut tx = pool.begin().await?;

    sqlx::query!(
        "UPDATE accounts SET balance = balance - $1 WHERE id = $2 AND balance >= $1",
        amount, from
    )
    .execute(&mut *tx)
    .await?;

    sqlx::query!(
        "UPDATE accounts SET balance = balance + $1 WHERE id = $2",
        amount, to
    )
    .execute(&mut *tx)
    .await?;

    // если функция вернётся через ? до commit — tx откатится автоматически при Drop
    tx.commit().await?;
    Ok(())
}

Для уровней изоляции выше READ COMMITTED (например, SERIALIZABLE) нужно быть готовым к ошибкам сериализации и реализовать retry с экспоненциальной задержкой:

async fn with_serializable_retry<F, Fut, T>(pool: &sqlx::PgPool, mut op: F) -> Result<T, sqlx::Error>
where
    F: FnMut(sqlx::PgPool) -> Fut,
    Fut: std::future::Future<Output = Result<T, sqlx::Error>>,
{
    for attempt in 0..3 {
        match op(pool.clone()).await {
            Ok(v) => return Ok(v),
            Err(e) if is_serialization_failure(&e) && attempt < 2 => {
                tokio::time::sleep(std::time::Duration::from_millis(50 * 2u64.pow(attempt))).await;
            }
            Err(e) => return Err(e),
        }
    }
    unreachable!()
}

Миграции схемы

sqlx-cli и sqlx::migrate! позволяют версионировать схему БД вместе с кодом и применять миграции при старте сервиса или через отдельный CI-шаг.

// миграции лежат в ./migrations/0001_create_users.sql, 0002_add_index.sql...
// каждая миграция необратима по умолчанию — для отката пишут отдельный .down.sql

#[tokio::main]
async fn main() -> Result<(), sqlx::Error> {
    let pool = sqlx::PgPool::connect(&database_url).await?;

    // применяет все pending-миграции, записывает версию в _sqlx_migrations
    sqlx::migrate!("./migrations").run(&pool).await?;

    Ok(())
}
  • Идемпотентность DDL — используйте CREATE TABLE IF NOT EXISTS, CREATE INDEX CONCURRENTLY для больших таблиц без блокировки
  • Обратная совместимость — деплой миграции и деплой кода разнесены во времени, новая колонка должна быть nullable или иметь default до того, как старый код перестанет работать
  • CI-проверкаcargo sqlx prepare --check в pipeline гарантирует, что .sqlx-кеш синхронизирован со схемой

Repository pattern — изоляция слоя данных

Репозиторий скрывает детали конкретной СУБД за трейтом, что упрощает тестирование бизнес-логики через моки и позволяет заменить БД без переписывания сервисного слоя.

use async_trait::async_trait;

#[async_trait]
trait UserRepository: Send + Sync {
    async fn find_by_id(&self, id: i64) -> Result<Option<User>, RepoError>;
    async fn create(&self, user: NewUser) -> Result<User, RepoError>;
}

struct PgUserRepository {
    pool: sqlx::PgPool,
}

#[async_trait]
impl UserRepository for PgUserRepository {
    async fn find_by_id(&self, id: i64) -> Result<Option<User>, RepoError> {
        let user = sqlx::query_as!(User, "SELECT id, name, email FROM users WHERE id = $1", id)
            .fetch_optional(&self.pool)
            .await?;
        Ok(user)
    }

    async fn create(&self, new_user: NewUser) -> Result<User, RepoError> {
        let user = sqlx::query_as!(
            User,
            "INSERT INTO users (name, email) VALUES ($1, $2) RETURNING id, name, email",
            new_user.name, new_user.email
        )
        .fetch_one(&self.pool)
        .await?;
        Ok(user)
    }
}

Сервисный слой зависит от Arc<dyn UserRepository>, а не от конкретной реализации, что позволяет в тестах подставлять InMemoryUserRepository без сети и БД.

NoSQL и key-value хранилища

Для Redis используется redis-rs с async-мультиплексированным соединением, для MongoDB — официальный mongodb driver.

use redis::AsyncCommands;

async fn cache_user(client: &redis::Client, id: i64, json: &str) -> redis::RedisResult<()> {
    let mut conn = client.get_multiplexed_async_connection().await?;
    conn.set_ex::<_, _, ()>(format!("user:{id}"), json, 300).await?;
    Ok(())
}
🔑 Ключевое: get_multiplexed_async_connection в redis-rs безопасно клонируется и переиспользуется между задачами tokio — не нужен отдельный пул, как для SQL-соединений, потому что мультиплексирование команд идёт по одному TCP-соединению.

Диагностика медленных запросов

Проблема N+1 актуальна и для Rust-приложений: цикл с отдельным запросом на каждую запись — частая ошибка при переносе логики из ORM-стиля в SQL.

// плохо: N+1 — отдельный запрос на каждого пользователя
for user in &users {
    let orders = sqlx::query_as!(Order, "SELECT * FROM orders WHERE user_id = $1", user.id)
        .fetch_all(&pool)
        .await?;
}

// хорошо: один запрос с ANY($1) и группировка на стороне приложения
let ids: Vec<i64> = users.iter().map(|u| u.id).collect();
let orders = sqlx::query_as!(
    Order,
    "SELECT * FROM orders WHERE user_id = ANY($1)",
    &ids
)
.fetch_all(&pool)
.await?;

Для профилирования используйте EXPLAIN (ANALYZE, BUFFERS) и включённый в sqlx tracing-инстументацию: каждый запрос логируется со временем выполнения на уровне DEBUG.

🚫 Опасно: Не используйте format! для подстановки значений в SQL-строку — это открывает дорогу SQL-инъекциям. Всегда передавайте параметры через плейсхолдеры ($1, $2), даже с бинарными данными.