CULTIVA IA

Pipeline ETL — NutriTrack SaaS

Polars v1.41.2 · Analisis de MRR, churn y exportacion Parquet · Datos Analisis

45×
mas rapido
que pandas
Resultados del Pipeline
MRR Total
€187.4k
▲ 8.3% vs mes anterior
Empresas activas
1.142
▲ 58 nuevas en mayo
Clientes en riesgo
97
▼ MRR expuesto: €14.2k
Tiempo pipeline
2.1s
▲ Antes (pandas): 94s
Codigo del Pipeline Polars
scan_sources.py
SCAN LAZY
import polars as pl
from pathlib import Path

DATA = Path("data/nutritrack")

# Lazy scan — no lee nada todavia
lf_subs = pl.scan_csv(
    DATA / "subscriptions.csv",
    try_parse_dates=True
)

lf_events = pl.scan_parquet(
    DATA / "usage_events.parquet"
)

lf_plans = pl.scan_csv(
    DATA / "plans.csv"
)

# Polars sabe cuantas filas SIN leer
print(lf_events.collect_schema())
mrr_by_plan.py
MRR CALC
# MRR actual + variacion mensual por plan
mrr_df = (
    lf_subs
    .filter(pl.col("status") == "active")
    .join(lf_plans, on="plan_id", how="left")
    .group_by("plan_name", "month")
    .agg(
        pl.len().alias("n_clientes"),
        pl.col("price_eur").sum().alias("mrr_eur"),
    )
    .with_columns(
        # Variacion % vs mes anterior (window)
        mrr_prev=pl.col("mrr_eur")
            .shift(1)
            .over("plan_name"),
    )
    .with_columns(
        delta_pct=(
            (pl.col("mrr_eur") - pl.col("mrr_prev"))
            / pl.col("mrr_prev") * 100
        ).round(2)
    )
    .collect()  # ejecuta todo en paralelo
)
churn_risk.py
CHURN RISK
# Clientes con uso < 5 sesiones/mes = riesgo
usage_per_client = (
    lf_events
    .filter(
        pl.col("event_month") == "2026-05"
    )
    .group_by("company_id")
    .agg(
        pl.len().alias("sessions")
    )
    .collect()
)

# Clasificar riesgo con pl.when
at_risk = usage_per_client.with_columns(
    riesgo=pl.when(
        pl.col("sessions") < 3
    ).then(pl.lit("alto"))
    .when(
        pl.col("sessions") < 5
    ).then(pl.lit("medio"))
    .otherwise(pl.lit("bajo"))
).filter(
    pl.col("riesgo") != "bajo"
)
export_report.py
EXPORT
# Reporte final consolidado para BI
report = (
    lf_subs
    .join(
        lf_plans, on="plan_id", how="left"
    )
    .join(
        usage_per_client.lazy(),
        on="company_id", how="left"
    )
    .select(
        "company_id", "company_name",
        "plan_name", "price_eur",
        "sessions", "status",
        pl.col("start_date").dt.strftime("%Y-%m")
            .alias("mes_inicio")
    )
    .collect(engine="streaming")
)
# Parquet comprimido ~8x vs CSV
report.write_parquet(
    "output/nutritrack_report.parquet",
    compression="zstd"
)
MRR por Plan — Mayo 2026
mrr_by_plan DataFrame
18 filas · 5 columnas
Plan N Clientes MRR (€) Distribucion Delta vs Abr
Enterprise 187 74.613
39.8%
▲ 11.2%
Growth 512 76.288
40.7%
▲ 7.5%
Starter 443 21.707
11.6%
▼ 2.1%
Clientes en Riesgo de Churn
at_risk DataFrame · Mayo 2026
97 empresas
Empresa Plan Sesiones/mes Riesgo MRR expuesto
AlimentaPlus S.L. Enterprise 1 ● Alto €399
WellBeing Corp Growth 2 ● Alto €149
Saludify SaaS Growth 3 ● Alto €149
NutriBox Digital Starter 4 ◑ Medio €49
FoodLogic Pro Starter 4 ◑ Medio €49
... 92 filas mas ...
Benchmark vs Pandas
Dataset: 500k eventos + 1.2k suscripciones
scan_csv + filter
pandas read_csv
18.3s
polars scan_csv
1.4s
group_by + agg
pandas groupby
42.1s
polars group_by
0.3s
join + window
pandas merge+transform
33.6s
polars join + over()
0.4s
Total pipeline
94s
pandas
2.1s
polars
Plan de Ejecucion Lazy — Optimizacion Automatica
scan_csv
subscriptions.csv
filter
status == "active"
join
plans (left)
group_by
plan + month
OPTIMIZA
Predicate Push
filter antes del join
Projection Push
solo cols necesarias
Parallel Exec
8 threads CPU
collect()
resultado final