•15 min read

Building High-Throughput Microservices with Rust and Axum: Complete Production Guide

Building High-Throughput Microservices with Rust and Axum: Complete Production Guide

This guide details the construction of high-throughput microservices using Rust, Axum, and Tokio. We will cover project structure, advanced Tower middleware for critical operational concerns, robust database interaction with SQLx, graceful shutdown mechanisms, and optimized Docker multi-stage builds for minimal production images.

Audio Briefing
0:00 / 0:00

Project Structure and Dependencies

A well-organized project structure is paramount for maintainability and scalability. We advocate for a modular approach, separating concerns into distinct crates or modules.

.
├── Cargo.toml
├── src
│   ├── main.rs
│   ├── config.rs
│   ├── handlers.rs
│   ├── models.rs
│   ├── middleware
│   │   ├── mod.rs
│   │   ├── rate_limit.rs
│   │   └── tracing.rs
│   └── db.rs
└── Dockerfile

Our Cargo.toml will include essential dependencies:

# Cargo.toml
[package]
name = "high-throughput-service"
version = "0.1.0"
edition = "2021"

[dependencies]
# Web framework
axum = { version = "0.7", features = ["macros"] }
tokio = { version = "1", features = ["full"] }
tower = { version = "0.4", features = ["full"] }
tower-http = { version = "0.5", features = ["full"] }

# Database
sqlx = { version = "0.7", features = ["runtime-tokio-rustls", "postgres", "uuid", "chrono"] }
deadpool-redis = { version = "0.13", features = ["rt-tokio-1"] } # For rate limiting

# Serialization/Deserialization
serde = { version = "1", features = ["derive"] }
serde_json = "1"

# Configuration
config = { version = "0.13", features = ["toml", "yaml"] } # Or `envy` for env vars

# Logging and Tracing
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] }
opentelemetry = { version = "0.21", features = ["rt-tokio"] }
opentelemetry-sdk = { version = "0.21", features = ["rt-tokio"] }
opentelemetry-stdout = { version = "0.16" } # For local dev
opentelemetry-otlp = { version = "0.14", features = ["grpc-tonic", "reqwest-client", "tokio"] }
tracing-opentelemetry = "0.22"

# Utilities
uuid = { version = "1", features = ["v4", "serde"] }
chrono = { version = "0.4", features = ["serde"] }
anyhow = "1"
Advertisement

Configuration Management

Externalized configuration is critical. We'll use the config crate to load settings from environment variables and a configuration file.

// src/config.rs
use serde::Deserialize;
use std::net::SocketAddr;

#[derive(Debug, Deserialize, Clone)]
pub struct AppConfig {
    pub server_address: SocketAddr,
    pub database_url: String,
    pub redis_url: String,
    pub service_name: String,
    pub otel_exporter_otlp_endpoint: Option<String>,
}

impl AppConfig {
    pub fn load() -> Result<Self, config::ConfigError> {
        let s = config::Config::builder()
            .add_source(config::File::with_name("config.toml").required(false))
            .add_source(config::Environment::with_prefix("APP"))
            .build()?;
        s.try_deserialize()
    }
}

Database Connection Pooling with SQLx

SQLx provides compile-time checked queries and robust connection pooling. We'll use deadpool-redis for Redis.

// src/db.rs
use sqlx::postgres::PgPoolOptions;
use sqlx::PgPool;
use deadpool_redis::{Pool as RedisPool, Config as RedisConfig, Runtime};
use std::time::Duration;
use anyhow::Result;

pub async fn setup_pg_pool(database_url: &str) -> Result<PgPool> {
    let pool = PgPoolOptions::new()
        .max_connections(50)
        .min_connections(10)
        .acquire_timeout(Duration::from_secs(5))
        .connect_timeout(Duration::from_secs(5))
        .idle_timeout(Duration::from_secs(30 * 60)) // 30 minutes
        .test_before_acquire(true)
        .connect(database_url)
        .await?;
    Ok(pool)
}

pub fn setup_redis_pool(redis_url: &str) -> Result<RedisPool> {
    let cfg = RedisConfig::from_url(redis_url);
    let pool = cfg.create_pool(Some(Runtime::Tokio1))?;
    Ok(pool)
}

Tower Middleware for Operational Concerns

Tower's service-oriented architecture is powerful. We'll implement rate limiting, distributed tracing, and request validation.

Distributed Tracing

OpenTelemetry integration is crucial for observability.

// src/middleware/tracing.rs
use axum::{
    extract::Request,
    middleware::Next,
    response::Response,
};
use opentelemetry::{
    global,
    trace::{SpanKind, Tracer},
    Context, KeyValue,
};
use opentelemetry_sdk::{
    trace::{Config, TracerProvider},
    Resource,
};
use opentelemetry_otlp::WithExportConfig;
use tracing::{info_span, Instrument};
use tracing_opentelemetry::OpenTelemetrySpanExt;
use anyhow::Result;

pub fn init_tracer(service_name: &str, otlp_endpoint: Option<&str>) -> Result<()> {
    let resource = Resource::new(vec![
        KeyValue::new("service.name", service_name.to_string()),
        KeyValue::new("service.version", env!("CARGO_PKG_VERSION").to_string()),
    ]);

    let tracer_provider = if let Some(endpoint) = otlp_endpoint {
        opentelemetry_otlp::new_exporter()
            .tonic()
            .with_endpoint(endpoint)
            .build_tracer_provider_with_config(
                Config::default().with_resource(resource)
            )
    } else {
        // Fallback to stdout for local development if no OTLP endpoint is configured
        opentelemetry_stdout::new_pipeline()
            .with_trace_config(Config::default().with_resource(resource))
            .install_simple()
            .expect("Failed to install stdout tracer provider")
    };

    global::set_tracer_provider(tracer_provider);
    Ok(())
}

pub async fn trace_layer(req: Request, next: Next) -> Response {
    let path = req.uri().path().to_string();
    let method = req.method().to_string();

    let tracer = global::tracer("axum-server");
    let parent_cx = global::get_text_map_propagator(|propagator| {
        propagator.extract(&opentelemetry_http::HeaderExtractor(req.headers()))
    });

    let span = tracer
        .span_builder(format!("HTTP {} {}", method, path))
        .with_kind(SpanKind::Server)
        .with_parent_context(parent_cx)
        .start(&tracer);

    let cx = Context::current_with_span(span);
    let _guard = cx.attach();

    let response = next.run(req).instrument(info_span!("request_processing")).await;

    let span = cx.span();
    span.set_attribute(KeyValue::new("http.method", method));
    span.set_attribute(KeyValue::new("http.target", path));
    span.set_attribute(KeyValue::new("http.status_code", response.status().as_u16() as i64));
    span.end();

    response
}

Rate Limiting

A distributed rate limiter using Redis is essential for protecting resources.

// src/middleware/rate_limit.rs
use axum::{
    extract::{Request, State},
    http::StatusCode,
    middleware::Next,
    response::Response,
};
use deadpool_redis::redis::AsyncCommands;
use deadpool_redis::Pool as RedisPool;
use std::time::Duration;
use tracing::{error, info};

#[derive(Clone)]
pub struct RateLimiter {
    pub redis_pool: RedisPool,
    pub limit: u64,
    pub window_seconds: u64,
}

impl RateLimiter {
    pub fn new(redis_pool: RedisPool, limit: u64, window_seconds: u64) -> Self {
        Self {
            redis_pool,
            limit,
            window_seconds,
        }
    }

    pub async fn layer(State(rate_limiter): State<RateLimiter>, request: Request, next: Next) -> Result<Response, StatusCode> {
        let ip_address = request
            .headers()
            .get("X-Forwarded-For")
            .and_then(|h| h.to_str().ok())
            .and_then(|s| s.split(',').next()) // Take the first IP in case of multiple proxies
            .unwrap_or("unknown")
            .to_string();

        let key = format!("rate_limit:{}", ip_address);

        let mut conn = rate_limiter.redis_pool.get().await.map_err(|e| {
            error!("Failed to get Redis connection: {:?}", e);
            StatusCode::INTERNAL_SERVER_ERROR
        })?;

        let (count, _): (u64, u64) = deadpool_redis::redis::pipe()
            .atomic()
            .incr(&key, 1)
            .expire(&key, rate_limiter.window_seconds as usize)
            .query_async(&mut *conn)
            .await
            .map_err(|e| {
                error!("Failed to execute Redis rate limit command: {:?}", e);
                StatusCode::INTERNAL_SERVER_ERROR
            })?;

        if count > rate_limiter.limit {
            info!("Rate limit exceeded for IP: {}", ip_address);
            return Err(StatusCode::TOO_MANY_REQUESTS);
        }

        Ok(next.run(request).await)
    }
}
Advertisement

Handlers and Routing

Axum's handler functions are straightforward. We'll define a simple health check and a data endpoint.

// src/handlers.rs
use axum::{
    extract::{Path, State},
    http::StatusCode,
    response::{IntoResponse, Json},
};
use serde::{Deserialize, Serialize};
use sqlx::PgPool;
use uuid::Uuid;
use chrono::{DateTime, Utc};
use tracing::{info, error};

#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct Item {
    pub id: Uuid,
    pub name: String,
    pub description: Option<String>,
    pub created_at: DateTime<Utc>,
}

#[derive(Debug, Deserialize)]
pub struct CreateItem {
    pub name: String,
    pub description: Option<String>,
}

#[derive(Clone)]
pub struct AppState {
    pub pg_pool: PgPool,
    // Other shared state like RedisPool, etc.
}

pub async fn health_check() -> impl IntoResponse {
    (StatusCode::OK, "OK")
}

pub async fn create_item(
    State(state): State<AppState>,
    Json(payload): Json<CreateItem>,
) -> Result<Json<Item>, StatusCode> {
    info!("Attempting to create item: {}", payload.name);
    let new_item = sqlx::query_as!(
        Item,
        r#"
        INSERT INTO items (id, name, description, created_at)
        VALUES ($1, $2, $3, $4)
        RETURNING id, name, description, created_at
        "#,
        Uuid::new_v4(),
        payload.name,
        payload.description,
        Utc::now()
    )
    .fetch_one(&state.pg_pool)
    .await
    .map_err(|e| {
        error!("Failed to insert item: {:?}", e);
        StatusCode::INTERNAL_SERVER_ERROR
    })?;

    info!("Successfully created item with ID: {}", new_item.id);
    Ok(Json(new_item))
}

pub async fn get_item(
    State(state): State<AppState>,
    Path(item_id): Path<Uuid>,
) -> Result<Json<Item>, StatusCode> {
    info!("Attempting to retrieve item with ID: {}", item_id);
    let item = sqlx::query_as!(
        Item,
        r#"
        SELECT id, name, description, created_at
        FROM items
        WHERE id = $1
        "#,
        item_id
    )
    .fetch_optional(&state.pg_pool)
    .await
    .map_err(|e| {
        error!("Failed to query item: {:?}", e);
        StatusCode::INTERNAL_SERVER_ERROR
    })?
    .ok_or(StatusCode::NOT_FOUND)?;

    info!("Successfully retrieved item with ID: {}", item_id);
    Ok(Json(item))
}

Graceful Shutdown

Proper shutdown ensures no in-flight requests are abruptly terminated and resources are released.

// src/main.rs (excerpt)
// ... imports ...
use tokio::signal;
use tracing::info;

async fn shutdown_signal() {
    let ctrl_c = async {
        signal::ctrl_c()
            .await
            .expect("failed to install Ctrl+C handler");
    };

    #[cfg(unix)]
    let terminate = async {
        signal::unix::signal(signal::unix::SignalKind::terminate())
            .expect("failed to install SIGTERM handler")
            .recv()
            .await;
    };

    #[cfg(not(unix))]
    let terminate = std::future::pending::<()>();

    tokio::select! {
        _ = ctrl_c => {},
        _ = terminate => {},
    }

    info!("Shutdown signal received, initiating graceful shutdown...");
}

#[tokio::main]
async fn main() -> anyhow::Result<()> {
    // ... config loading, tracing setup, db setup ...

    let app = Router::new()
        .route("/health", get(health_check))
        .route("/items", post(create_item))
        .route("/items/:id", get(get_item))
        .with_state(app_state.clone())
        .layer(middleware::from_fn_with_state(rate_limiter.clone(), RateLimiter::layer))
        .layer(middleware::from_fn(trace_layer))
        .layer(TraceLayer::new_for_http()) // Axum's built-in tracing for request/response logging
        .layer(SetRequestIdLayer::new(
            Make<Uuid>::new(),
            PropagateRequestIdLayer::new(HeaderName::from_static("x-request-id")),
        ))
        .layer(TimeoutLayer::new(Duration::from_secs(30))); // Request timeout

    let listener = tokio::net::TcpListener::bind(&config.server_address).await?;
    info!("Listening on {}", config.server_address);

    axum::serve(listener, app.into_make_service())
        .with_graceful_shutdown(shutdown_signal())
        .await?;

    info!("Server gracefully shut down.");
    Ok(())
}

Docker Multi-Stage Builds

Optimized Docker images are crucial for fast deployments and reduced attack surface. We'll use a multi-stage build to produce a minimal scratch image.

# Dockerfile

# Stage 1: Builder
FROM rust:1.78-slim-bookworm AS builder

# Install build dependencies for SQLx and OpenSSL
RUN apt-get update && apt-get install -y \
    pkg-config \
    libssl-dev \
    postgresql-client \
    musl-tools \
    && rm -rf /var/lib/apt/lists/*

# Set the target for musl (static linking)
RUN rustup target add x86_64-unknown-linux-musl

WORKDIR /app

# Copy Cargo.toml and Cargo.lock first to leverage Docker cache
COPY Cargo.toml Cargo.lock ./

# Create a dummy src directory and main.rs to cache dependencies
RUN mkdir src && echo "fn main() {}" > src/main.rs
# Build dependencies only
RUN cargo build --release --target x86_64-unknown-linux-musl
# Remove dummy files
RUN rm -rf src target/x86_64-unknown-linux-musl/release/high-throughput-service

# Copy the actual source code
COPY src ./src

# Build the application
RUN cargo build --release --target x86_64-unknown-linux-musl

# Stage 2: Runtime
FROM scratch

# Set timezone data (important for many applications)
# FROM debian:bookworm-slim # Alternative if scratch is too restrictive for your needs
# RUN apt-get update && apt-get install -y tzdata && rm -rf /var/lib/apt/lists/*

WORKDIR /app

# Copy the compiled binary from the builder stage
COPY --from=builder /app/target/x86_64-unknown-linux-musl/release/high-throughput-service ./high-throughput-service

# Copy configuration file if present
COPY config.toml ./config.toml

# Expose the port the application listens on
EXPOSE 8080

# Set environment variables for configuration
ENV APP_SERVER_ADDRESS="0.0.0.0:8080"
ENV APP_DATABASE_URL="postgres://user:password@host:port/database"
ENV APP_REDIS_URL="redis://:password@host:port/"
ENV APP_SERVICE_NAME="high-throughput-service"
# ENV APP_OTEL_EXPORTER_OTLP_ENDPOINT="http://otel-collector:4317"

# Run the application
CMD ["./high-throughput-service"]

This Dockerfile produces an image typically under 20MB, containing only the statically linked binary and necessary configuration.

Architecture & Tradeoffs Comparison

Feature/AspectRust/Axum/TokioNode.js/Express/FastifyGo/Gin/Echo
PerformanceExcellent (near bare-metal)Good (event-driven, single-threaded bottleneck)Excellent (goroutines, compiled)
Memory FootprintVery LowModerate to High (V8 engine overhead)Low to Moderate
Concurrency ModelAsync/await (Tokio runtime)Event Loop (single-threaded, non-blocking I/O)Goroutines & Channels
Type SafetyStrong (compile-time)Weak (runtime, TypeScript improves this)Strong (compile-time)
Error HandlingResult enum, anyhow/thiserrorCallbacks, Promises, try/catchMultiple return values, error interface
Ecosystem MaturityGrowing rapidly, production-readyVery Mature, vast npm ecosystemMature, robust standard library
Learning CurveSteep (ownership, borrow checker, async)Moderate (JavaScript nuances, async patterns)Moderate (concurrency primitives, interfaces)
Binary SizeVery Small (static linking)Large (Node.js runtime included)Small (static linking)
Use CaseHigh-performance APIs, low-latency services, embedded, systems programmingRapid prototyping, I/O-bound services, web appsMicroservices, CLI tools, network services
TradeoffsHigher initial development cost, steep learning curve, excellent runtime safety and performance.Faster development, larger runtime, potential for runtime errors, less CPU-bound performance.Good balance of performance and development speed, simpler concurrency than Rust, less expressive type system.

Production Gotchas & Troubleshooting

1. database connection timed out or connection refused

Failure Mode: The application fails to connect to the PostgreSQL database on startup or during operation. Fixes:

  • Network Connectivity: Verify the database host is reachable from the microservice's environment (e.g., ping, telnet <db_host> <db_port>).
  • Firewall Rules: Ensure necessary ports (default 5432 for PostgreSQL) are open.
  • Database Credentials: Double-check APP_DATABASE_URL for correct username, password, host, port, and database name.
  • Database Availability: Confirm the PostgreSQL server is running and accepting connections.
  • Connection Pool Exhaustion: If errors occur during operation, increase max_connections in PgPoolOptions and monitor database connection limits.

2. Redis connection timed out or connection refused

Failure Mode: Rate limiting or other Redis-dependent features fail. Fixes:

  • Network Connectivity: Verify Redis host is reachable (e.g., ping, telnet <redis_host> <redis_port>).
  • Firewall Rules: Ensure Redis port (default 6379) is open.
  • Redis Credentials: Check APP_REDIS_URL for correct password and host/port.
  • Redis Availability: Confirm the Redis server is running.

3. High CPU Usage / Latency Spikes

Failure Mode: The service experiences unexpected high CPU load or slow response times under moderate load. Fixes:

  • Database Query Optimization: Analyze slow queries using EXPLAIN ANALYZE in PostgreSQL. Ensure proper indexing.
  • N+1 Query Problem: Identify and refactor handlers making multiple database calls in a loop. Use JOINs or batching.
  • CPU-Bound Operations: Profile Rust code to identify hot spots. Consider offloading heavy computations to background workers or optimizing algorithms.
  • Logging Verbosity: Excessive DEBUG or TRACE level logging in production can add significant overhead. Adjust RUST_LOG environment variable.

4. Too Many Open Files Error

Failure Mode: The service crashes with an OS-level error indicating it cannot open more files (sockets are files). Fixes:

  • ulimit: Increase the nofile limit for the user running the service. In Docker, this can be set with --ulimit nofile=65536:65536.
  • Connection Pool Settings: Ensure max_connections for database and Redis pools are reasonable and not excessively high, consuming too many file descriptors.
  • Resource Leaks: Investigate if connections or file handles are not being properly closed. Rust's RAII helps, but manual drop or close might be needed for external resources.

5. Tracing Data Not Appearing in Collector

Failure Mode: OpenTelemetry traces are not visible in your APM system (e.g., Jaeger, Datadog). Fixes:

  • APP_OTEL_EXPORTER_OTLP_ENDPOINT: Verify the endpoint is correctly configured and reachable from the microservice.
  • Collector Availability: Ensure the OpenTelemetry collector or APM agent is running and listening on the specified port.
  • Firewall: Check firewall rules between the service and the collector.
  • Sampling: If using a sampler, ensure traces are not being aggressively sampled away.
  • Service Name: Confirm APP_SERVICE_NAME is set, as it's crucial for identifying traces.

6. Docker Image Size Issues

Failure Mode: The final Docker image is unexpectedly large, despite using multi-stage builds. Fixes:

  • FROM scratch: Ensure the final stage truly uses FROM scratch or a minimal base image like distroless/static.
  • Unnecessary Files: Double-check that only the compiled binary and essential config are copied in the final stage. Avoid copying target directories or source code.
  • Static Linking: Ensure x86_64-unknown-linux-musl target is used for static linking to avoid needing glibc in the final image.
  • Build Cache Invalidation: If Cargo.toml or Cargo.lock change frequently, the dependency build cache might be invalidated. Structure the Dockerfile to copy these first.

Frequently Asked Questions

Q1: How do I handle request validation beyond basic JSON deserialization?

A1: For complex validation, integrate a validation library like validator. You can create a custom Axum extractor that uses validator to check incoming data.

// Example: Custom validator extractor
use axum::{
    async_trait,
    extract::{FromRequest, Request},
    http::StatusCode,
    response::{IntoResponse, Response},
    Json,
};
use serde::de::DeserializeOwned;
use validator::Validate;

pub struct ValidatedJson<T>(pub T);

#[async_trait]
impl<T, S> FromRequest<S> for ValidatedJson<T>
where
    T: DeserializeOwned + Validate,
    S: Send + Sync,
{
    type Rejection = Response;

    async fn from_request(req: Request, state: &S) -> Result<Self, Self::Rejection> {
        let Json(value) = Json::<T>::from_request(req, state)
            .await
            .map_err(|err| err.into_response())?;

        value.validate().map_err(|err| {
            (StatusCode::BAD_REQUEST, Json(err)).into_response()
        })?;

        Ok(ValidatedJson(value))
    }
}

// Usage in handler:
// pub async fn create_item(ValidatedJson(payload): ValidatedJson<CreateItem>) -> ...

Q2: What's the best way to manage shared state (like database pools) across handlers?

A2: Axum's State extractor is the idiomatic way. Define a struct (e.g., AppState) that holds all your shared resources, clone it, and pass it to the router's .with_state() method. Each handler can then extract State<AppState>. Ensure your state struct implements Clone.

Q3: How can I implement authentication and authorization?

A3: Implement authentication and authorization as Tower middleware. For authentication, a middleware can extract a token (e.g., JWT) from headers, validate it, and then insert user information into the request extensions. Subsequent authorization middleware or handler logic can then retrieve this user data. Libraries like jsonwebtoken are useful for JWT handling.

Q4: My service is crashing with SIGSEGV or other low-level errors. What should I do?

A4: SIGSEGV (segmentation fault) in Rust is rare but indicates memory corruption, often due to unsafe code or FFI (Foreign Function Interface) interactions.

  1. Review unsafe blocks: Carefully audit any unsafe code for correctness.
  2. FFI Boundaries: If interacting with C/C++ libraries, ensure correct memory management and type conversions across the FFI boundary.
  3. Dependency Issues: Check for known issues in your dependencies, especially those that use unsafe internally.
  4. Memory Sanitizers: While harder to integrate with Rust, tools like AddressSanitizer (ASan) can sometimes help debug native memory issues.
  5. Reproducibility: Try to create a minimal reproducible example to isolate the problem.

Q5: How do I handle database migrations in a production environment?

A5: Use sqlx-cli for managing migrations.

  1. Create Migrations: sqlx migrate add <migration_name>
  2. Apply Migrations: In your application's startup code, use sqlx::migrate!().run(&pool).await; to automatically apply pending migrations. This ensures your database schema is always up-to-date with your application code.
  3. Separate Migration Service: For more complex deployments, consider running migrations as a separate, pre-deployment step or as a dedicated Kubernetes init container.
// src/main.rs (excerpt for migrations)
// ...
#[tokio::main]
async fn main() -> anyhow::Result<()> {
    // ...
    let pg_pool = db::setup_pg_pool(&config.database_url).await?;

    // Apply database migrations
    info!("Running database migrations...");
    sqlx::migrate!("./migrations") // Path to your migrations directory
        .run(&pg_pool)
        .await?;
    info!("Database migrations applied successfully.");
    // ...
}
Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement