Construindo CLIs e Serviços de Dados
Entregue ETL rápido como CLIs baseadas em clap e serviços Tokio/Axum com observabilidade e modos de falha claros.
Busque em todas as páginas da documentação
Entregue ETL rápido como CLIs baseadas em clap e serviços Tokio/Axum com observabilidade e modos de falha claros.
Cartão de receita de referência rápida - pronto para copiar e colar.
use anyhow::{Context, Result};
use clap::Parser;
use polars::prelude::*;
use tracing::info;
#[derive(Parser)]
#[command(name = "csv2pq", about = "Converte CSV para Parquet com filtro opcional")]
struct Cli {
#[arg(long)]
input: String,
#[arg(long)]
output: String,
#[arg(long, default_value_t = 0)]
min_amount: i64,
}
fn main() -> Result<()> {
tracing_subscriber::fmt::init();
let cli = Cli::parse();
run(&cli).context("etl falhou")?;
Ok(())
}
fn run(cli: &Cli) -> PolarsResult<()> {
info!(input = %cli.input, "iniciando");
let mut df = CsvReadOptions::default()
.try_into_reader_with_file_path(Some(cli.input.clone().into()))?
.finish()?;
if cli.min_amount > 0 {
df = df.lazy().filter(col("amount").gt(lit(cli.min_amount))).collect()?;
}
ParquetWriter::new(std::fs::File::create(&cli.output)?).finish(&mut df)?;
info!(rows = df.height(), "concluído");
Ok(())
}Quando usar isso:
--help consistenteuse axum::{extract::Query, routing::get, Json, Router};
use polars::prelude::*;
use serde::Deserialize;
use std::net::SocketAddr;
use std::sync::Arc;
use tokio::sync::RwLock;
#[derive(Clone)]
struct AppState {
df: Arc<RwLock<DataFrame>>,
}
#[derive(Deserialize)]
struct Params {
region: String,
}
async fn summary(
Query(p): Query<Params>,
axum::extract::State(state): axum::extract::State<AppState>,
) -> Result<Json<serde_json::Value>, String> {
let df = state.df.read().await;
let filtered = df
.clone()
.lazy()
.filter(col("region").eq(lit(p.region.clone())))
.select([col("amount").sum().alias("total")])
.collect()
.map_err(|e| e.to_string())?;
let total: i64 = filtered.column("total").unwrap().get(0).unwrap().try_extract().unwrap();
Ok(Json(serde_json::json!({ "region": p.region, "total": total })))
}
#[tokio::main]
async fn main() -> PolarsResult<()> {
let df = CsvReadOptions::default()
.try_into_reader_with_file_path(Some("sales.csv".into()))?
.finish()?;
let state = AppState { df: Arc::new(RwLock::new(df)) };
let app = Router::new().route("/summary", get(summary)).with_state(state);
let addr = SocketAddr::from(([127, 0, 0, 1], 3000));
let listener = tokio::net::TcpListener::bind(addr).await.unwrap();
axum::serve(listener, app).await.unwrap();
Ok(())
}O que isso demonstra:
anyhow e logs tracingArc<RwLock<_>> para datasets recarregáveisargv uma vez, executam síncronamente ou bloqueiam em async com #[tokio::main].mmap Parquet para consultas repetidas.tracing) correlacionam tempos de estágio com caminhos de entrada e contagens de linhas.| Forma | Melhor para | Nota de Operações |
|---|---|---|
| CLI | Cron, transformações únicas | Versione o binário no repositório de artefatos |
| Serviço Axum | APIs internas de baixo QPS | Adicione autenticação, timeouts, limites de payload |
// Distinguir códigos de saída no main da CLI:
// std::process::exit(2) para argumentos ruins, 1 para erros de dados - ajuda orquestradores a tentar novamente corretamente.tracing ou indicatif.Result através de main e mapear para código de saída.mTLS ou middleware de token desde o início.| Alternativa | Usar Quando | Não Usar Quando |
|---|---|---|
| Python + Click | Equipe só conhece Python | Requisito de deploy de binário único |
| Airflow orquestrando CLI Rust | Dependências complexas de DAG | Trabalho simples de uma etapa noturna |
| gRPC + Arrow Flight | Service mesh de alta vazão | API JSON interna rápida |
| DuckDB CLI | SQL ad hoc em arquivos | Transformações Rust personalizadas no mesmo trabalho |
CLIs com foco em arquivos são frequentemente Polars síncrono com tokio opcional apenas para serviços. Não use async para collect limitado por CPU sem spawn_blocking.
Troque Arc<DataFrame> em um timer ou evento S3, ou mapeie em memória Parquet para benefícios de cache do SO.
Use tracing com tracing-subscriber JSON em produção; emparelhe spans por estágio do pipeline.
cargo build --release mais cargo deb opcional ou binários musl compilados cruzadamente - veja a seção de operações de deploy para padrões de distribuição.
Sim - #[arg(long, env = "INPUT_PATH")] ajuda contêineres sem longas listas de argv.
Use assert_cmd para executar o binário com CSVs de fixture temporários e afirmar o esquema Parquet de saída.
JSON para humanos e resumos pequenos; Arrow IPC para clientes de máquina puxando grandes conjuntos de resultados.
Leia credenciais S3 de variáveis de ambiente ou roles IAM - nunca incorpore chaves no binário.
Serde TOML para definições de pipeline; flags da CLI substituem os padrões do arquivo para execuções únicas.
Incorpore SQL na camada de serviço - veja DataFusion para manipuladores de consulta.
clapVersões da Stack: Esta página foi escrita para Rust 1.97.0 (edição 2024), Tokio 1.x, Axum 0.8, serde 1.0, sqlx 0.8, clap 4, e Polars 0.46+.
Revisado por Chris St. John·Última atualização: 16 de jul. de 2026